From 4bd1e13913e78be4cfa06bce609a520ffcd62c98 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Sun, 27 Sep 2026 18:42:54 -0700 Subject: [PATCH 1/2] refactor(db): require a transaction for transaction-scoped advisory locks --- .../lib/application/authorized-use-case.ts | 100 ++++++++++++------ .../access-requests/lib/application/review.ts | 20 ++-- .../ee/scim/lib/application/admin/mappings.ts | 4 +- .../lib/application/groups/manage-groups.ts | 4 +- .../scim/lib/application/users/update-user.ts | 6 +- apps/sim/ee/scim/lib/projection/auto-map.ts | 4 +- .../ee/scim/lib/projection/reconcile-user.ts | 12 +-- .../workspace-forking/application/revision.ts | 4 +- .../lib/copy/workflow-mcp-attachments.ts | 4 +- .../workspace-forking/lib/lineage/lineage.ts | 12 ++- apps/sim/lib/billing/organization.ts | 4 +- .../organizations/billing-identity-lock.ts | 7 +- .../organizations/create-organization.ts | 4 +- .../lib/billing/organizations/membership.ts | 17 +-- .../billing/organizations/provision-seat.ts | 8 +- apps/sim/lib/credential-groups/enrollments.ts | 6 +- apps/sim/lib/credential-groups/service.ts | 14 +-- .../lib/credential-groups/shared-slack-app.ts | 4 +- apps/sim/lib/credentials/env-locks.ts | 8 +- apps/sim/lib/credentials/environment.ts | 10 +- apps/sim/lib/credentials/personal-tokens.ts | 21 ++-- apps/sim/lib/db/advisory-locks.ts | 10 +- apps/sim/lib/db/transaction.ts | 4 +- apps/sim/lib/folders/locks.ts | 4 +- apps/sim/lib/invitations/core.ts | 6 +- apps/sim/lib/invitations/locks.ts | 4 +- apps/sim/lib/invitations/resend-policy.ts | 4 +- apps/sim/lib/invitations/send.ts | 4 +- .../lib/invitations/workspace-invitations.ts | 4 +- apps/sim/lib/mcp/server-locks.ts | 9 +- apps/sim/lib/mcp/workflow-mcp-sync.ts | 8 +- .../lib/organizations/members/lifecycle.ts | 6 +- .../application/group-membership.ts | 6 +- apps/sim/lib/permission-groups/locks.ts | 4 +- apps/sim/lib/sim-search/live/member-setup.ts | 4 +- .../workspace-file-folder-manager.ts | 4 +- .../persistence/deployment-operations.ts | 6 +- apps/sim/lib/workspaces/create.ts | 6 +- .../sim/lib/workspaces/operations/receipts.ts | 4 +- .../lib/workspaces/organization-workspaces.ts | 6 +- apps/sim/lib/workspaces/policy.ts | 4 +- .../scripts/migrate-gitlab-personal-tokens.ts | 8 +- 42 files changed, 226 insertions(+), 162 deletions(-) diff --git a/apps/sim/ee/access-requests/lib/application/authorized-use-case.ts b/apps/sim/ee/access-requests/lib/application/authorized-use-case.ts index 98a5241a7e4..2780b5748d7 100644 --- a/apps/sim/ee/access-requests/lib/application/authorized-use-case.ts +++ b/apps/sim/ee/access-requests/lib/application/authorized-use-case.ts @@ -8,7 +8,7 @@ import { import type { OperationUseCase } from '@/lib/core/application/operation' import { requireAllowedWorkspacePrincipal } from '@/lib/core/application/workspace-authorization' import { runWithOutboundOrganization } from '@/lib/core/network/context.server' -import type { DbOrTx } from '@/lib/db/types' +import type { DbClient, DbOrTx, DbTransaction } from '@/lib/db/types' import { type AccessRequestContext, authorizeAccessRequestScope, @@ -25,25 +25,40 @@ interface AccessRequestPreparationArgs { context: AccessRequestContext } -interface AccessRequestUseCaseArgs extends AccessRequestPreparationArgs { - executor: DbOrTx +/** A mutation executes inside the funnel's transaction; a read runs on the pool-level client. */ +interface AccessRequestUseCaseArgs + extends AccessRequestPreparationArgs { + executor: E } interface AccessRequestUseCaseDefinition { operation: AccessRequestOperation scope(input: I): AccessRequestScope - mutation?: boolean projectAudit?(args: AccessRequestUseCaseArgs & { result: R }): WorkspaceUseCaseAuditEntry[] } -interface PreparedAccessRequestUseCase extends AccessRequestUseCaseDefinition { +interface PreparedAccessRequestUseCase + extends AccessRequestUseCaseDefinition { prepare(args: AccessRequestPreparationArgs): Promise

- execute(args: AccessRequestUseCaseArgs & { prepared: P }): Promise + execute(args: AccessRequestUseCaseArgs & { prepared: P }): Promise } -interface UnpreparedAccessRequestUseCase extends AccessRequestUseCaseDefinition { +interface UnpreparedAccessRequestUseCase + extends AccessRequestUseCaseDefinition { prepare?: never - execute(args: AccessRequestUseCaseArgs & { prepared: undefined }): Promise + execute(args: AccessRequestUseCaseArgs & { prepared: undefined }): Promise +} + +type AccessRequestUseCase = + | PreparedAccessRequestUseCase + | UnpreparedAccessRequestUseCase + +type MutationAccessRequestUseCase = AccessRequestUseCase & { + mutation: true +} + +type ReadAccessRequestUseCase = AccessRequestUseCase & { + mutation?: false } function requireAccessRequestPrincipal( @@ -53,15 +68,33 @@ function requireAccessRequestPrincipal( requireAllowedWorkspacePrincipal(principal, operation) } +/** Runs preparation before any transaction opens and binds its result to `execute`. */ +async function prepareExecution( + definition: AccessRequestUseCase, + args: AccessRequestPreparationArgs +): Promise<(args: AccessRequestUseCaseArgs) => Promise> { + if (definition.prepare) { + const prepared = await definition.prepare(args) + const executePrepared = definition.execute + return (executeArgs) => executePrepared({ ...executeArgs, prepared }) + } + const executeUnprepared = definition.execute + return (executeArgs) => executeUnprepared({ ...executeArgs, prepared: undefined }) +} + export function defineAuthorizedAccessRequestUseCase( - definition: PreparedAccessRequestUseCase + definition: + | (PreparedAccessRequestUseCase & { mutation: true }) + | (PreparedAccessRequestUseCase & { mutation?: false }) ): OperationUseCase export function defineAuthorizedAccessRequestUseCase( - definition: UnpreparedAccessRequestUseCase + definition: + | (UnpreparedAccessRequestUseCase & { mutation: true }) + | (UnpreparedAccessRequestUseCase & { mutation?: false }) ): OperationUseCase /** Shared human-credential funnel; preparation finishes before any transaction acquires locks. */ export function defineAuthorizedAccessRequestUseCase( - definition: PreparedAccessRequestUseCase | UnpreparedAccessRequestUseCase + definition: MutationAccessRequestUseCase | ReadAccessRequestUseCase ): OperationUseCase { return { operation: definition.operation, @@ -74,32 +107,29 @@ export function defineAuthorizedAccessRequestUseCase( const scope = definition.scope(input) const initial = await authorizeAccessRequestScope(principal, definition.operation, scope) return runWithOutboundOrganization(initial.organizationId, async () => { - let execute: (args: AccessRequestUseCaseArgs) => Promise - if (definition.prepare) { - const prepared = await definition.prepare({ principal, input, context: initial }) - const executePrepared = definition.execute - execute = (args) => executePrepared({ ...args, prepared }) + const preparation = { principal, input, context: initial } + let context = initial + let result: R + if (definition.mutation) { + const execute = await prepareExecution(definition, preparation) + result = await db.transaction(async (executor) => { + if (initial.organizationId) { + await acquireOrganizationMutationLock(executor, initial.organizationId) + } + context = await authorizeAccessRequestScope( + principal, + definition.operation, + scope, + executor, + true, + initial + ) + return execute({ principal, input, context, executor }) + }) } else { - const executeUnprepared = definition.execute - execute = (args) => executeUnprepared({ ...args, prepared: undefined }) + const execute = await prepareExecution(definition, preparation) + result = await execute({ principal, input, context, executor: db }) } - let context = initial - const result = definition.mutation - ? await db.transaction(async (executor) => { - if (initial.organizationId) { - await acquireOrganizationMutationLock(executor, initial.organizationId) - } - context = await authorizeAccessRequestScope( - principal, - definition.operation, - scope, - executor, - true, - initial - ) - return execute({ principal, input, context, executor }) - }) - : await execute({ principal, input, context, executor: db }) if (definition.projectAudit) { recordProjectedUseCaseAuditEntries( definition.operation, diff --git a/apps/sim/ee/access-requests/lib/application/review.ts b/apps/sim/ee/access-requests/lib/application/review.ts index e1a45306dc3..19d68dabe3a 100644 --- a/apps/sim/ee/access-requests/lib/application/review.ts +++ b/apps/sim/ee/access-requests/lib/application/review.ts @@ -8,7 +8,7 @@ import { setOrgMemberUsageLimit } from '@/lib/billing/organizations/member-limit import type { WorkspaceUseCaseAuditEntry } from '@/lib/core/application/authorized-workspace-use-case' import { OrchestrationError } from '@/lib/core/orchestration/types' import { enqueueOutboxEvent } from '@/lib/core/outbox/service' -import type { DbOrTx } from '@/lib/db/types' +import type { DbTransaction } from '@/lib/db/types' import { acquirePermissionGroupOrgLock } from '@/lib/permission-groups/locks' import { loadAccessRequestMembership } from '@/ee/access-requests/lib/application/authorization' import { defineAuthorizedAccessRequestUseCase } from '@/ee/access-requests/lib/application/authorized-use-case' @@ -52,13 +52,17 @@ const organizationScope = (input: ReviewInput): AccessRequestScope => ({ organizationId: input.organizationId, }) -/** Checks the requester's present scope before inspecting or modifying any governing policy. */ +/** + * Checks the requester's present scope before inspecting or modifying any governing policy. + * Given `lockingTx`, it locks the scope rows and the governing policy in that transaction. + */ async function loadReviewPreview( - executor: DbOrTx, row: StoredAccessRequest, prepared: PreparedAccessRequestPolicy | null, - forUpdate = false + lockingTx?: DbTransaction ) { + const executor = lockingTx ?? db + const forUpdate = Boolean(lockingTx) const request = await presentAccessRequest(executor, row) if (row.status === 'fulfilled' && row.decision) { const snapshot = storedAccessRequestDecisionSchema.parse(row.decision) @@ -118,8 +122,8 @@ async function loadReviewPreview( membershipId: membership?.membershipId ?? '', role: membership?.role ?? ('read' as const), } - if (forUpdate) - await acquirePermissionGroupOrgLock(executor, row.organizationId, { + if (lockingTx) + await acquirePermissionGroupOrgLock(lockingTx, row.organizationId, { lockTimeoutAlreadyBounded: true, }) const catalog = prepared.catalog @@ -233,7 +237,7 @@ export const previewAccessRequest = defineAuthorizedAccessRequestUseCase({ }, async execute({ input, executor, prepared }) { const row = await loadStoredAccessRequest(executor, input.organizationId, input.requestId) - return (await loadReviewPreview(executor, row, prepared)).preview + return (await loadReviewPreview(row, prepared)).preview }, }) @@ -275,7 +279,7 @@ export const resolveAccessRequest = defineAuthorizedAccessRequestUseCase({ } if (!prepared) throw new OrchestrationError('internal', 'Request preview preparation is missing') - const { preview, policy } = await loadReviewPreview(executor, row, prepared, true) + const { preview, policy } = await loadReviewPreview(row, prepared, executor) if (!preview.canApply) throw new OrchestrationError( 'conflict', diff --git a/apps/sim/ee/scim/lib/application/admin/mappings.ts b/apps/sim/ee/scim/lib/application/admin/mappings.ts index 563afe2c758..62b8e5d762b 100644 --- a/apps/sim/ee/scim/lib/application/admin/mappings.ts +++ b/apps/sim/ee/scim/lib/application/admin/mappings.ts @@ -13,7 +13,7 @@ import { and, count, eq, sql } from 'drizzle-orm' import type { ScimGroupMappingView } from '@/lib/api/contracts/organization-scim' import { acquireOrganizationMutationLock } from '@/lib/billing/organizations/membership' import { OrchestrationError } from '@/lib/core/orchestration/types' -import type { DbOrTx } from '@/lib/db/types' +import type { DbTransaction } from '@/lib/db/types' import { acquirePermissionGroupOrgLock } from '@/lib/permission-groups/locks' import { assertWorkspaceInOrganization, @@ -147,7 +147,7 @@ async function requireGroup(connectionId: string, groupId: string) { } async function assertPermissionGroupTarget( - tx: DbOrTx, + tx: DbTransaction, organizationId: string, permissionGroupId: string ) { diff --git a/apps/sim/ee/scim/lib/application/groups/manage-groups.ts b/apps/sim/ee/scim/lib/application/groups/manage-groups.ts index 41f71151eaa..16a303b2871 100644 --- a/apps/sim/ee/scim/lib/application/groups/manage-groups.ts +++ b/apps/sim/ee/scim/lib/application/groups/manage-groups.ts @@ -4,7 +4,7 @@ import { scimGroup } from '@sim/db/schema' import { and, eq, ne } from 'drizzle-orm' import type { ScimPatchOperation } from '@/lib/api/contracts/scim' import { acquireOrganizationMutationLock } from '@/lib/billing/organizations/membership' -import type { DbOrTx } from '@/lib/db/types' +import type { DbOrTx, DbTransaction } from '@/lib/db/types' import { defineAuthorizedScimUseCase, type ScimUseCaseArgs, @@ -54,7 +54,7 @@ import { */ function withGroupWrite( context: ScimUseCaseContext, - work: (tx: DbOrTx) => Promise + work: (tx: DbTransaction) => Promise ): Promise { return db.transaction(async (tx) => { await acquireOrganizationMutationLock(tx, context.organizationId) diff --git a/apps/sim/ee/scim/lib/application/users/update-user.ts b/apps/sim/ee/scim/lib/application/users/update-user.ts index 0b69ff1f617..900f814a5be 100644 --- a/apps/sim/ee/scim/lib/application/users/update-user.ts +++ b/apps/sim/ee/scim/lib/application/users/update-user.ts @@ -4,7 +4,7 @@ import type { ScimUserAttributes } from '@sim/db/schema' import { normalizeEmail } from '@sim/utils/string' import type { ScimPatchOperation } from '@/lib/api/contracts/scim' import { acquireOrganizationUserMutationLocks } from '@/lib/billing/organizations/membership' -import type { DbOrTx } from '@/lib/db/types' +import type { DbOrTx, DbTransaction } from '@/lib/db/types' import { suspendMemberTx, unsuspendMemberTx } from '@/lib/organizations/members/lifecycle' import { invalidateAfterSessionRevocation, @@ -49,7 +49,7 @@ export interface UpdateOutcome { } async function applyUserUpdate( - tx: DbOrTx, + tx: DbTransaction, context: ScimUseCaseContext, current: ScimUserRecord, next: ScimUserAttributes @@ -135,7 +135,7 @@ export interface UpdateScimUserResult { * which take the advisory locks first and then touch rows referencing this one. */ async function loadUserForUpdate( - tx: DbOrTx, + tx: DbTransaction, context: ScimUseCaseContext, scimUserId: string ): Promise { diff --git a/apps/sim/ee/scim/lib/projection/auto-map.ts b/apps/sim/ee/scim/lib/projection/auto-map.ts index e1924c760ac..629940e52a0 100644 --- a/apps/sim/ee/scim/lib/projection/auto-map.ts +++ b/apps/sim/ee/scim/lib/projection/auto-map.ts @@ -1,7 +1,7 @@ import { permissionGroup, scimGroupMapping } from '@sim/db/schema' import { generateId } from '@sim/utils/id' import { and, eq, inArray, ne } from 'drizzle-orm' -import type { DbOrTx } from '@/lib/db/types' +import type { DbOrTx, DbTransaction } from '@/lib/db/types' import { acquirePermissionGroupOrgLock } from '@/lib/permission-groups/locks' /** @@ -94,7 +94,7 @@ export async function autoMapPermissionGroupByName( * mode, and both commit together. */ export async function settleMappedPermissionGroupsExplicit( - tx: DbOrTx, + tx: DbTransaction, params: { organizationId: string; scimGroupId: string } ): Promise { const inheriting = await tx diff --git a/apps/sim/ee/scim/lib/projection/reconcile-user.ts b/apps/sim/ee/scim/lib/projection/reconcile-user.ts index 69260fe34a7..696f4290945 100644 --- a/apps/sim/ee/scim/lib/projection/reconcile-user.ts +++ b/apps/sim/ee/scim/lib/projection/reconcile-user.ts @@ -16,7 +16,7 @@ import { generateId } from '@sim/utils/id' import { and, eq, inArray } from 'drizzle-orm' import { acquireOrganizationUserMutationLocks } from '@/lib/billing/organizations/membership' import { OrchestrationError } from '@/lib/core/orchestration/types' -import type { DbOrTx } from '@/lib/db/types' +import type { DbOrTx, DbTransaction } from '@/lib/db/types' import { changeMemberRoleTx } from '@/lib/organizations/members/lifecycle' import { addPermissionGroupMemberTx, @@ -195,7 +195,7 @@ async function findForeignWorkspaces( /** Applies one grant. `skipped` means the grant describes nothing this server can apply. */ async function applyGrant( - tx: DbOrTx, + tx: DbTransaction, params: { organizationId: string userId: string @@ -259,7 +259,7 @@ async function applyGrant( * is simply inert for that one person. */ async function setOrganizationRole( - tx: DbOrTx, + tx: DbTransaction, organizationId: string, userId: string, role: 'admin' | 'member' @@ -293,7 +293,7 @@ async function setOrganizationRole( * has made the directory the source of truth. */ async function withdrawGrant( - tx: DbOrTx, + tx: DbTransaction, params: { organizationId: string userId: string @@ -398,7 +398,7 @@ async function withdrawGrant( * without a dry-run mode. */ export async function reconcileUserProjection( - tx: DbOrTx, + tx: DbTransaction, params: { connectionId: string organizationId: string @@ -576,7 +576,7 @@ export async function reconcileUserProjection( /** Reconciles several users, in a stable order so concurrent syncs cannot deadlock. */ export async function reconcileUsersProjection( - tx: DbOrTx, + tx: DbTransaction, params: { connectionId: string organizationId: string diff --git a/apps/sim/ee/workspace-forking/application/revision.ts b/apps/sim/ee/workspace-forking/application/revision.ts index 3c27d3c592e..cb3fb69beeb 100644 --- a/apps/sim/ee/workspace-forking/application/revision.ts +++ b/apps/sim/ee/workspace-forking/application/revision.ts @@ -23,7 +23,7 @@ import { workspaceSandbox, } from '@sim/db/schema' import { and, type SQL, sql } from 'drizzle-orm' -import type { DbOrTx } from '@/lib/db/types' +import type { DbOrTx, DbTransaction } from '@/lib/db/types' import { acquireFolderMutationLock } from '@/lib/folders/locks' import { activeWorkspaceFileConditions } from '@/lib/workspace-files/query-scope' import { @@ -141,7 +141,7 @@ export async function loadForkPreviewRevision( } /** Locks normalized graph rows as well as workflow metadata, including realtime-only writes. */ -export async function lockForkRevision(tx: DbOrTx, scope: ForkRevisionScope): Promise { +export async function lockForkRevision(tx: DbTransaction, scope: ForkRevisionScope): Promise { const workspaceIds = [ ...new Set([ scope.sourceWorkspaceId, diff --git a/apps/sim/ee/workspace-forking/lib/copy/workflow-mcp-attachments.ts b/apps/sim/ee/workspace-forking/lib/copy/workflow-mcp-attachments.ts index 2a5168dfa47..affca75de72 100644 --- a/apps/sim/ee/workspace-forking/lib/copy/workflow-mcp-attachments.ts +++ b/apps/sim/ee/workspace-forking/lib/copy/workflow-mcp-attachments.ts @@ -1,7 +1,7 @@ import { workflowMcpServer, workflowMcpTool } from '@sim/db/schema' import { generateId } from '@sim/utils/id' import { and, eq, inArray, isNull } from 'drizzle-orm' -import type { DbOrTx } from '@/lib/db/types' +import type { DbOrTx, DbTransaction } from '@/lib/db/types' import { acquireWorkflowMcpServerLock } from '@/lib/mcp/server-locks' import { validateMcpToolMetadataForStorage } from '@/lib/mcp/tool-limits' import { getEdgeMappingRows } from '@/ee/workspace-forking/lib/mapping/mapping-store' @@ -102,7 +102,7 @@ export async function copyForkWorkflowMcpAttachments(params: { * Returns the affected target server ids so the caller can notify them post-commit. */ export async function reconcileForkWorkflowMcpAttachments(params: { - tx: DbOrTx + tx: DbTransaction childWorkspaceId: string /** True when the sync SOURCE is the parent workspace (a pull). */ sourceIsParent: boolean diff --git a/apps/sim/ee/workspace-forking/lib/lineage/lineage.ts b/apps/sim/ee/workspace-forking/lib/lineage/lineage.ts index a3471fa0d17..b25f4689b92 100644 --- a/apps/sim/ee/workspace-forking/lib/lineage/lineage.ts +++ b/apps/sim/ee/workspace-forking/lib/lineage/lineage.ts @@ -2,7 +2,7 @@ import { db } from '@sim/db' import { workspace } from '@sim/db/schema' import { and, desc, eq, isNull, sql } from 'drizzle-orm' import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' -import type { DbOrTx } from '@/lib/db/types' +import type { DbOrTx, DbTransaction } from '@/lib/db/types' export interface ForkLineageNode { id: string @@ -110,7 +110,10 @@ export async function setForkLockTimeout(tx: DbOrTx): Promise { * between distinct keys astronomically unlikely; a collision would only cause * unnecessary serialization, never a correctness issue. */ -export async function acquireForkEdgeLock(tx: DbOrTx, childWorkspaceId: string): Promise { +export async function acquireForkEdgeLock( + tx: DbTransaction, + childWorkspaceId: string +): Promise { await acquireAdvisoryXactLock(tx, 'fork_edge', `fork-edge:${childWorkspaceId}`) } @@ -121,6 +124,9 @@ export async function acquireForkEdgeLock(tx: DbOrTx, childWorkspaceId: string): * interleaving and keeping rollback's "newest sync" check race-free. Always acquire * this BEFORE {@link acquireForkEdgeLock} so the two are taken in a consistent order. */ -export async function acquireForkTargetLock(tx: DbOrTx, targetWorkspaceId: string): Promise { +export async function acquireForkTargetLock( + tx: DbTransaction, + targetWorkspaceId: string +): Promise { await acquireAdvisoryXactLock(tx, 'fork_target', `fork-target:${targetWorkspaceId}`) } diff --git a/apps/sim/lib/billing/organization.ts b/apps/sim/lib/billing/organization.ts index fc6219b250c..c937b6d6eb9 100644 --- a/apps/sim/lib/billing/organization.ts +++ b/apps/sim/lib/billing/organization.ts @@ -14,7 +14,7 @@ import { acquireOrganizationMutationLock } from '@/lib/billing/organizations/mem import { isEnterprise, isOrgPlan, isPaid } from '@/lib/billing/plan-helpers' import { ENTITLED_SUBSCRIPTION_STATUSES } from '@/lib/billing/subscriptions/utils' import { toDecimal } from '@/lib/billing/utils/decimal' -import type { DbOrTx } from '@/lib/db/types' +import type { DbTransaction } from '@/lib/db/types' import { attachOwnedWorkspacesToOrganization, attachOwnedWorkspacesToOrganizationTx, @@ -350,7 +350,7 @@ export async function ensureOrganizationForTeamSubscription( * resolution and workspace attachment through the caller's transaction. */ export async function ensureOrganizationForTeamSubscriptionTx( - tx: DbOrTx, + tx: DbTransaction, subscription: SubscriptionData & { workspaceIdsToAttach: string[] } ): Promise { if (!isOrgPlan(subscription.plan)) { diff --git a/apps/sim/lib/billing/organizations/billing-identity-lock.ts b/apps/sim/lib/billing/organizations/billing-identity-lock.ts index 3be8789aae2..5c615bde155 100644 --- a/apps/sim/lib/billing/organizations/billing-identity-lock.ts +++ b/apps/sim/lib/billing/organizations/billing-identity-lock.ts @@ -1,6 +1,6 @@ import { sql } from 'drizzle-orm' import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' -import type { DbOrTx } from '@/lib/db/types' +import type { DbTransaction } from '@/lib/db/types' const USER_BILLING_IDENTITY_LOCK_TIMEOUT_MS = 5_000 @@ -9,7 +9,10 @@ const USER_BILLING_IDENTITY_LOCK_TIMEOUT_MS = 5_000 * organization billed. Organization locks alone are insufficient because a * personal credit grant does not have an organization id when it begins. */ -export async function acquireUserBillingIdentityLock(tx: DbOrTx, userId: string): Promise { +export async function acquireUserBillingIdentityLock( + tx: DbTransaction, + userId: string +): Promise { await tx.execute( sql`select set_config('lock_timeout', ${`${USER_BILLING_IDENTITY_LOCK_TIMEOUT_MS}ms`}, true)` ) diff --git a/apps/sim/lib/billing/organizations/create-organization.ts b/apps/sim/lib/billing/organizations/create-organization.ts index 1947a333745..41eaf6f600d 100644 --- a/apps/sim/lib/billing/organizations/create-organization.ts +++ b/apps/sim/lib/billing/organizations/create-organization.ts @@ -3,7 +3,7 @@ import { member, organization } from '@sim/db/schema' import { generateId } from '@sim/utils/id' import { and, eq, ne } from 'drizzle-orm' import { acquireUserBillingIdentityLock } from '@/lib/billing/organizations/billing-identity-lock' -import type { DbOrTx } from '@/lib/db/types' +import type { DbTransaction } from '@/lib/db/types' const ORGANIZATION_SLUG_REGEX = /^[a-z0-9-_]+$/ @@ -80,7 +80,7 @@ export async function createOrganizationWithOwner( * check and the insert across two transactions allows duplicate slugs. */ export async function createOrganizationWithOwnerTx( - tx: DbOrTx, + tx: DbTransaction, { ownerUserId, name, slug, metadata = {} }: CreateOrganizationWithOwnerParams ): Promise { validateOrganizationSlugOrThrow(slug) diff --git a/apps/sim/lib/billing/organizations/membership.ts b/apps/sim/lib/billing/organizations/membership.ts index 676444484d2..ce5181d8a27 100644 --- a/apps/sim/lib/billing/organizations/membership.ts +++ b/apps/sim/lib/billing/organizations/membership.ts @@ -53,7 +53,7 @@ import { enqueueOutboxEvent } from '@/lib/core/outbox/service' import { revokeWorkspaceCredentialMembershipsTx } from '@/lib/credentials/access' import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' import { isRetryableTransactionError } from '@/lib/db/transaction' -import type { DbOrTx } from '@/lib/db/types' +import type { DbOrTx, DbTransaction } from '@/lib/db/types' import { acquireInvitationMutationLocks } from '@/lib/invitations/locks' import { requireMemberManagementAuthority } from '@/lib/organizations/members/authority' import { @@ -78,7 +78,7 @@ export const MEMBER_BILLING_RECONCILIATION_EVENT_TYPE = 'billing.reconcile-membe /** Serializes organization-wide owner, seat, move, and membership decisions. */ export async function acquireOrganizationMutationLock( - tx: DbOrTx, + tx: DbTransaction, organizationId: string ): Promise { await tx.execute( @@ -104,7 +104,7 @@ export async function acquireOrganizationMutationLock( * the wait (it raises SQLSTATE 55P03 instead of hanging) if a holder is stuck. */ export async function acquireOrgMembershipLock( - tx: DbOrTx, + tx: DbTransaction, userId: string, organizationId: string ): Promise { @@ -126,7 +126,7 @@ export async function acquireOrgMembershipLock( * transfer and refuses the insert. */ export async function acquireOrganizationUserMutationLocks( - tx: DbOrTx, + tx: DbTransaction, params: { userId: string; organizationIds: string[] } ): Promise { const organizationIds = [...new Set(params.organizationIds)].sort() @@ -637,7 +637,7 @@ interface MembershipValidationResult { * back together in the caller's transaction. */ export async function ensureUserInOrganizationTx( - tx: DbOrTx, + tx: DbTransaction, params: AddMemberParams ): Promise { const { @@ -883,7 +883,7 @@ async function applyPaidOrgJoinBillingTx( * and the personal-Pro transition. */ export async function reapplyPaidOrgJoinBillingForExistingMemberTx( - tx: DbOrTx, + tx: DbTransaction, userId: string, organizationId: string, options: { sourceOperationId?: string } = {} @@ -966,7 +966,10 @@ export async function withInvitationSafeOrganizationAccessMutation( scope: InvitationRemovalScope additionalOrganizationIds?: string[] }, - operation: (tx: DbOrTx, locked: { workspaceIds: string[]; invitationIds: string[] }) => Promise + operation: ( + tx: DbTransaction, + locked: { workspaceIds: string[]; invitationIds: string[] } + ) => Promise ): Promise { let candidate = await getInvitationRemovalLockSnapshot(db, params) diff --git a/apps/sim/lib/billing/organizations/provision-seat.ts b/apps/sim/lib/billing/organizations/provision-seat.ts index 5d045a5fa5b..5749b34f808 100644 --- a/apps/sim/lib/billing/organizations/provision-seat.ts +++ b/apps/sim/lib/billing/organizations/provision-seat.ts @@ -18,7 +18,7 @@ import { getPlanByName } from '@/lib/billing/plans' import { hasUsableSubscriptionStatus } from '@/lib/billing/subscriptions/utils' import { OUTBOX_EVENT_TYPES } from '@/lib/billing/webhooks/outbox-handlers' import { enqueueOutboxEvent } from '@/lib/core/outbox/service' -import type { DbOrTx } from '@/lib/db/types' +import type { DbOrTx, DbTransaction } from '@/lib/db/types' const logger = createLogger('ProvisionSeat') @@ -49,7 +49,7 @@ interface EnsureTeamOrganizationParams { billingOwnerUserId: string workspaceOrganizationId: string | null /** Transaction that also accepts the invitation and grants permissions. */ - executor: DbOrTx + executor: DbTransaction /** Workspace rows already covered by the caller's invitation/workspace locks. */ workspaceIdsToAttach: string[] } @@ -113,7 +113,7 @@ export async function ensureTeamOrganizationForAcceptance( async function ensureOrganizationOnTeamPlan( organizationId: string, actorId: string, - executor: DbOrTx + executor: DbTransaction ): Promise { await acquireOrganizationMutationLock(executor, organizationId) await assertNoUnresolvedEnterpriseIssuance(executor, organizationId) @@ -153,7 +153,7 @@ async function ensureOrganizationOnTeamPlan( async function convertPersonalSubscriptionToTeam( userId: string, workspaceIdsToAttach: string[], - executor: DbOrTx + executor: DbTransaction ): Promise { const personalSub = await getHighestPriorityPersonalSubscription(userId, { onError: 'throw', diff --git a/apps/sim/lib/credential-groups/enrollments.ts b/apps/sim/lib/credential-groups/enrollments.ts index 11660f08ea8..c2246d2f682 100644 --- a/apps/sim/lib/credential-groups/enrollments.ts +++ b/apps/sim/lib/credential-groups/enrollments.ts @@ -45,7 +45,7 @@ import type { InviteCredentialGroupEnrollmentsInput, } from '@/lib/credential-groups/types' import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' -import type { DbOrTx } from '@/lib/db/types' +import type { DbTransaction } from '@/lib/db/types' import { sendEmail } from '@/lib/messaging/email/mailer' import { getFromEmailAddress } from '@/lib/messaging/email/utils' @@ -174,7 +174,7 @@ export interface CredentialGroupEnrollmentCompletion { /** Serializes OAuth grant persistence and administrative revocation for one enrollment. */ export async function lockCredentialGroupEnrollmentLifecycle( - executor: DbOrTx, + executor: DbTransaction, enrollmentId: string ): Promise { if (!enrollmentId.trim()) throw new Error('Credential group enrollment ID is required') @@ -187,7 +187,7 @@ export async function lockCredentialGroupEnrollmentLifecycle( /** Serializes invitation issuance before an enrollment row is known or locked. */ async function lockCredentialGroupInvitationTarget( - executor: DbOrTx, + executor: DbTransaction, groupId: string, email: string ): Promise { diff --git a/apps/sim/lib/credential-groups/service.ts b/apps/sim/lib/credential-groups/service.ts index 91faed0ca50..e92b6684f6e 100644 --- a/apps/sim/lib/credential-groups/service.ts +++ b/apps/sim/lib/credential-groups/service.ts @@ -41,7 +41,7 @@ import { createWorkspaceAccountsGroup, } from '@/lib/credential-groups/workspace-accounts' import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' -import type { DbOrTx } from '@/lib/db/types' +import type { DbOrTx, DbTransaction } from '@/lib/db/types' type WorkspaceCredentialGroupRecord = CredentialGroupRecord & { workspaceId: string } type OrganizationCredentialGroupRecord = CredentialGroupRecord & { @@ -249,25 +249,25 @@ export function ensureWorkspaceAccountsGroup( scope: Extract, userId: string, option?: CredentialGroupOptionInput, - executor?: DbOrTx + executor?: DbTransaction ): Promise export function ensureWorkspaceAccountsGroup( workspaceId: string, userId: string, option?: CredentialGroupOptionInput, - executor?: DbOrTx + executor?: DbTransaction ): Promise export function ensureWorkspaceAccountsGroup( scope: ResourceScope, userId: string, option?: CredentialGroupOptionInput, - executor?: DbOrTx + executor?: DbTransaction ): Promise export async function ensureWorkspaceAccountsGroup( scopeInput: string | ResourceScope, userId: string, option?: CredentialGroupOptionInput, - executor?: DbOrTx + executor?: DbTransaction ): Promise { const scope = credentialGroupScope(scopeInput) if (option?.provider === 'slack') { @@ -275,7 +275,7 @@ export async function ensureWorkspaceAccountsGroup( } const preparedOption = option ? await buildOption(scope, { ...option, required: false }) : null let wasCreated = false - const provision = async (tx: DbOrTx) => { + const provision = async (tx: DbTransaction) => { await acquireAdvisoryXactLock( tx, 'search_accounts', @@ -400,7 +400,7 @@ export async function addOrganizationAccountProvider( organizationId: string, userId: string, option: { provider: CredentialGroupStandardOAuthProvider; label: string }, - executor: DbOrTx + executor: DbTransaction ): Promise<{ groupId: string; changed: boolean }> { const scope = { kind: 'organization', organizationId } as const const group = await ensureWorkspaceAccountsGroup(scope, userId, undefined, executor) diff --git a/apps/sim/lib/credential-groups/shared-slack-app.ts b/apps/sim/lib/credential-groups/shared-slack-app.ts index ce6fd46bf57..02f9761fd18 100644 --- a/apps/sim/lib/credential-groups/shared-slack-app.ts +++ b/apps/sim/lib/credential-groups/shared-slack-app.ts @@ -9,11 +9,11 @@ import { } from '@/lib/credential-groups/provider-configuration' import { ensureWorkspaceAccountsGroup } from '@/lib/credential-groups/service' import { SLACK_SEARCH_USER_SCOPES } from '@/lib/credential-groups/slack-managed-user-scopes' -import type { DbOrTx } from '@/lib/db/types' +import type { DbTransaction } from '@/lib/db/types' /** Configures personal consent atomically with the authorized admin's bot installation. */ export async function configureSharedSlackMemberApp( - tx: DbOrTx, + tx: DbTransaction, input: { organizationId: string userId: string diff --git a/apps/sim/lib/credentials/env-locks.ts b/apps/sim/lib/credentials/env-locks.ts index d9e9168ba2e..fb8bfb61541 100644 --- a/apps/sim/lib/credentials/env-locks.ts +++ b/apps/sim/lib/credentials/env-locks.ts @@ -1,6 +1,6 @@ import { sql } from 'drizzle-orm' import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' -import type { DbOrTx } from '@/lib/db/types' +import type { DbTransaction } from '@/lib/db/types' const ENV_MAP_LOCK_TIMEOUT_MS = 5_000 @@ -17,17 +17,17 @@ const ENV_MAP_LOCK_TIMEOUT_MS = 5_000 * already takes this lock — a prefixed key would be a different lock and would * serialize against nothing. */ -async function lockEnvMap(tx: DbOrTx, tag: string, lockKey: string): Promise { +async function lockEnvMap(tx: DbTransaction, tag: string, lockKey: string): Promise { await tx.execute(sql`SELECT set_config('lock_timeout', ${`${ENV_MAP_LOCK_TIMEOUT_MS}ms`}, true)`) await acquireAdvisoryXactLock(tx, tag, lockKey) } /** Serializes writers of one workspace's environment variables map. */ -export async function lockWorkspaceEnvMap(tx: DbOrTx, workspaceId: string): Promise { +export async function lockWorkspaceEnvMap(tx: DbTransaction, workspaceId: string): Promise { await lockEnvMap(tx, 'workspace_env_map', workspaceId) } /** Serializes writers of one user's personal environment variables map. */ -export async function lockPersonalEnvMap(tx: DbOrTx, userId: string): Promise { +export async function lockPersonalEnvMap(tx: DbTransaction, userId: string): Promise { await lockEnvMap(tx, 'personal_env_map', userId) } diff --git a/apps/sim/lib/credentials/environment.ts b/apps/sim/lib/credentials/environment.ts index 1840847871c..03c44c20729 100644 --- a/apps/sim/lib/credentials/environment.ts +++ b/apps/sim/lib/credentials/environment.ts @@ -17,7 +17,7 @@ import { and, asc, eq, inArray, isNotNull, isNull, notInArray, or, sql } from 'd import { acquireUserBillingIdentityLock } from '@/lib/billing/organizations/billing-identity-lock' import { isManagedCredentialGroupBindingLive } from '@/lib/credential-groups/credentials' import { lockPersonalEnvMap } from '@/lib/credentials/env-locks' -import type { DbOrTx } from '@/lib/db/types' +import type { DbOrTx, DbTransaction } from '@/lib/db/types' import { getEffectiveWorkspacePermission, hasWorkspaceAdminAccess, @@ -526,11 +526,11 @@ export async function upsertPersonalEnvCredentialForUser(params: { userId: string envKey: string updatedAt: Date - executor?: DbOrTx + executor?: DbTransaction }): Promise { const { userId, envKey, updatedAt } = params - const upsert = async (tx: DbOrTx) => { + const upsert = async (tx: DbTransaction) => { await acquireUserBillingIdentityLock(tx, userId) const workspaceIds = (await getUserWorkspaceIds(userId, tx)).sort() if (workspaceIds.length === 0) return @@ -655,11 +655,11 @@ export async function getPersonalEnvCredentialMetadata(params: { export async function deletePersonalEnvCredentialForUser(params: { userId: string envKey: string - executor?: DbOrTx + executor?: DbTransaction }): Promise { const { userId, envKey } = params - const remove = async (tx: DbOrTx) => { + const remove = async (tx: DbTransaction) => { await acquireUserBillingIdentityLock(tx, userId) await tx .delete(credential) diff --git a/apps/sim/lib/credentials/personal-tokens.ts b/apps/sim/lib/credentials/personal-tokens.ts index 00f300f9751..28983bca774 100644 --- a/apps/sim/lib/credentials/personal-tokens.ts +++ b/apps/sim/lib/credentials/personal-tokens.ts @@ -24,7 +24,7 @@ import { verifyGitLabPersonalToken, } from '@/lib/credentials/gitlab-personal-token' import type { CredentialRow } from '@/lib/credentials/queries' -import type { DbOrTx } from '@/lib/db/types' +import type { DbTransaction } from '@/lib/db/types' import { normalizeGitLabHost } from '@/tools/gitlab/utils' export interface PersonalTokenCredential { @@ -115,11 +115,14 @@ function liveEnrollmentConditions(workspaceId: string, userId: string) { ] } -/** Rechecks the canonical group and the verified person behind a bound token before every use. */ +/** + * Rechecks the canonical group and the verified person behind a bound token before every use. + * Given `lockingTx`, it serializes against the enrollment's lifecycle and holds the binding in + * that transaction. + */ export async function requirePersonalTokenEnrollment( input: ResourceOwner & { userId: string; enrollmentId: string | null }, - executor: DbOrTx = db, - lock = false + lockingTx?: DbTransaction ): Promise<{ credentialGroupId: string }> { const scope = resourceScopeFromOwner(input) if (!input.enrollmentId) @@ -127,7 +130,8 @@ export async function requirePersonalTokenEnrollment( 'forbidden', 'Reconnect your personal account in Connected accounts' ) - if (lock) await lockCredentialGroupEnrollmentLifecycle(executor, input.enrollmentId) + if (lockingTx) await lockCredentialGroupEnrollmentLifecycle(lockingTx, input.enrollmentId) + const executor = lockingTx ?? db const query = executor .select({ id: credentialGroupEnrollment.id, @@ -158,7 +162,7 @@ export async function requirePersonalTokenEnrollment( ) ) .limit(1) - const [binding] = await (lock + const [binding] = await (lockingTx ? query.for('share', { of: [credentialGroupEnrollment, credentialGroup, user] }) : query) if (!binding) @@ -228,8 +232,7 @@ export async function createPersonalTokenCredential(input: CreatePersonalTokenPa userId: input.userId, enrollmentId: enrollment.id, }, - tx, - true + tx ) await tx .update(credentialGroupEnrollment) @@ -369,7 +372,7 @@ export async function updatePersonalTokenCredential(input: UpdatePersonalTokenPa updatedFields.push('apiToken') } const updated = await db.transaction(async (tx) => { - await requirePersonalTokenEnrollment(enrollmentBinding, tx, true) + await requirePersonalTokenEnrollment(enrollmentBinding, tx) const [updated] = await tx .update(credential) .set(updates) diff --git a/apps/sim/lib/db/advisory-locks.ts b/apps/sim/lib/db/advisory-locks.ts index f0fb3c8174b..c5073e011f0 100644 --- a/apps/sim/lib/db/advisory-locks.ts +++ b/apps/sim/lib/db/advisory-locks.ts @@ -1,5 +1,5 @@ import { type SQL, sql } from 'drizzle-orm' -import type { DbOrTx } from '@/lib/db/types' +import type { DbTransaction } from '@/lib/db/types' const LOCK_TAG_PATTERN = /^[a-z][a-z0-9_]*$/ @@ -19,7 +19,11 @@ function lockTag(tag: string): SQL { * Blocks until the transaction-scoped advisory lock for `key` is held. The lock * releases on commit or rollback. */ -export async function acquireAdvisoryXactLock(tx: DbOrTx, tag: string, key: string): Promise { +export async function acquireAdvisoryXactLock( + tx: DbTransaction, + tag: string, + key: string +): Promise { await tx.execute(sql`SELECT pg_advisory_xact_lock(hashtextextended(${key}, 0)) ${lockTag(tag)}`) } @@ -28,7 +32,7 @@ export async function acquireAdvisoryXactLock(tx: DbOrTx, tag: string, key: stri * Returns whether the lock is now held. */ export async function tryAcquireAdvisoryXactLock( - tx: DbOrTx, + tx: DbTransaction, tag: string, key: string ): Promise { diff --git a/apps/sim/lib/db/transaction.ts b/apps/sim/lib/db/transaction.ts index 9e4e2e7f960..f0739bdbc5e 100644 --- a/apps/sim/lib/db/transaction.ts +++ b/apps/sim/lib/db/transaction.ts @@ -3,7 +3,7 @@ import { createLogger } from '@sim/logger' import { getPostgresErrorCode } from '@sim/utils/errors' import { sleep } from '@sim/utils/helpers' import { backoffWithJitter } from '@sim/utils/retry' -import type { DbOrTx } from '@/lib/db/types' +import type { DbTransaction } from '@/lib/db/types' const logger = createLogger('DbTransaction') @@ -48,7 +48,7 @@ export function isRetryableTransactionError(error: unknown): boolean { * than once, and only the committing attempt is durable. */ export async function withTransactionRetry( - fn: (tx: DbOrTx) => Promise, + fn: (tx: DbTransaction) => Promise, options: { attempts?: number; label?: string } = {} ): Promise { const attempts = options.attempts ?? DEFAULT_ATTEMPTS diff --git a/apps/sim/lib/folders/locks.ts b/apps/sim/lib/folders/locks.ts index 3f8e16882ba..e831fc1d9ad 100644 --- a/apps/sim/lib/folders/locks.ts +++ b/apps/sim/lib/folders/locks.ts @@ -2,13 +2,13 @@ import { db } from '@sim/db' import { sql } from 'drizzle-orm' import type { FolderResourceType } from '@/lib/api/contracts/folders' import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' -import type { DbOrTx } from '@/lib/db/types' +import type { DbOrTx, DbTransaction } from '@/lib/db/types' const FOLDER_MUTATION_LOCK_TIMEOUT_MS = 5_000 /** Serializes every writer for one workspace resource-folder tree. */ export async function acquireFolderMutationLock( - tx: DbOrTx, + tx: DbTransaction, workspaceId: string, resourceType: FolderResourceType ): Promise { diff --git a/apps/sim/lib/invitations/core.ts b/apps/sim/lib/invitations/core.ts index 4d7583295a6..4574b5343b9 100644 --- a/apps/sim/lib/invitations/core.ts +++ b/apps/sim/lib/invitations/core.ts @@ -38,7 +38,7 @@ import { hasUsableSubscriptionStatus } from '@/lib/billing/subscriptions/utils' import { ForbiddenOperationError } from '@/lib/core/application/forbidden' import { isBillingEnabled } from '@/lib/core/config/env-flags' import { syncWorkspaceEnvCredentials } from '@/lib/credentials/environment' -import type { DbOrTx } from '@/lib/db/types' +import type { DbOrTx, DbTransaction } from '@/lib/db/types' import { acquireInvitationMutationLocks } from '@/lib/invitations/locks' import { APP_ENTRY_PATH, organizationRoutes } from '@/lib/navigation/paths' import { captureServerEvent } from '@/lib/posthog/server' @@ -98,7 +98,7 @@ export async function getInvitationById( * use the protected state. */ export async function lockInvitationForMutation( - tx: DbOrTx, + tx: DbTransaction, invitationId: string, options?: { lockCurrentGrantWorkspaces?: boolean @@ -921,7 +921,7 @@ async function acceptLockedInvitation( input: AcceptInvitationInput, inv: InvitationWithGrants, lockPlan: InvitationAcceptanceLockPlan, - tx: DbOrTx, + tx: DbTransaction, effects: InvitationAcceptancePostCommitEffects ): Promise { let membershipAlreadyExists = false diff --git a/apps/sim/lib/invitations/locks.ts b/apps/sim/lib/invitations/locks.ts index 80d544c2780..0f8e0a69313 100644 --- a/apps/sim/lib/invitations/locks.ts +++ b/apps/sim/lib/invitations/locks.ts @@ -1,6 +1,6 @@ import { sql } from 'drizzle-orm' import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' -import type { DbOrTx } from '@/lib/db/types' +import type { DbTransaction } from '@/lib/db/types' const INVITATION_MUTATION_LOCK_TIMEOUT_MS = 10_000 @@ -13,7 +13,7 @@ const INVITATION_MUTATION_LOCK_TIMEOUT_MS = 10_000 * workspace mutations. */ export async function acquireInvitationMutationLocks( - tx: DbOrTx, + tx: DbTransaction, params: { invitationIds: string[]; workspaceIds: string[] } ): Promise { await tx.execute( diff --git a/apps/sim/lib/invitations/resend-policy.ts b/apps/sim/lib/invitations/resend-policy.ts index 4eeab9ad831..a6603cb0843 100644 --- a/apps/sim/lib/invitations/resend-policy.ts +++ b/apps/sim/lib/invitations/resend-policy.ts @@ -5,7 +5,7 @@ import { isEnterprise, isTeam } from '@/lib/billing/plan-helpers' import { hasUsableSubscriptionStatus } from '@/lib/billing/subscriptions/utils' import { isBillingEnabled } from '@/lib/core/config/env-flags' import { OrchestrationError } from '@/lib/core/orchestration/types' -import type { DbOrTx } from '@/lib/db/types' +import type { DbTransaction } from '@/lib/db/types' import { type InvitationWithGrants, requireInvitationResendAuthority, @@ -22,7 +22,7 @@ import { validateInvitationsAllowed } from '@/ee/access-control/utils/permission * precede organization and billing-identity locks; permission-group locks are leaves. */ export async function lockInvitationResendPolicy( - tx: DbOrTx, + tx: DbTransaction, invitation: InvitationWithGrants, actorUserId: string, assertedOrganizationId?: string diff --git a/apps/sim/lib/invitations/send.ts b/apps/sim/lib/invitations/send.ts index 1f8cc821802..44da5aae753 100644 --- a/apps/sim/lib/invitations/send.ts +++ b/apps/sim/lib/invitations/send.ts @@ -22,7 +22,7 @@ import { } from '@/components/emails' import { OrchestrationError } from '@/lib/core/orchestration/types' import { getBaseUrl } from '@/lib/core/utils/urls' -import type { DbOrTx } from '@/lib/db/types' +import type { DbOrTx, DbTransaction } from '@/lib/db/types' import { computeInvitationExpiry, lockInvitationForMutation } from '@/lib/invitations/core' import { InvitationNotPendingError } from '@/lib/invitations/errors' import { acquireInvitationMutationLocks } from '@/lib/invitations/locks' @@ -54,7 +54,7 @@ export interface CreatePendingInvitationInput { * and re-authorize stale preflight decisions. */ validateLockedContext?: (context: { - tx: DbOrTx + tx: DbTransaction organizationId: string | null workspaceIds: string[] }) => Promise diff --git a/apps/sim/lib/invitations/workspace-invitations.ts b/apps/sim/lib/invitations/workspace-invitations.ts index d29ccaf7771..efdc042feba 100644 --- a/apps/sim/lib/invitations/workspace-invitations.ts +++ b/apps/sim/lib/invitations/workspace-invitations.ts @@ -20,7 +20,7 @@ import { validateSeatAvailability } from '@/lib/billing/validation/seat-manageme import { isBillingEnabled } from '@/lib/core/config/env-flags' import type { OrchestrationRequestContext } from '@/lib/core/orchestration/types' import { PlatformEvents } from '@/lib/core/telemetry' -import type { DbOrTx } from '@/lib/db/types' +import type { DbTransaction } from '@/lib/db/types' import { DirectGrantContextChangedError, type DirectGrantOutcome, @@ -356,7 +356,7 @@ async function validateLockedWorkspaceInvitationContext({ inviteeEmail, validateLockedWorkspace, }: { - tx: DbOrTx + tx: DbTransaction context: WorkspaceInvitationContext workspaceIds: string[] organizationId: string | null diff --git a/apps/sim/lib/mcp/server-locks.ts b/apps/sim/lib/mcp/server-locks.ts index ff746711025..b07b36e0205 100644 --- a/apps/sim/lib/mcp/server-locks.ts +++ b/apps/sim/lib/mcp/server-locks.ts @@ -1,18 +1,21 @@ import { getPostgresErrorCode } from '@sim/utils/errors' import { sql } from 'drizzle-orm' import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' -import type { DbOrTx } from '@/lib/db/types' +import type { DbTransaction } from '@/lib/db/types' const MCP_SERVER_LOCK_TIMEOUT_MS = 3_000 const LOCK_NOT_AVAILABLE_SQLSTATE = '55P03' -export async function setWorkflowMcpTransactionLockTimeout(tx: DbOrTx): Promise { +export async function setWorkflowMcpTransactionLockTimeout(tx: DbTransaction): Promise { await tx.execute( sql`select set_config('lock_timeout', ${`${MCP_SERVER_LOCK_TIMEOUT_MS}ms`}, true)` ) } -export async function acquireWorkflowMcpServerLock(tx: DbOrTx, serverId: string): Promise { +export async function acquireWorkflowMcpServerLock( + tx: DbTransaction, + serverId: string +): Promise { await setWorkflowMcpTransactionLockTimeout(tx) await acquireAdvisoryXactLock(tx, 'workflow_mcp_server', serverId) } diff --git a/apps/sim/lib/mcp/workflow-mcp-sync.ts b/apps/sim/lib/mcp/workflow-mcp-sync.ts index 6fb0b627643..3e979ef7f16 100644 --- a/apps/sim/lib/mcp/workflow-mcp-sync.ts +++ b/apps/sim/lib/mcp/workflow-mcp-sync.ts @@ -2,7 +2,7 @@ import { db, workflowMcpServer, workflowMcpTool } from '@sim/db' import { createLogger } from '@sim/logger' import { and, asc, desc, eq, gt, inArray, isNotNull, isNull, notExists } from 'drizzle-orm' import { alias } from 'drizzle-orm/pg-core' -import type { DbOrTx } from '@/lib/db/types' +import type { DbOrTx, DbTransaction } from '@/lib/db/types' import { MAX_MCP_SERVERS_PER_WORKFLOW, MAX_MCP_TOOLS_PER_SERVER } from '@/lib/mcp/constants' import { acquireWorkflowMcpServerLock } from '@/lib/mcp/server-locks' import { @@ -268,7 +268,7 @@ async function getRestoreSkipReason( * neither deadlock against each other nor race the checks. */ async function restoreArchivedMcpToolsForWorkflow( - tx: DbOrTx, + tx: DbTransaction, workflowId: string, requestId: string ): Promise { @@ -434,7 +434,7 @@ interface SyncOptionsBase { */ type SyncOptions = SyncOptionsBase & ( - | { tx: DbOrTx; state: { blocks?: Record }; notify?: false } + | { tx: DbTransaction; state: { blocks?: Record }; notify?: false } | { tx?: undefined; state?: { blocks?: Record }; notify?: boolean } ) @@ -607,7 +607,7 @@ export async function syncMcpToolsForWorkflow( export async function removeMcpToolsForWorkflow( workflowId: string, requestId: string, - tx?: DbOrTx, + tx?: DbTransaction, throwOnError = false ): Promise> { if (!tx) { diff --git a/apps/sim/lib/organizations/members/lifecycle.ts b/apps/sim/lib/organizations/members/lifecycle.ts index 9bc6851266f..57ed51cbb85 100644 --- a/apps/sim/lib/organizations/members/lifecycle.ts +++ b/apps/sim/lib/organizations/members/lifecycle.ts @@ -3,7 +3,7 @@ import { createLogger } from '@sim/logger' import { and, eq, isNull } from 'drizzle-orm' import { acquireOrganizationUserMutationLocks } from '@/lib/billing/organizations/membership' import { OrchestrationError } from '@/lib/core/orchestration/types' -import type { DbOrTx } from '@/lib/db/types' +import type { DbOrTx, DbTransaction } from '@/lib/db/types' import { revokeUserSessionsTx } from '@/lib/organizations/members/revocation' const logger = createLogger('OrganizationMemberLifecycle') @@ -38,7 +38,7 @@ export interface SuspendMemberResult { * restores every automation exactly as it was. */ export async function suspendMemberTx( - tx: DbOrTx, + tx: DbTransaction, params: { userId: string; organizationId: string; source: SuspensionSource } ): Promise { await acquireOrganizationUserMutationLocks(tx, { @@ -100,7 +100,7 @@ export type ChangeMemberRoleResult = * moves billing and the last-owner guarantee with it, which is its own operation. */ export async function changeMemberRoleTx( - tx: DbOrTx, + tx: DbTransaction, params: { organizationId: string; userId: string; role: OrganizationMemberRole } ): Promise { await acquireOrganizationUserMutationLocks(tx, { diff --git a/apps/sim/lib/permission-groups/application/group-membership.ts b/apps/sim/lib/permission-groups/application/group-membership.ts index a46b68ce1c1..8d5cf0ea0be 100644 --- a/apps/sim/lib/permission-groups/application/group-membership.ts +++ b/apps/sim/lib/permission-groups/application/group-membership.ts @@ -8,7 +8,7 @@ import { } from '@sim/db/schema' import { generateId } from '@sim/utils/id' import { and, asc, count, eq, inArray, ne, sql } from 'drizzle-orm' -import type { DbOrTx } from '@/lib/db/types' +import type { DbOrTx, DbTransaction } from '@/lib/db/types' import { acquirePermissionGroupOrgLock } from '@/lib/permission-groups/locks' /** @@ -206,7 +206,7 @@ export type AddPermissionGroupMemberResult = 'added' | 'already-member' * further advisory lock may follow it. The membership has no human author. */ export async function addPermissionGroupMemberTx( - tx: DbOrTx, + tx: DbTransaction, params: { organizationId: string; groupId: string; userId: string } ): Promise { await acquirePermissionGroupOrgLock(tx, params.organizationId, { @@ -252,7 +252,7 @@ export type RemovePermissionGroupMemberResult = 'removed' | 'not-a-member' /** Removes a user from a permission group; the caller holds the organization lock. */ export async function removePermissionGroupMemberTx( - tx: DbOrTx, + tx: DbTransaction, params: { organizationId: string; groupId: string; userId: string } ): Promise { await acquirePermissionGroupOrgLock(tx, params.organizationId, { diff --git a/apps/sim/lib/permission-groups/locks.ts b/apps/sim/lib/permission-groups/locks.ts index 6d1c4986531..d52c8c552b5 100644 --- a/apps/sim/lib/permission-groups/locks.ts +++ b/apps/sim/lib/permission-groups/locks.ts @@ -1,6 +1,6 @@ import { sql } from 'drizzle-orm' import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' -import type { DbOrTx } from '@/lib/db/types' +import type { DbTransaction } from '@/lib/db/types' const PERMISSION_GROUP_LOCK_TIMEOUT_MS = 5_000 @@ -46,7 +46,7 @@ const PERMISSION_GROUP_LOCK_TIMEOUT_MS = 5_000 * acquires it, and `lib/` must not import from `app/api/**`. */ export async function acquirePermissionGroupOrgLock( - tx: DbOrTx, + tx: DbTransaction, organizationId: string, options?: { lockTimeoutAlreadyBounded?: boolean } ): Promise { diff --git a/apps/sim/lib/sim-search/live/member-setup.ts b/apps/sim/lib/sim-search/live/member-setup.ts index cd3ecf5790f..9339e7b5fd4 100644 --- a/apps/sim/lib/sim-search/live/member-setup.ts +++ b/apps/sim/lib/sim-search/live/member-setup.ts @@ -7,7 +7,7 @@ import { ManagedMcpConnectorError, } from '@/lib/credential-groups/managed-mcp-service' import { ensureWorkspaceAccountsGroup } from '@/lib/credential-groups/service' -import type { DbOrTx } from '@/lib/db/types' +import type { DbTransaction } from '@/lib/db/types' import type { ManagedSearchMcpProvider } from '@/lib/sim-search/live/managed-mcp-config' /** Joins source approval's transaction, serializing concurrent setup through the accounts lock. */ @@ -15,7 +15,7 @@ export async function addOrganizationSearchMcpProvider( organizationId: string, userId: string, provider: ManagedSearchMcpProvider, - executor: DbOrTx + executor: DbTransaction ): Promise<{ groupId: string; changed: boolean }> { const group = await ensureWorkspaceAccountsGroup( { kind: 'organization', organizationId }, diff --git a/apps/sim/lib/uploads/contexts/workspace/workspace-file-folder-manager.ts b/apps/sim/lib/uploads/contexts/workspace/workspace-file-folder-manager.ts index f021fd7c622..62dd595c92b 100644 --- a/apps/sim/lib/uploads/contexts/workspace/workspace-file-folder-manager.ts +++ b/apps/sim/lib/uploads/contexts/workspace/workspace-file-folder-manager.ts @@ -6,7 +6,7 @@ import { generateId } from '@sim/utils/id' import { and, eq, inArray, isNull, min, sql } from 'drizzle-orm' import { type ListSortOrder, listOrderBy } from '@/lib/api/list-query' import { OrchestrationError } from '@/lib/core/orchestration/types' -import type { DbOrTx } from '@/lib/db/types' +import type { DbOrTx, DbTransaction } from '@/lib/db/types' import { acquireFolderMutationLock } from '@/lib/folders/locks' import { deduplicateFolderName } from '@/lib/folders/naming' import { @@ -233,7 +233,7 @@ export function workspaceFileNameFolderCondition(folderId?: string | null) { return sql`coalesce(${workspaceFiles.folderId}, '') = ${folderId ?? ''}` } -async function acquireWorkspaceFileFolderMutationLock(tx: DbOrTx, workspaceId: string) { +async function acquireWorkspaceFileFolderMutationLock(tx: DbTransaction, workspaceId: string) { await acquireFolderMutationLock(tx, workspaceId, FILE_FOLDER_RESOURCE_TYPE) } diff --git a/apps/sim/lib/workflows/persistence/deployment-operations.ts b/apps/sim/lib/workflows/persistence/deployment-operations.ts index 3d5a95517a6..6bac6ef19fa 100644 --- a/apps/sim/lib/workflows/persistence/deployment-operations.ts +++ b/apps/sim/lib/workflows/persistence/deployment-operations.ts @@ -4,6 +4,7 @@ import type { DbOrTx } from '@sim/workflow-persistence/types' import type { WorkflowState } from '@sim/workflow-types/workflow' import type { InferSelectModel } from 'drizzle-orm' import { and, desc, eq, inArray, or, sql } from 'drizzle-orm' +import type { DbTransaction } from '@/lib/db/types' import { canTransitionDeploymentOperation, createDeploymentReadiness, @@ -95,7 +96,10 @@ export interface MarkDeploymentComponentReadinessParams extends DeploymentOperat } export interface ActivateDeploymentOperationParams extends DeploymentOperationGeneration { - onActivateTransaction?: (tx: DbOrTx, operation: WorkflowDeploymentOperation) => Promise + onActivateTransaction?: ( + tx: DbTransaction, + operation: WorkflowDeploymentOperation + ) => Promise } interface PrepareOperationContext { diff --git a/apps/sim/lib/workspaces/create.ts b/apps/sim/lib/workspaces/create.ts index 8e5d8f46c0b..adae897329a 100644 --- a/apps/sim/lib/workspaces/create.ts +++ b/apps/sim/lib/workspaces/create.ts @@ -4,7 +4,7 @@ import { createLogger } from '@sim/logger' import { getPostgresConstraintName, getPostgresErrorCode } from '@sim/utils/errors' import { generateId } from '@sim/utils/id' import { PlatformEvents } from '@/lib/core/telemetry' -import type { DbOrTx } from '@/lib/db/types' +import type { DbTransaction } from '@/lib/db/types' import { buildDefaultWorkflowArtifacts } from '@/lib/workflows/defaults' import { saveWorkflowToNormalizedTables } from '@/lib/workflows/persistence/utils' import { @@ -88,7 +88,7 @@ export interface TransactionalCreateWorkspaceParams extends CreateWorkspaceParam * permission and optional starter workflow atomically. */ export async function createWorkspaceInTransaction( - tx: DbOrTx, + tx: DbTransaction, { userId, observedOrganizationId, @@ -263,7 +263,7 @@ export async function createWorkspace(params: CreateWorkspaceParams) { * transaction already holds. */ export async function createDefaultPersonalWorkspaceInTransaction( - tx: DbOrTx, + tx: DbTransaction, params: { userId: string; userName: string | null | undefined } ): Promise { const firstName = params.userName?.split(' ')[0] || null diff --git a/apps/sim/lib/workspaces/operations/receipts.ts b/apps/sim/lib/workspaces/operations/receipts.ts index f0ccdd8611b..b88bc9e7923 100644 --- a/apps/sim/lib/workspaces/operations/receipts.ts +++ b/apps/sim/lib/workspaces/operations/receipts.ts @@ -5,7 +5,7 @@ import { sortObjectKeysDeep } from '@sim/utils/object' import { and, eq } from 'drizzle-orm' import { OrchestrationError } from '@/lib/core/orchestration/types' import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' -import type { DbOrTx } from '@/lib/db/types' +import type { DbOrTx, DbTransaction } from '@/lib/db/types' import type { DeploymentOperationStatus } from '@/lib/workflows/deployment-lifecycle' import type { ImportedWorkflowBlock } from '@/lib/workflows/operations/import-workflow' import type { CreateForkResult } from '@/ee/workspace-forking/lib/create-fork' @@ -91,7 +91,7 @@ export function workflowOperationFingerprint(value: unknown): string { /** Serializes absent receipts as well as existing ones without a separately committed claim. */ export async function lockWorkspaceOperationRequest( - tx: DbOrTx, + tx: DbTransaction, workspaceId: string, requestId: string ): Promise { diff --git a/apps/sim/lib/workspaces/organization-workspaces.ts b/apps/sim/lib/workspaces/organization-workspaces.ts index 95324c679f7..8659ab20424 100644 --- a/apps/sim/lib/workspaces/organization-workspaces.ts +++ b/apps/sim/lib/workspaces/organization-workspaces.ts @@ -12,7 +12,7 @@ import { } from '@/lib/billing/organizations/membership' import { changeWorkspaceStoragePayersInTx } from '@/lib/billing/storage/payer-transfer' import { OrchestrationError } from '@/lib/core/orchestration/types' -import type { DbOrTx } from '@/lib/db/types' +import type { DbOrTx, DbTransaction } from '@/lib/db/types' import { acquireInvitationMutationLocks } from '@/lib/invitations/locks' import { invalidateWorkspaceTableLimitsCache } from '@/lib/table/billing' import { getOrganizationOwnerId, WORKSPACE_MODE } from '@/lib/workspaces/policy' @@ -205,7 +205,7 @@ export async function attachOwnedWorkspacesToOrganization({ * transaction. */ export async function attachOwnedWorkspacesToOrganizationTx( - tx: DbOrTx, + tx: DbTransaction, { ownerUserId, organizationId, @@ -424,7 +424,7 @@ export async function detachOrganizationWorkspaces( * describing detachments that a later rollback undid. */ export async function detachOrganizationWorkspacesTx( - tx: DbOrTx, + tx: DbTransaction, organizationId: string ): Promise { const organizationWorkspacesWhere = and( diff --git a/apps/sim/lib/workspaces/policy.ts b/apps/sim/lib/workspaces/policy.ts index e1b97e1ff66..74c3fe02727 100644 --- a/apps/sim/lib/workspaces/policy.ts +++ b/apps/sim/lib/workspaces/policy.ts @@ -14,7 +14,7 @@ import type { PlanCategory } from '@/lib/billing/plan-helpers' import { getPlanType, isEnterprise, isMaxTier, isPro, isTeam } from '@/lib/billing/plan-helpers' import { hasUsableSubscriptionStatus } from '@/lib/billing/subscriptions/utils' import { isBillingEnabled } from '@/lib/core/config/env-flags' -import type { DbOrTx } from '@/lib/db/types' +import type { DbOrTx, DbTransaction } from '@/lib/db/types' import { capabilityDeniedBy, capabilityRefusal, @@ -225,7 +225,7 @@ export async function resolveGoverningPermissionGroupOrganization(params: { * lapse, which is the condition the admin must fix first anyway. */ export async function lockWorkspaceCreationContext( - tx: DbOrTx, + tx: DbTransaction, { userId, organizationId, diff --git a/apps/sim/scripts/migrate-gitlab-personal-tokens.ts b/apps/sim/scripts/migrate-gitlab-personal-tokens.ts index 4ea212c4a25..cd90bb0d45a 100644 --- a/apps/sim/scripts/migrate-gitlab-personal-tokens.ts +++ b/apps/sim/scripts/migrate-gitlab-personal-tokens.ts @@ -24,7 +24,7 @@ import { and, asc, eq, gt, isNull, ne, or, sql } from 'drizzle-orm' import { lockCredentialGroupEnrollmentLifecycle } from '@/lib/credential-groups/enrollments' import { requireOrganizationAccountsSetup } from '@/lib/credential-groups/organization-setup' import { decryptPersonalToken, encryptPersonalToken } from '@/lib/credentials/gitlab-personal-token' -import type { DbOrTx } from '@/lib/db/types' +import type { DbOrTx, DbTransaction } from '@/lib/db/types' const logger = createLogger('MigrateGitLabPersonalTokens') const BATCH_SIZE = 100 @@ -62,7 +62,11 @@ async function assertUniqueIdentities(executor: DbOrTx, organizationId: string) ) } -async function migrateToken(executor: DbOrTx, credentialId: string, options: MigrationOptions) { +async function migrateToken( + executor: DbTransaction, + credentialId: string, + options: MigrationOptions +) { const [initial] = await executor .select() .from(credential) From 56a3764bf1b606eda658515a2785f041b532e449 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Sun, 27 Sep 2026 18:59:10 -0700 Subject: [PATCH 2/2] chore(db): document why the advisory lock helpers take a transaction --- apps/sim/lib/db/advisory-locks.ts | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/apps/sim/lib/db/advisory-locks.ts b/apps/sim/lib/db/advisory-locks.ts index c5073e011f0..608ba70f76e 100644 --- a/apps/sim/lib/db/advisory-locks.ts +++ b/apps/sim/lib/db/advisory-locks.ts @@ -17,7 +17,9 @@ function lockTag(tag: string): SQL { /** * Blocks until the transaction-scoped advisory lock for `key` is held. The lock - * releases on commit or rollback. + * releases on commit or rollback. It takes a transaction, never the pool: on a + * pooled connection the statement autocommits, releasing the lock before the + * caller's work runs. */ export async function acquireAdvisoryXactLock( tx: DbTransaction, @@ -29,7 +31,8 @@ export async function acquireAdvisoryXactLock( /** * Takes the transaction-scoped advisory lock for `key` without waiting. - * Returns whether the lock is now held. + * Returns whether the lock is now held. Like {@link acquireAdvisoryXactLock}, + * it takes a transaction so the lock outlives the statement. */ export async function tryAcquireAdvisoryXactLock( tx: DbTransaction,