diff --git a/apps/sim/app/api/schedules/execute/route.ts b/apps/sim/app/api/schedules/execute/route.ts index fa298ff9d00..8e001717d37 100644 --- a/apps/sim/app/api/schedules/execute/route.ts +++ b/apps/sim/app/api/schedules/execute/route.ts @@ -44,6 +44,7 @@ import { import { runDetached } from '@/lib/core/utils/background' import { generateRequestId } from '@/lib/core/utils/request' import { withRouteHandler } from '@/lib/core/utils/with-route-handler' +import { tryAcquireAdvisoryXactLock } from '@/lib/db/advisory-locks' import type { DbOrTx } from '@/lib/db/types' import { registerManualExecutionAborter, @@ -885,10 +886,12 @@ async function recoverStaleDatabaseScheduleJobs(now: Date): Promise { const disabledScheduleIds = new Set() await db.transaction(async (tx) => { - const [lock] = await tx.execute<{ acquired: boolean }>( - sql`SELECT pg_try_advisory_xact_lock(hashtextextended(${SCHEDULE_EXECUTION_QUEUE_NAME}, 0)) AS acquired` + const acquired = await tryAcquireAdvisoryXactLock( + tx, + 'schedule_execution_queue', + SCHEDULE_EXECUTION_QUEUE_NAME ) - if (!lock?.acquired) { + if (!acquired) { logger.info( 'Skipped stale database schedule job recovery because another worker holds the lock' ) @@ -1065,10 +1068,12 @@ async function tryStartDatabaseScheduleJob(jobId: string): Promise { - const [lock] = await tx.execute<{ acquired: boolean }>( - sql`SELECT pg_try_advisory_xact_lock(hashtextextended(${SCHEDULE_EXECUTION_QUEUE_NAME}, 0)) AS acquired` + const acquired = await tryAcquireAdvisoryXactLock( + tx, + 'schedule_execution_queue', + SCHEDULE_EXECUTION_QUEUE_NAME ) - if (!lock?.acquired) return 'capacity_full' + if (!acquired) return 'capacity_full' const [row] = await tx .select({ diff --git a/apps/sim/app/api/workspaces/[id]/byok-keys/route.ts b/apps/sim/app/api/workspaces/[id]/byok-keys/route.ts index 6cdd60590d7..96438089796 100644 --- a/apps/sim/app/api/workspaces/[id]/byok-keys/route.ts +++ b/apps/sim/app/api/workspaces/[id]/byok-keys/route.ts @@ -30,6 +30,7 @@ import { getSession } from '@/lib/auth' import { encryptSecret } from '@/lib/core/security/encryption' import { generateRequestId } from '@/lib/core/utils/request' import { withRouteHandler } from '@/lib/core/utils/with-route-handler' +import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' import { captureServerEvent } from '@/lib/posthog/server' import { getUserEntityPermissions } from '@/lib/workspaces/permissions/utils' @@ -158,9 +159,7 @@ export const POST = withRouteHandler( await tx.execute( sql`SELECT set_config('lock_timeout', ${`${WORKSPACE_BYOK_LOCK_TIMEOUT_MS}ms`}, true)` ) - await tx.execute( - sql`SELECT pg_advisory_xact_lock(hashtextextended(${`byok:${workspaceId}:${providerId}`}, 0))` - ) + await acquireAdvisoryXactLock(tx, 'workspace_byok', `byok:${workspaceId}:${providerId}`) const [{ keyCount }] = await tx .select({ keyCount: count() }) diff --git a/apps/sim/ee/workspace-forking/lib/lineage/lineage.ts b/apps/sim/ee/workspace-forking/lib/lineage/lineage.ts index c7449df0082..a3471fa0d17 100644 --- a/apps/sim/ee/workspace-forking/lib/lineage/lineage.ts +++ b/apps/sim/ee/workspace-forking/lib/lineage/lineage.ts @@ -1,6 +1,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' export interface ForkLineageNode { @@ -110,9 +111,7 @@ export async function setForkLockTimeout(tx: DbOrTx): Promise { * unnecessary serialization, never a correctness issue. */ export async function acquireForkEdgeLock(tx: DbOrTx, childWorkspaceId: string): Promise { - await tx.execute( - sql`select pg_advisory_xact_lock(hashtextextended(${`fork-edge:${childWorkspaceId}`}, 0))` - ) + await acquireAdvisoryXactLock(tx, 'fork_edge', `fork-edge:${childWorkspaceId}`) } /** @@ -123,7 +122,5 @@ export async function acquireForkEdgeLock(tx: DbOrTx, childWorkspaceId: string): * this BEFORE {@link acquireForkEdgeLock} so the two are taken in a consistent order. */ export async function acquireForkTargetLock(tx: DbOrTx, targetWorkspaceId: string): Promise { - await tx.execute( - sql`select pg_advisory_xact_lock(hashtextextended(${`fork-target:${targetWorkspaceId}`}, 0))` - ) + await acquireAdvisoryXactLock(tx, 'fork_target', `fork-target:${targetWorkspaceId}`) } diff --git a/apps/sim/lib/api-key/application/organization-byok-keys.ts b/apps/sim/lib/api-key/application/organization-byok-keys.ts index c4bff1d162b..63477770349 100644 --- a/apps/sim/lib/api-key/application/organization-byok-keys.ts +++ b/apps/sim/lib/api-key/application/organization-byok-keys.ts @@ -26,6 +26,7 @@ import { authorizeOrganizationOperation } from '@/lib/core/application/organizat import type { OrchestrationRequestContext } from '@/lib/core/orchestration/types' import { OrchestrationError } from '@/lib/core/orchestration/types' import { decryptSecret, encryptSecret } from '@/lib/core/security/encryption' +import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' import { captureServerEvent } from '@/lib/posthog/server' import { loadActiveWorkspaceApplicationContext } from '@/lib/workspaces/application/workspace-context' import type { BYOKProviderId } from '@/tools/types' @@ -335,8 +336,10 @@ export const saveOrganizationByokKey = defineAuthorizedOrganizationByokUseCase({ await tx.execute( sql`SELECT set_config('lock_timeout', ${`${ORGANIZATION_BYOK_LOCK_TIMEOUT_MS}ms`}, true)` ) - await tx.execute( - sql`SELECT pg_advisory_xact_lock(hashtextextended(${`byok:organization:${context.organizationId}:${input.providerId}`}, 0))` + await acquireAdvisoryXactLock( + tx, + 'organization_byok', + `byok:organization:${context.organizationId}:${input.providerId}` ) const [{ keyCount }] = await tx diff --git a/apps/sim/lib/billing/core/usage-log.ts b/apps/sim/lib/billing/core/usage-log.ts index e44e1ba314f..5a845acad67 100644 --- a/apps/sim/lib/billing/core/usage-log.ts +++ b/apps/sim/lib/billing/core/usage-log.ts @@ -27,6 +27,7 @@ import { isOrgScopedSubscription } from '@/lib/billing/subscriptions/utils' import type { InternalUsageLogSource } from '@/lib/billing/usage-sources' import { asOrchestrationError, OrchestrationError } from '@/lib/core/orchestration/types' import { HttpError } from '@/lib/core/utils/http-error' +import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' import type { DbClient, DbOrTx } from '@/lib/db/types' const logger = createLogger('UsageLog') @@ -740,7 +741,7 @@ export async function recordCumulativeUsage( set_config('lock_timeout', ${`${CUMULATIVE_FLUSH_LOCK_TIMEOUT_MS}ms`}, true) `) enterStage('lock') - await tx.execute(sql`select pg_advisory_xact_lock(hashtextextended(${eventKey}, 0))`) + await acquireAdvisoryXactLock(tx, 'usage_log_event', eventKey) enterStage('read') const [existing] = await tx diff --git a/apps/sim/lib/billing/enterprise-owner-claim.ts b/apps/sim/lib/billing/enterprise-owner-claim.ts index 6f407357aa6..3a9cb50ef00 100644 --- a/apps/sim/lib/billing/enterprise-owner-claim.ts +++ b/apps/sim/lib/billing/enterprise-owner-claim.ts @@ -30,6 +30,7 @@ import { processOutboxEventById, } from '@/lib/core/outbox/service' import { getBaseUrl } from '@/lib/core/utils/urls' +import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' import type { DbOrTx } from '@/lib/db/types' import { computeInvitationExpiry, INVITATION_EXPIRY_DAYS } from '@/lib/invitations/expiry' import { MAX_INVITE_EMAILS, MAX_INVITE_WORKSPACES } from '@/lib/invitations/limits' @@ -492,8 +493,10 @@ export async function createEnterpriseOwnerClaim( }) const requestKey = buildClaimRequestKey(input, normalized) const result = await db.transaction(async (tx) => { - await tx.execute( - sql`select pg_advisory_xact_lock(hashtextextended(${`enterprise-owner-claim:${normalized.ownerEmail}`}, 0))` + await acquireAdvisoryXactLock( + tx, + 'enterprise_owner_claim', + `enterprise-owner-claim:${normalized.ownerEmail}` ) const [accountCreatedDuringReview] = await tx .select({ id: user.id }) diff --git a/apps/sim/lib/billing/organizations/billing-identity-lock.ts b/apps/sim/lib/billing/organizations/billing-identity-lock.ts index 381545e1ce7..3be8789aae2 100644 --- a/apps/sim/lib/billing/organizations/billing-identity-lock.ts +++ b/apps/sim/lib/billing/organizations/billing-identity-lock.ts @@ -1,4 +1,5 @@ import { sql } from 'drizzle-orm' +import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' import type { DbOrTx } from '@/lib/db/types' const USER_BILLING_IDENTITY_LOCK_TIMEOUT_MS = 5_000 @@ -12,7 +13,5 @@ export async function acquireUserBillingIdentityLock(tx: DbOrTx, userId: string) await tx.execute( sql`select set_config('lock_timeout', ${`${USER_BILLING_IDENTITY_LOCK_TIMEOUT_MS}ms`}, true)` ) - await tx.execute( - sql`select pg_advisory_xact_lock(hashtextextended(${`user-billing-identity:${userId}`}, 0))` - ) + await acquireAdvisoryXactLock(tx, 'user_billing_identity', `user-billing-identity:${userId}`) } diff --git a/apps/sim/lib/billing/organizations/membership.ts b/apps/sim/lib/billing/organizations/membership.ts index c28c416b8b2..676444484d2 100644 --- a/apps/sim/lib/billing/organizations/membership.ts +++ b/apps/sim/lib/billing/organizations/membership.ts @@ -51,6 +51,7 @@ import { isBillingEnabled } from '@/lib/core/config/env-flags' import { OrchestrationError } from '@/lib/core/orchestration/types' 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 { acquireInvitationMutationLocks } from '@/lib/invitations/locks' @@ -83,8 +84,10 @@ export async function acquireOrganizationMutationLock( await tx.execute( sql`select set_config('lock_timeout', ${`${ORG_MEMBERSHIP_LOCK_TIMEOUT_MS}ms`}, true)` ) - await tx.execute( - sql`select pg_advisory_xact_lock(hashtextextended(${`organization-mutation:${organizationId}`}, 0))` + await acquireAdvisoryXactLock( + tx, + 'organization_mutation', + `organization-mutation:${organizationId}` ) } @@ -108,9 +111,7 @@ export async function acquireOrgMembershipLock( await tx.execute( sql`select set_config('lock_timeout', ${`${ORG_MEMBERSHIP_LOCK_TIMEOUT_MS}ms`}, true)` ) - await tx.execute( - sql`select pg_advisory_xact_lock(hashtextextended(${`${userId}:${organizationId}`}, 0))` - ) + await acquireAdvisoryXactLock(tx, 'organization_membership', `${userId}:${organizationId}`) } /** diff --git a/apps/sim/lib/billing/webhooks/enterprise.ts b/apps/sim/lib/billing/webhooks/enterprise.ts index 3f5d6baaa68..e66807b7afd 100644 --- a/apps/sim/lib/billing/webhooks/enterprise.ts +++ b/apps/sim/lib/billing/webhooks/enterprise.ts @@ -43,6 +43,7 @@ import { enqueueOutboxEvents, patchOutboxEventPayload, } from '@/lib/core/outbox/service' +import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' import { sendEmail } from '@/lib/messaging/email/mailer' import { getFromEmailAddress } from '@/lib/messaging/email/utils' import { captureServerEvent } from '@/lib/posthog/server' @@ -173,8 +174,10 @@ async function reconcileManualEnterpriseSubscription( const coreResult = await db.transaction(async (tx) => { await acquireOrganizationMutationLock(tx, referenceId) - await tx.execute( - sql`select pg_advisory_xact_lock(hashtextextended(${`stripe-subscription:${stripeSubscription.id}`}, 0))` + await acquireAdvisoryXactLock( + tx, + 'stripe_subscription', + `stripe-subscription:${stripeSubscription.id}` ) // The authoritative Stripe read happened under a durable subscription // lease. Fence the write before touching billing state so a crashed holder diff --git a/apps/sim/lib/credential-groups/enrollments.ts b/apps/sim/lib/credential-groups/enrollments.ts index dce3439294c..11660f08ea8 100644 --- a/apps/sim/lib/credential-groups/enrollments.ts +++ b/apps/sim/lib/credential-groups/enrollments.ts @@ -44,6 +44,7 @@ import type { CredentialGroupEnrollmentRecord, InviteCredentialGroupEnrollmentsInput, } from '@/lib/credential-groups/types' +import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' import type { DbOrTx } from '@/lib/db/types' import { sendEmail } from '@/lib/messaging/email/mailer' import { getFromEmailAddress } from '@/lib/messaging/email/utils' @@ -177,8 +178,10 @@ export async function lockCredentialGroupEnrollmentLifecycle( enrollmentId: string ): Promise { if (!enrollmentId.trim()) throw new Error('Credential group enrollment ID is required') - await executor.execute( - sql`SELECT pg_advisory_xact_lock(hashtextextended(${`credential-group-enrollment:${enrollmentId}`}, 0))` + await acquireAdvisoryXactLock( + executor, + 'credential_group_enrollment', + `credential-group-enrollment:${enrollmentId}` ) } @@ -190,8 +193,10 @@ async function lockCredentialGroupInvitationTarget( ): Promise { if (!groupId.trim()) throw new Error('Credential group ID is required') if (!email.trim()) throw new Error('Credential group enrollment email is required') - await executor.execute( - sql`SELECT pg_advisory_xact_lock(hashtextextended(${`credential-group-invitation:${groupId}:${email}`}, 0))` + await acquireAdvisoryXactLock( + executor, + 'credential_group_invitation', + `credential-group-invitation:${groupId}:${email}` ) } diff --git a/apps/sim/lib/credential-groups/oauth.ts b/apps/sim/lib/credential-groups/oauth.ts index a58729ffd93..6f5b50dd512 100644 --- a/apps/sim/lib/credential-groups/oauth.ts +++ b/apps/sim/lib/credential-groups/oauth.ts @@ -4,7 +4,7 @@ import { createLogger } from '@sim/logger' import { sha256Hex } from '@sim/security/hash' import { getErrorMessage } from '@sim/utils/errors' import { generateId } from '@sim/utils/id' -import { and, eq, ne, sql } from 'drizzle-orm' +import { and, eq, ne } from 'drizzle-orm' import { resourceScopeColumns, resourceScopeFields, @@ -41,6 +41,7 @@ import { decryptManagedOAuthTokenSet, encryptManagedOAuthTokenSet, } from '@/lib/credentials/managed-oauth' +import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' function scopesEqual(left: string[], right: string[]): boolean { const normalizedLeft = [...new Set(left)].sort() @@ -160,8 +161,10 @@ async function persistGrant( const completion: CredentialGroupOAuthCompletion = await db.transaction(async (tx) => { if (!context.credentialOwnerId) throw new CredentialGroupInvitationUnavailableError() await lockCredentialGroupEnrollmentLifecycle(tx, context.enrollmentId) - await tx.execute( - sql`SELECT pg_advisory_xact_lock(hashtextextended(${`credential-group-oauth:${context.enrollmentId}:${context.option.id}`}, 0))` + await acquireAdvisoryXactLock( + tx, + 'credential_group_oauth', + `credential-group-oauth:${context.enrollmentId}:${context.option.id}` ) const [enrollment] = await tx .select({ diff --git a/apps/sim/lib/credential-groups/service.ts b/apps/sim/lib/credential-groups/service.ts index ee3b4a889f1..91faed0ca50 100644 --- a/apps/sim/lib/credential-groups/service.ts +++ b/apps/sim/lib/credential-groups/service.ts @@ -10,7 +10,7 @@ import { resourcePolicy, } from '@sim/db/schema' import { generateId } from '@sim/utils/id' -import { and, asc, eq, inArray, isNull, sql } from 'drizzle-orm' +import { and, asc, eq, inArray, isNull } from 'drizzle-orm' import { OrchestrationError } from '@/lib/core/orchestration/types' import { type ResourceScope, @@ -40,6 +40,7 @@ import { createOrganizationAccountsGroup, createWorkspaceAccountsGroup, } from '@/lib/credential-groups/workspace-accounts' +import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' import type { DbOrTx } from '@/lib/db/types' type WorkspaceCredentialGroupRecord = CredentialGroupRecord & { workspaceId: string } @@ -275,8 +276,10 @@ export async function ensureWorkspaceAccountsGroup( const preparedOption = option ? await buildOption(scope, { ...option, required: false }) : null let wasCreated = false const provision = async (tx: DbOrTx) => { - await tx.execute( - sql`SELECT pg_advisory_xact_lock(hashtextextended(${`search-accounts:${resourceScopeKey(scope)}`}, 0))` + await acquireAdvisoryXactLock( + tx, + 'search_accounts', + `search-accounts:${resourceScopeKey(scope)}` ) const [existing] = await tx .select() diff --git a/apps/sim/lib/credential-groups/slack-managed-users.ts b/apps/sim/lib/credential-groups/slack-managed-users.ts index 3ebc8102c0b..ce7982f217f 100644 --- a/apps/sim/lib/credential-groups/slack-managed-users.ts +++ b/apps/sim/lib/credential-groups/slack-managed-users.ts @@ -11,7 +11,7 @@ import { createLogger } from '@sim/logger' import { sha256Hex } from '@sim/security/hash' import { getErrorMessage } from '@sim/utils/errors' import { generateId } from '@sim/utils/id' -import { and, eq, inArray, isNull, or, sql } from 'drizzle-orm' +import { and, eq, inArray, isNull, or } from 'drizzle-orm' import { getRedisClient } from '@/lib/core/config/redis' import { resourceScopeFields, resourceScopeFromOwner } from '@/lib/core/resource-scope' import { resourceScopeCondition } from '@/lib/core/resource-scope.server' @@ -27,6 +27,7 @@ import { SLACK_MANAGED_USER_CONFIGURATION_CALLBACK_PATH, SLACK_SEARCH_USER_SCOPES, } from '@/lib/credential-groups/slack-managed-user-scopes' +import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' import type { DbOrTx } from '@/lib/db/types' import { SLACK_CUSTOM_BOT_PROVIDER_ID, SLACK_CUSTOM_BOT_SECRET_TYPE } from '@/lib/oauth/types' import { resolveSlackAppCredentials } from '@/lib/slack-search/app-configuration' @@ -743,8 +744,10 @@ export async function exchangeAndConfigureSlackManagedUsers(params: { 'invalid_state' ) } - await tx.execute( - sql`SELECT pg_advisory_xact_lock(hashtextextended(${`slack-managed-users:${params.attempt.credentialGroupId}`}, 0))` + await acquireAdvisoryXactLock( + tx, + 'slack_managed_users', + `slack-managed-users:${params.attempt.credentialGroupId}` ) const [group] = await tx .select() diff --git a/apps/sim/lib/credentials/env-locks.ts b/apps/sim/lib/credentials/env-locks.ts index 9bb425ffabf..d9e9168ba2e 100644 --- a/apps/sim/lib/credentials/env-locks.ts +++ b/apps/sim/lib/credentials/env-locks.ts @@ -1,4 +1,5 @@ import { sql } from 'drizzle-orm' +import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' import type { DbOrTx } from '@/lib/db/types' const ENV_MAP_LOCK_TIMEOUT_MS = 5_000 @@ -16,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, lockKey: string): Promise { +async function lockEnvMap(tx: DbOrTx, tag: string, lockKey: string): Promise { await tx.execute(sql`SELECT set_config('lock_timeout', ${`${ENV_MAP_LOCK_TIMEOUT_MS}ms`}, true)`) - await tx.execute(sql`SELECT pg_advisory_xact_lock(hashtextextended(${lockKey}, 0))`) + await acquireAdvisoryXactLock(tx, tag, lockKey) } /** Serializes writers of one workspace's environment variables map. */ export async function lockWorkspaceEnvMap(tx: DbOrTx, workspaceId: string): Promise { - await lockEnvMap(tx, workspaceId) + 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 { - await lockEnvMap(tx, userId) + await lockEnvMap(tx, 'personal_env_map', userId) } diff --git a/apps/sim/lib/credentials/managed-oauth.ts b/apps/sim/lib/credentials/managed-oauth.ts index fb4c6bc5945..51a12431fcb 100644 --- a/apps/sim/lib/credentials/managed-oauth.ts +++ b/apps/sim/lib/credentials/managed-oauth.ts @@ -2,7 +2,7 @@ import { db } from '@sim/db' import { credential, credentialGroupEnrollment } from '@sim/db/schema' import { createLogger } from '@sim/logger' import { getErrorMessage } from '@sim/utils/errors' -import { and, eq, sql } from 'drizzle-orm' +import { and, eq } from 'drizzle-orm' import type { WorkspaceAuthorizationContext } from '@/lib/core/application' import { type ResourceOwner, @@ -22,6 +22,7 @@ import { } from '@/lib/credential-groups/provider-adapter' import { getCredentialGroupProviderAdapterByProviderId } from '@/lib/credential-groups/provider-registry' import { isScopedCredentialGroupsAvailable } from '@/lib/credential-groups/scoped-availability' +import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' import { loadActiveWorkspaceApplicationContext } from '@/lib/workspaces/application/workspace-context' const logger = createLogger('ManagedOAuthCredential') @@ -414,9 +415,7 @@ export async function resolveManagedOAuthToken( } const refreshOutcome = await db.transaction(async (tx) => { - await tx.execute( - sql`SELECT pg_advisory_xact_lock(hashtextextended(${`managed-oauth:${params.credentialId}`}, 0))` - ) + await acquireAdvisoryXactLock(tx, 'managed_oauth', `managed-oauth:${params.credentialId}`) const current = await getManagedCredential(tx, params.credentialId, params) if (!current) { return { diff --git a/apps/sim/lib/db/advisory-locks.integration.ts b/apps/sim/lib/db/advisory-locks.integration.ts new file mode 100644 index 00000000000..0bd6802fed7 --- /dev/null +++ b/apps/sim/lib/db/advisory-locks.integration.ts @@ -0,0 +1,116 @@ +/** Tagged advisory locks must keep blocking semantics on real PostgreSQL and carry their tag to the server. */ +import * as schema from '@sim/db/schema' +import { readTestDatabaseUrl } from '@sim/db/testing/test-infrastructure' +import { getPostgresErrorCode } from '@sim/utils/errors' +import { sleep } from '@sim/utils/helpers' +import { generateId } from '@sim/utils/id' +import { sql } from 'drizzle-orm' +import { drizzle } from 'drizzle-orm/postgres-js' +import postgres from 'postgres' +import { afterAll, describe, expect, it } from 'vitest' +import { acquireAdvisoryXactLock, tryAcquireAdvisoryXactLock } from '@/lib/db/advisory-locks' + +const connection = postgres(readTestDatabaseUrl(), { max: 4, prepare: false, onnotice: () => {} }) +const db = drizzle(connection, { schema }) + +/** + * Holds `key` in an open transaction until the returned release function is + * called. Rejects if the holder transaction fails before taking the lock. + */ +async function holdLock(tag: string, key: string): Promise<() => Promise> { + let release!: () => void + let acquired!: () => void + const released = new Promise((resolve) => { + release = resolve + }) + const held = new Promise((resolve) => { + acquired = resolve + }) + const transaction = db.transaction(async (tx) => { + await acquireAdvisoryXactLock(tx, tag, key) + acquired() + await released + }) + await Promise.race([held, transaction]) + return async () => { + release() + await transaction + } +} + +async function acquireWithTimeout(tag: string, key: string, timeoutMs: number): Promise { + await db.transaction(async (tx) => { + await tx.execute(sql`SELECT set_config('lock_timeout', ${`${timeoutMs}ms`}, true)`) + await acquireAdvisoryXactLock(tx, tag, key) + }) +} + +describe('advisory xact locks', () => { + afterAll(async () => { + await connection.end() + }) + + it('blocks a second transaction on the same key until the holder commits', async () => { + const key = `lock-test:${generateId()}` + const release = await holdLock('lock_test', key) + try { + const error = await acquireWithTimeout('lock_test', key, 200).catch((e: unknown) => e) + expect(getPostgresErrorCode(error)).toBe('55P03') + } finally { + await release() + } + await expect(acquireWithTimeout('lock_test', key, 200)).resolves.toBeUndefined() + }) + + it('does not block a transaction on a different key', async () => { + const release = await holdLock('lock_test', `lock-test:${generateId()}`) + try { + await expect( + acquireWithTimeout('lock_test', `lock-test:${generateId()}`, 200) + ).resolves.toBeUndefined() + } finally { + await release() + } + }) + + it('reports whether the non-blocking variant acquired the lock', async () => { + const key = `lock-test:${generateId()}` + const release = await holdLock('lock_test', key) + let whileHeld: boolean + try { + whileHeld = await db.transaction((tx) => tryAcquireAdvisoryXactLock(tx, 'lock_test', key)) + } finally { + await release() + } + const afterRelease = await db.transaction((tx) => + tryAcquireAdvisoryXactLock(tx, 'lock_test', key) + ) + expect({ whileHeld, afterRelease }).toEqual({ whileHeld: false, afterRelease: true }) + }) + + it('sends the caller tag with the waiting statement', async () => { + const key = `lock-test:${generateId()}` + const release = await holdLock('lock_test', key) + const waiter = acquireWithTimeout('lock_test_waiter', key, 5_000) + let waitingQuery: string | undefined + try { + for (let attempt = 0; attempt < 50 && !waitingQuery; attempt++) { + const [row] = await connection<{ query: string }[]>` + SELECT query FROM pg_stat_activity + WHERE wait_event_type = 'Lock' AND wait_event = 'advisory' AND query LIKE ${'%lock_test_waiter%'}` + waitingQuery = row?.query + if (!waitingQuery) await sleep(20) + } + } finally { + await release() + await waiter + } + expect(waitingQuery).toMatch(/pg_advisory_xact_lock\(.*\) \/\*lock='lock_test_waiter'\*\/$/) + }) + + it('rejects a tag that would escape the SQL comment', async () => { + await expect( + db.transaction((tx) => acquireAdvisoryXactLock(tx, "x'*/; SELECT 1; /*", 'lock-test')) + ).rejects.toThrow('Invalid advisory lock tag') + }) +}) diff --git a/apps/sim/lib/db/advisory-locks.ts b/apps/sim/lib/db/advisory-locks.ts new file mode 100644 index 00000000000..f0fb3c8174b --- /dev/null +++ b/apps/sim/lib/db/advisory-locks.ts @@ -0,0 +1,39 @@ +import { type SQL, sql } from 'drizzle-orm' +import type { DbOrTx } from '@/lib/db/types' + +const LOCK_TAG_PATTERN = /^[a-z][a-z0-9_]*$/ + +/** + * Renders a SQLCommenter tag naming the lock's caller family. Every advisory + * lock shares one normalized statement, so the tag is what lets query insights + * attribute waits to a caller. It trails the statement because tags must + * precede any terminating semicolon, and it is interpolated raw, so the name is + * restricted to a static snake_case identifier. + */ +function lockTag(tag: string): SQL { + if (!LOCK_TAG_PATTERN.test(tag)) throw new Error(`Invalid advisory lock tag: ${tag}`) + return sql.raw(`/*lock='${tag}'*/`) +} + +/** + * 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 { + await tx.execute(sql`SELECT pg_advisory_xact_lock(hashtextextended(${key}, 0)) ${lockTag(tag)}`) +} + +/** + * Takes the transaction-scoped advisory lock for `key` without waiting. + * Returns whether the lock is now held. + */ +export async function tryAcquireAdvisoryXactLock( + tx: DbOrTx, + tag: string, + key: string +): Promise { + const [lock] = await tx.execute<{ acquired: boolean }>( + sql`SELECT pg_try_advisory_xact_lock(hashtextextended(${key}, 0)) AS acquired ${lockTag(tag)}` + ) + return Boolean(lock?.acquired) +} diff --git a/apps/sim/lib/folders/locks.ts b/apps/sim/lib/folders/locks.ts index 4942d3aaa68..3f8e16882ba 100644 --- a/apps/sim/lib/folders/locks.ts +++ b/apps/sim/lib/folders/locks.ts @@ -1,6 +1,7 @@ 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' const FOLDER_MUTATION_LOCK_TIMEOUT_MS = 5_000 @@ -14,8 +15,10 @@ export async function acquireFolderMutationLock( await tx.execute( sql`select set_config('lock_timeout', ${`${FOLDER_MUTATION_LOCK_TIMEOUT_MS}ms`}, true)` ) - await tx.execute( - sql`select pg_advisory_xact_lock(hashtextextended(${`resource_folders:${resourceType}:${workspaceId}`}, 0))` + await acquireAdvisoryXactLock( + tx, + 'resource_folders', + `resource_folders:${resourceType}:${workspaceId}` ) } diff --git a/apps/sim/lib/invitations/locks.ts b/apps/sim/lib/invitations/locks.ts index 3446c24f3dc..80d544c2780 100644 --- a/apps/sim/lib/invitations/locks.ts +++ b/apps/sim/lib/invitations/locks.ts @@ -1,4 +1,5 @@ import { sql } from 'drizzle-orm' +import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' import type { DbOrTx } from '@/lib/db/types' const INVITATION_MUTATION_LOCK_TIMEOUT_MS = 10_000 @@ -27,6 +28,6 @@ export async function acquireInvitationMutationLocks( ].sort() for (const key of keys) { - await tx.execute(sql`select pg_advisory_xact_lock(hashtextextended(${key}, 0))`) + await acquireAdvisoryXactLock(tx, 'invitation_mutation', key) } } diff --git a/apps/sim/lib/knowledge/application/github-installations.ts b/apps/sim/lib/knowledge/application/github-installations.ts index 82c4a035bf9..3766cd602c9 100644 --- a/apps/sim/lib/knowledge/application/github-installations.ts +++ b/apps/sim/lib/knowledge/application/github-installations.ts @@ -17,6 +17,7 @@ import { ManagedOAuthCredentialError, resolveManagedOAuthToken, } from '@/lib/credentials/managed-oauth' +import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' import type { DbOrTx } from '@/lib/db/types' import { requireOrganizationSearchAvailable } from '@/lib/knowledge/access/availability' import { defineAuthorizedKnowledgeUseCase } from '@/lib/knowledge/application/authorized-knowledge-use-case' @@ -172,8 +173,10 @@ export const connectGitHubSearchInstallation = defineAuthorizedKnowledgeUseCase( }) const { encrypted } = await encryptSecret(JSON.stringify(binding)) return db.transaction(async (tx) => { - await tx.execute( - sql`SELECT pg_advisory_xact_lock(hashtextextended(${`github-search:${context.organizationId}:${binding.installationId}`}, 0))` + await acquireAdvisoryXactLock( + tx, + 'github_search', + `github-search:${context.organizationId}:${binding.installationId}` ) const [admin] = await tx .select({ id: member.id }) diff --git a/apps/sim/lib/knowledge/application/slack-search/installations.ts b/apps/sim/lib/knowledge/application/slack-search/installations.ts index fd1fbd17ecb..744f1ba2921 100644 --- a/apps/sim/lib/knowledge/application/slack-search/installations.ts +++ b/apps/sim/lib/knowledge/application/slack-search/installations.ts @@ -2,8 +2,9 @@ import { AuditAction, AuditResourceType } from '@sim/audit' import { db } from '@sim/db' import { credential, slackApp, slackSearchInstallation } from '@sim/db/schema' import { generateId } from '@sim/utils/id' -import { and, eq, ne, sql } from 'drizzle-orm' +import { and, eq, ne } from 'drizzle-orm' import { OrchestrationError } from '@/lib/core/orchestration/types' +import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' import { SlackSearchConfigurationError, SlackSearchProviderError, @@ -134,9 +135,7 @@ export const configureSlackSearchInstallation = defineAuthorizedKnowledgeUseCase } return db.transaction(async (tx) => { if (identity) - await tx.execute( - sql`SELECT pg_advisory_xact_lock(hashtextextended(${`slack-search:${identity.teamId}`}, 0))` - ) + await acquireAdvisoryXactLock(tx, 'slack_search', `slack-search:${identity.teamId}`) const [current] = await tx .select() .from(credential) diff --git a/apps/sim/lib/knowledge/application/slack-search/setup.ts b/apps/sim/lib/knowledge/application/slack-search/setup.ts index 57890df1601..cc6a1dac795 100644 --- a/apps/sim/lib/knowledge/application/slack-search/setup.ts +++ b/apps/sim/lib/knowledge/application/slack-search/setup.ts @@ -3,7 +3,7 @@ import type { SessionPrincipal } from '@sim/auth/principal' import { db } from '@sim/db' import { credential, slackApp, slackSearchInstallation } from '@sim/db/schema' import { generateId } from '@sim/utils/id' -import { and, eq, ne, sql } from 'drizzle-orm' +import { and, eq, ne } from 'drizzle-orm' import { authorizeOrganizationOperation } from '@/lib/core/application/organization-authorization' import { OrchestrationError } from '@/lib/core/orchestration/types' import { decryptSecret, encryptSecret } from '@/lib/core/security/encryption' @@ -13,6 +13,7 @@ import { loadOrganizationSlackMemberApps, } from '@/lib/credential-groups/organization-slack-app' import { configureSharedSlackMemberApp } from '@/lib/credential-groups/shared-slack-app' +import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' import type { DbOrTx } from '@/lib/db/types' import { buildSlackAppCreationUrl } from '@/lib/integrations/slack-manifest' import { @@ -232,12 +233,8 @@ async function revokeUninstalledSharedGrant( ) { try { await db.transaction(async (tx) => { - await tx.execute( - sql`SELECT pg_advisory_xact_lock(hashtextextended(${`slack-search:${grant.team.id}`}, 0))` - ) - await tx.execute( - sql`SELECT pg_advisory_xact_lock(hashtextextended(${`slack-app:${grant.app_id}`}, 0))` - ) + await acquireAdvisoryXactLock(tx, 'slack_search', `slack-search:${grant.team.id}`) + await acquireAdvisoryXactLock(tx, 'slack_app', `slack-app:${grant.app_id}`) const [installation] = await tx .select({ id: slackSearchInstallation.id }) .from(slackSearchInstallation) @@ -298,12 +295,8 @@ async function saveInstallation( ) await db.transaction(async (tx) => { /** Serialize app/workspace installs before checking ownership or inserting missing rows. */ - await tx.execute( - sql`SELECT pg_advisory_xact_lock(hashtextextended(${`slack-search:${identity.teamId}`}, 0))` - ) - await tx.execute( - sql`SELECT pg_advisory_xact_lock(hashtextextended(${`slack-app:${identity.appId}`}, 0))` - ) + await acquireAdvisoryXactLock(tx, 'slack_search', `slack-search:${identity.teamId}`) + await acquireAdvisoryXactLock(tx, 'slack_app', `slack-app:${identity.appId}`) const [existingApp] = await tx .select() .from(slackApp) diff --git a/apps/sim/lib/knowledge/chunks/service.ts b/apps/sim/lib/knowledge/chunks/service.ts index f471c084f81..b276f06398a 100644 --- a/apps/sim/lib/knowledge/chunks/service.ts +++ b/apps/sim/lib/knowledge/chunks/service.ts @@ -14,6 +14,7 @@ import { searchFilter, textKey, } from '@/lib/api/list-query' +import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' import type { DurableSecretProvenance } from '@/lib/execution/durable-secret-provenance' import { knowledgeAccessCondition } from '@/lib/knowledge/access/predicate' import type { KnowledgeAccessScope } from '@/lib/knowledge/access/types' @@ -239,9 +240,7 @@ export async function createChunk( await tx.execute( sql`select set_config('lock_timeout', ${`${KB_CHUNK_LOCK_TIMEOUT_MS}ms`}, true)` ) - await tx.execute( - sql`select pg_advisory_xact_lock(hashtextextended(${`kb_chunk_seq:${documentId}`}, 0))` - ) + await acquireAdvisoryXactLock(tx, 'kb_chunk_seq', `kb_chunk_seq:${documentId}`) const activeDocument = await tx .select({ id: document.id }) diff --git a/apps/sim/lib/logs/execution/logger.ts b/apps/sim/lib/logs/execution/logger.ts index a63334e1ecd..9d1deb2753a 100644 --- a/apps/sim/lib/logs/execution/logger.ts +++ b/apps/sim/lib/logs/execution/logger.ts @@ -33,6 +33,7 @@ import { checkAndBillPayerOverageThreshold } from '@/lib/billing/threshold-billi import { isBillingEnabled } from '@/lib/core/config/env-flags' import { redactApiKeys } from '@/lib/core/security/redaction' import { filterForDisplay } from '@/lib/core/utils/display-filters' +import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' import { collectLargeValueReferenceKeys, replaceLargeValueReferenceKeysWithClient, @@ -1786,7 +1787,7 @@ export class ExecutionLogger implements IExecutionLoggerService { await tx.execute( sql`select set_config('lock_timeout', ${`${USAGE_RECONCILE_LOCK_TIMEOUT_MS}ms`}, true)` ) - await tx.execute(sql`select pg_advisory_xact_lock(hashtextextended(${executionId}, 0))`) + await acquireAdvisoryXactLock(tx, 'execution_usage_reconcile', executionId) // Already-billed for this execution, scoped to the rows this path owns // (source='workflow') so a same-executionId row from another source diff --git a/apps/sim/lib/mcp/server-locks.ts b/apps/sim/lib/mcp/server-locks.ts index 02811294699..ff746711025 100644 --- a/apps/sim/lib/mcp/server-locks.ts +++ b/apps/sim/lib/mcp/server-locks.ts @@ -1,5 +1,6 @@ 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' const MCP_SERVER_LOCK_TIMEOUT_MS = 3_000 @@ -13,7 +14,7 @@ export async function setWorkflowMcpTransactionLockTimeout(tx: DbOrTx): Promise< export async function acquireWorkflowMcpServerLock(tx: DbOrTx, serverId: string): Promise { await setWorkflowMcpTransactionLockTimeout(tx) - await tx.execute(sql`select pg_advisory_xact_lock(hashtextextended(${serverId}, 0))`) + await acquireAdvisoryXactLock(tx, 'workflow_mcp_server', serverId) } export function isWorkflowMcpServerLockTimeout(error: unknown): boolean { diff --git a/apps/sim/lib/memory/locks.ts b/apps/sim/lib/memory/locks.ts index d8a61dec875..da56b8cb9ac 100644 --- a/apps/sim/lib/memory/locks.ts +++ b/apps/sim/lib/memory/locks.ts @@ -1,4 +1,4 @@ -import { sql } from 'drizzle-orm' +import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' import type { DbTransaction } from '@/lib/db/types' /** Serializes first writes too, before an absent conversation has a row that can be locked. */ @@ -8,5 +8,5 @@ export async function lockMemoryConversationInTx( key: string ): Promise { const lockKey = JSON.stringify(['memory', workspaceId, key]) - await tx.execute(sql`SELECT pg_advisory_xact_lock(hashtextextended(${lockKey}, 0))`) + await acquireAdvisoryXactLock(tx, 'memory_conversation', lockKey) } diff --git a/apps/sim/lib/mothership/async-runs/repository.ts b/apps/sim/lib/mothership/async-runs/repository.ts index fae022827de..28cc32397f9 100644 --- a/apps/sim/lib/mothership/async-runs/repository.ts +++ b/apps/sim/lib/mothership/async-runs/repository.ts @@ -27,6 +27,7 @@ import { sql, } from 'drizzle-orm' import { type ResourceOwner, resourceScopeFromOwner } from '@/lib/core/resource-scope' +import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' import type { SessionProcessIdentity } from '@/lib/execution/remote-sandbox/session-process' import { AsyncToolCallOwnershipError } from '@/lib/mothership/async-runs/errors' import { @@ -132,7 +133,7 @@ export async function withRunAdmissionLock( return db.transaction(async (tx) => { await tx.execute(sql`SET LOCAL statement_timeout = '10s'`) const key = JSON.stringify(['copilot-run-admission', userId, streamId]) - await tx.execute(sql`SELECT pg_advisory_xact_lock(hashtextextended(${key}, 0))`) + await acquireAdvisoryXactLock(tx, 'copilot_run_admission', key) return action(tx) }) } diff --git a/apps/sim/lib/mothership/chat/delegation.ts b/apps/sim/lib/mothership/chat/delegation.ts index 95f69ef2571..928353270a0 100644 --- a/apps/sim/lib/mothership/chat/delegation.ts +++ b/apps/sim/lib/mothership/chat/delegation.ts @@ -1,9 +1,10 @@ import { apiKey } from '@sim/db/schema' import { createLogger } from '@sim/logger' import { generateShortId } from '@sim/utils/id' -import { and, eq, gt, sql } from 'drizzle-orm' +import { and, eq, gt } from 'drizzle-orm' import { createApiKey } from '@/lib/api-key/auth' import { decryptApiKey, hashApiKey } from '@/lib/api-key/crypto' +import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' import { traceMothershipQuery, traceMothershipTransaction, @@ -36,7 +37,7 @@ export async function mintDelegationToken(params: { return await traceMothershipTransaction('delegation', async (tx) => { const name = delegationName(params.userId) await traceMothershipQuery('advisory_lock', 'api_key', () => - tx.execute(sql`SELECT pg_advisory_xact_lock(hashtextextended(${name}, 0))`) + acquireAdvisoryXactLock(tx, 'mothership_delegation', name) ) const [existing] = await traceMothershipQuery('SELECT', 'api_key', () => tx diff --git a/apps/sim/lib/organizations/instance-org.ts b/apps/sim/lib/organizations/instance-org.ts index c9c309dc7d6..92099f42745 100644 --- a/apps/sim/lib/organizations/instance-org.ts +++ b/apps/sim/lib/organizations/instance-org.ts @@ -27,6 +27,7 @@ import { } from '@/lib/billing/organizations/create-organization' import { env } from '@/lib/core/config/env' import { isBillingEnabled } from '@/lib/core/config/env-flags' +import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' import type { DbOrTx } from '@/lib/db/types' const logger = createLogger('InstanceOrganization') @@ -203,8 +204,10 @@ export async function ensureInstanceOrganization( await tx.execute( sql`select set_config('lock_timeout', ${`${INSTANCE_ORG_LOCK_TIMEOUT_MS}ms`}, true)` ) - await tx.execute( - sql`select pg_advisory_xact_lock(hashtextextended(${`instance-organization:${config.slug}`}, 0))` + await acquireAdvisoryXactLock( + tx, + 'instance_organization', + `instance-organization:${config.slug}` ) const resolved = await resolveInstanceOrganizationBySlug(tx, config.slug) diff --git a/apps/sim/lib/permission-groups/locks.ts b/apps/sim/lib/permission-groups/locks.ts index c1cb101f9f5..6d1c4986531 100644 --- a/apps/sim/lib/permission-groups/locks.ts +++ b/apps/sim/lib/permission-groups/locks.ts @@ -1,4 +1,5 @@ import { sql } from 'drizzle-orm' +import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' import type { DbOrTx } from '@/lib/db/types' const PERMISSION_GROUP_LOCK_TIMEOUT_MS = 5_000 @@ -54,7 +55,5 @@ export async function acquirePermissionGroupOrgLock( sql`select set_config('lock_timeout', ${`${PERMISSION_GROUP_LOCK_TIMEOUT_MS}ms`}, true)` ) } - await tx.execute( - sql`select pg_advisory_xact_lock(hashtextextended(${`permission_group:${organizationId}`}, 0))` - ) + await acquireAdvisoryXactLock(tx, 'permission_group', `permission_group:${organizationId}`) } diff --git a/apps/sim/lib/table/rows/ordering.ts b/apps/sim/lib/table/rows/ordering.ts index c0c1dc08d3f..86f8e8e3a62 100644 --- a/apps/sim/lib/table/rows/ordering.ts +++ b/apps/sim/lib/table/rows/ordering.ts @@ -10,6 +10,7 @@ import { db } from '@sim/db' import { userTableRows } from '@sim/db/schema' import { createLogger } from '@sim/logger' import { and, asc, desc, eq, gt, inArray, lt, lte, type SQL, sql } from 'drizzle-orm' +import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' import type { DbOrTx } from '@/lib/db/types' import { getDeleteSnapshotBatchSize, TABLE_LIMITS } from '@/lib/table/constants' import type { MutationProof } from '@/lib/table/mutation-locks' @@ -156,9 +157,7 @@ export async function nextImportStartOrderKey(tableId: string): Promise { if (!revalidate) return undefined - await trx.execute( - sql`SELECT pg_advisory_xact_lock(hashtextextended(${`user_table_schema:${tableId}`}, 0))` - ) + await acquireAdvisoryXactLock(trx, 'user_table_schema', `user_table_schema:${tableId}`) return revalidate(trx) } diff --git a/apps/sim/lib/table/service.ts b/apps/sim/lib/table/service.ts index 87b6de8e544..0f749b86c07 100644 --- a/apps/sim/lib/table/service.ts +++ b/apps/sim/lib/table/service.ts @@ -31,6 +31,7 @@ import { import { ForbiddenOperationError } from '@/lib/core/application/forbidden' import { OrchestrationError } from '@/lib/core/orchestration/types' import { generateRestoreName } from '@/lib/core/utils/restore-name' +import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' import type { DbOrTx } from '@/lib/db/types' import { resolveRestoredFolderId } from '@/lib/folders/queries' import { notifyWorkspaceTablesChanged } from '@/lib/realtime/notify' @@ -135,9 +136,7 @@ export async function withLockedTable( ): Promise { return db.transaction(async (trx) => { await setTableTxTimeouts(trx) - await trx.execute( - sql`SELECT pg_advisory_xact_lock(hashtextextended(${`user_table_schema:${tableId}`}, 0))` - ) + await acquireAdvisoryXactLock(trx, 'user_table_schema', `user_table_schema:${tableId}`) const table = await getTableById(tableId, { tx: trx, includeArchived: opts?.includeArchived }) if (!table || (opts?.expectedWorkspaceId && table.workspaceId !== opts.expectedWorkspaceId)) { throw new OrchestrationError('not_found', 'Table not found') diff --git a/apps/sim/lib/table/views/service.ts b/apps/sim/lib/table/views/service.ts index e538454adae..550cb617713 100644 --- a/apps/sim/lib/table/views/service.ts +++ b/apps/sim/lib/table/views/service.ts @@ -15,6 +15,7 @@ import { tableViews } from '@sim/db/schema' import { createLogger } from '@sim/logger' import { generateId } from '@sim/utils/id' import { and, asc, count, eq, ne, sql } from 'drizzle-orm' +import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' import { buildColumnIdByName, getColumnId, @@ -404,9 +405,7 @@ async function withTableViewsLock( ): Promise { return db.transaction(async (trx) => { await setTableTxTimeouts(trx) - await trx.execute( - sql`SELECT pg_advisory_xact_lock(hashtextextended(${`user_table_views:${tableId}`}, 0))` - ) + await acquireAdvisoryXactLock(trx, 'user_table_views', `user_table_views:${tableId}`) return write(trx) }) } diff --git a/apps/sim/lib/workspace-files/search/dispatcher.ts b/apps/sim/lib/workspace-files/search/dispatcher.ts index 52b5fcd4e85..377380f0ed0 100644 --- a/apps/sim/lib/workspace-files/search/dispatcher.ts +++ b/apps/sim/lib/workspace-files/search/dispatcher.ts @@ -30,6 +30,7 @@ import { import { isTriggerDevEnabled } from '@/lib/core/config/env-flags' import { isInsideTriggerRun } from '@/lib/core/config/trigger-runtime' import { runDetached } from '@/lib/core/utils/background' +import { tryAcquireAdvisoryXactLock } from '@/lib/db/advisory-locks' import type { DbTransaction } from '@/lib/db/types' import { FILE_SEARCH_BACKFILL_PAGE_SIZE, @@ -461,10 +462,12 @@ export async function prepareWorkspaceFileSearchDispatch( }) ) return runDispatchPhase('prepare', async () => { - const [lock] = await tx.execute<{ acquired: boolean }>( - sql`SELECT pg_try_advisory_xact_lock(hashtextextended(${DISPATCH_LOCK_NAME}, 0)) AS acquired` + const acquired = await tryAcquireAdvisoryXactLock( + tx, + 'workspace_file_search_dispatch', + DISPATCH_LOCK_NAME ) - if (!lock?.acquired) { + if (!acquired) { return { payloads: [], backfilledFiles: 0, diff --git a/apps/sim/lib/workspaces/operations/receipts.ts b/apps/sim/lib/workspaces/operations/receipts.ts index 9bb74418ab9..f0ccdd8611b 100644 --- a/apps/sim/lib/workspaces/operations/receipts.ts +++ b/apps/sim/lib/workspaces/operations/receipts.ts @@ -2,8 +2,9 @@ import { createHash } from 'node:crypto' import { db } from '@sim/db' import { workspaceOperationReceipt } from '@sim/db/schema' import { sortObjectKeysDeep } from '@sim/utils/object' -import { and, eq, sql } from 'drizzle-orm' +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 { DeploymentOperationStatus } from '@/lib/workflows/deployment-lifecycle' import type { ImportedWorkflowBlock } from '@/lib/workflows/operations/import-workflow' @@ -94,8 +95,10 @@ export async function lockWorkspaceOperationRequest( workspaceId: string, requestId: string ): Promise { - await tx.execute( - sql`SELECT pg_advisory_xact_lock(hashtextextended(${JSON.stringify(['workspace-operation', workspaceId, requestId])}, 0))` + await acquireAdvisoryXactLock( + tx, + 'workspace_operation', + JSON.stringify(['workspace-operation', workspaceId, requestId]) ) } diff --git a/apps/sim/scripts/register-platform-slack-app.ts b/apps/sim/scripts/register-platform-slack-app.ts index 185a7b9aa8a..2eff959d7a8 100644 --- a/apps/sim/scripts/register-platform-slack-app.ts +++ b/apps/sim/scripts/register-platform-slack-app.ts @@ -2,8 +2,9 @@ import { db } from '@sim/db' import { slackApp } from '@sim/db/schema' import { createLogger } from '@sim/logger' import { generateId } from '@sim/utils/id' -import { eq, sql } from 'drizzle-orm' +import { eq } from 'drizzle-orm' import { encryptSecret } from '@/lib/core/security/encryption' +import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks' const logger = createLogger('RegisterPlatformSlackApp') @@ -26,9 +27,7 @@ async function main() { encryptSecret(signingSecret), ]) await db.transaction(async (tx) => { - await tx.execute( - sql`SELECT pg_advisory_xact_lock(hashtextextended(${`slack-app:${appId}`}, 0))` - ) + await acquireAdvisoryXactLock(tx, 'slack_app', `slack-app:${appId}`) const [existing] = await tx .select({ kind: slackApp.kind }) .from(slackApp)