Skip to content

Commit 3984171

Browse files
committed
improvement(db): tag advisory lock statements by caller
1 parent a0c93d6 commit 3984171

35 files changed

Lines changed: 279 additions & 109 deletions

File tree

‎apps/sim/app/api/schedules/execute/route.ts‎

Lines changed: 11 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -44,6 +44,7 @@ import {
4444
import { runDetached } from '@/lib/core/utils/background'
4545
import { generateRequestId } from '@/lib/core/utils/request'
4646
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
47+
import { tryAcquireAdvisoryXactLock } from '@/lib/db/advisory-locks'
4748
import type { DbOrTx } from '@/lib/db/types'
4849
import {
4950
registerManualExecutionAborter,
@@ -885,10 +886,12 @@ async function recoverStaleDatabaseScheduleJobs(now: Date): Promise<void> {
885886
const disabledScheduleIds = new Set<string>()
886887

887888
await db.transaction(async (tx) => {
888-
const [lock] = await tx.execute<{ acquired: boolean }>(
889-
sql`SELECT pg_try_advisory_xact_lock(hashtextextended(${SCHEDULE_EXECUTION_QUEUE_NAME}, 0)) AS acquired`
889+
const acquired = await tryAcquireAdvisoryXactLock(
890+
tx,
891+
'schedule_execution_queue',
892+
SCHEDULE_EXECUTION_QUEUE_NAME
890893
)
891-
if (!lock?.acquired) {
894+
if (!acquired) {
892895
logger.info(
893896
'Skipped stale database schedule job recovery because another worker holds the lock'
894897
)
@@ -1065,10 +1068,12 @@ async function tryStartDatabaseScheduleJob(jobId: string): Promise<DatabaseSched
10651068
const now = new Date()
10661069

10671070
return db.transaction(async (tx) => {
1068-
const [lock] = await tx.execute<{ acquired: boolean }>(
1069-
sql`SELECT pg_try_advisory_xact_lock(hashtextextended(${SCHEDULE_EXECUTION_QUEUE_NAME}, 0)) AS acquired`
1071+
const acquired = await tryAcquireAdvisoryXactLock(
1072+
tx,
1073+
'schedule_execution_queue',
1074+
SCHEDULE_EXECUTION_QUEUE_NAME
10701075
)
1071-
if (!lock?.acquired) return 'capacity_full'
1076+
if (!acquired) return 'capacity_full'
10721077

10731078
const [row] = await tx
10741079
.select({

‎apps/sim/app/api/workspaces/[id]/byok-keys/route.ts‎

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,7 @@ import { getSession } from '@/lib/auth'
3030
import { encryptSecret } from '@/lib/core/security/encryption'
3131
import { generateRequestId } from '@/lib/core/utils/request'
3232
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
33+
import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks'
3334
import { captureServerEvent } from '@/lib/posthog/server'
3435
import { getUserEntityPermissions } from '@/lib/workspaces/permissions/utils'
3536

@@ -158,9 +159,7 @@ export const POST = withRouteHandler(
158159
await tx.execute(
159160
sql`SELECT set_config('lock_timeout', ${`${WORKSPACE_BYOK_LOCK_TIMEOUT_MS}ms`}, true)`
160161
)
161-
await tx.execute(
162-
sql`SELECT pg_advisory_xact_lock(hashtextextended(${`byok:${workspaceId}:${providerId}`}, 0))`
163-
)
162+
await acquireAdvisoryXactLock(tx, 'workspace_byok', `byok:${workspaceId}:${providerId}`)
164163

165164
const [{ keyCount }] = await tx
166165
.select({ keyCount: count() })

‎apps/sim/ee/workspace-forking/lib/lineage/lineage.ts‎

Lines changed: 3 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
import { db } from '@sim/db'
22
import { workspace } from '@sim/db/schema'
33
import { and, desc, eq, isNull, sql } from 'drizzle-orm'
4+
import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks'
45
import type { DbOrTx } from '@/lib/db/types'
56

67
export interface ForkLineageNode {
@@ -110,9 +111,7 @@ export async function setForkLockTimeout(tx: DbOrTx): Promise<void> {
110111
* unnecessary serialization, never a correctness issue.
111112
*/
112113
export async function acquireForkEdgeLock(tx: DbOrTx, childWorkspaceId: string): Promise<void> {
113-
await tx.execute(
114-
sql`select pg_advisory_xact_lock(hashtextextended(${`fork-edge:${childWorkspaceId}`}, 0))`
115-
)
114+
await acquireAdvisoryXactLock(tx, 'fork_edge', `fork-edge:${childWorkspaceId}`)
116115
}
117116

118117
/**
@@ -123,7 +122,5 @@ export async function acquireForkEdgeLock(tx: DbOrTx, childWorkspaceId: string):
123122
* this BEFORE {@link acquireForkEdgeLock} so the two are taken in a consistent order.
124123
*/
125124
export async function acquireForkTargetLock(tx: DbOrTx, targetWorkspaceId: string): Promise<void> {
126-
await tx.execute(
127-
sql`select pg_advisory_xact_lock(hashtextextended(${`fork-target:${targetWorkspaceId}`}, 0))`
128-
)
125+
await acquireAdvisoryXactLock(tx, 'fork_target', `fork-target:${targetWorkspaceId}`)
129126
}

‎apps/sim/lib/api-key/application/organization-byok-keys.ts‎

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@ import { authorizeOrganizationOperation } from '@/lib/core/application/organizat
2626
import type { OrchestrationRequestContext } from '@/lib/core/orchestration/types'
2727
import { OrchestrationError } from '@/lib/core/orchestration/types'
2828
import { decryptSecret, encryptSecret } from '@/lib/core/security/encryption'
29+
import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks'
2930
import { captureServerEvent } from '@/lib/posthog/server'
3031
import { loadActiveWorkspaceApplicationContext } from '@/lib/workspaces/application/workspace-context'
3132
import type { BYOKProviderId } from '@/tools/types'
@@ -335,8 +336,10 @@ export const saveOrganizationByokKey = defineAuthorizedOrganizationByokUseCase({
335336
await tx.execute(
336337
sql`SELECT set_config('lock_timeout', ${`${ORGANIZATION_BYOK_LOCK_TIMEOUT_MS}ms`}, true)`
337338
)
338-
await tx.execute(
339-
sql`SELECT pg_advisory_xact_lock(hashtextextended(${`byok:organization:${context.organizationId}:${input.providerId}`}, 0))`
339+
await acquireAdvisoryXactLock(
340+
tx,
341+
'organization_byok',
342+
`byok:organization:${context.organizationId}:${input.providerId}`
340343
)
341344

342345
const [{ keyCount }] = await tx

‎apps/sim/lib/billing/core/usage-log.ts‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,7 @@ import { isOrgScopedSubscription } from '@/lib/billing/subscriptions/utils'
2727
import type { InternalUsageLogSource } from '@/lib/billing/usage-sources'
2828
import { asOrchestrationError, OrchestrationError } from '@/lib/core/orchestration/types'
2929
import { HttpError } from '@/lib/core/utils/http-error'
30+
import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks'
3031
import type { DbClient, DbOrTx } from '@/lib/db/types'
3132

3233
const logger = createLogger('UsageLog')
@@ -740,7 +741,7 @@ export async function recordCumulativeUsage(
740741
set_config('lock_timeout', ${`${CUMULATIVE_FLUSH_LOCK_TIMEOUT_MS}ms`}, true)
741742
`)
742743
enterStage('lock')
743-
await tx.execute(sql`select pg_advisory_xact_lock(hashtextextended(${eventKey}, 0))`)
744+
await acquireAdvisoryXactLock(tx, 'usage_log_event', eventKey)
744745

745746
enterStage('read')
746747
const [existing] = await tx

‎apps/sim/lib/billing/enterprise-owner-claim.ts‎

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,7 @@ import {
3030
processOutboxEventById,
3131
} from '@/lib/core/outbox/service'
3232
import { getBaseUrl } from '@/lib/core/utils/urls'
33+
import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks'
3334
import type { DbOrTx } from '@/lib/db/types'
3435
import { computeInvitationExpiry, INVITATION_EXPIRY_DAYS } from '@/lib/invitations/expiry'
3536
import { MAX_INVITE_EMAILS, MAX_INVITE_WORKSPACES } from '@/lib/invitations/limits'
@@ -492,8 +493,10 @@ export async function createEnterpriseOwnerClaim(
492493
})
493494
const requestKey = buildClaimRequestKey(input, normalized)
494495
const result = await db.transaction(async (tx) => {
495-
await tx.execute(
496-
sql`select pg_advisory_xact_lock(hashtextextended(${`enterprise-owner-claim:${normalized.ownerEmail}`}, 0))`
496+
await acquireAdvisoryXactLock(
497+
tx,
498+
'enterprise_owner_claim',
499+
`enterprise-owner-claim:${normalized.ownerEmail}`
497500
)
498501
const [accountCreatedDuringReview] = await tx
499502
.select({ id: user.id })

‎apps/sim/lib/billing/organizations/billing-identity-lock.ts‎

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
import { sql } from 'drizzle-orm'
2+
import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks'
23
import type { DbOrTx } from '@/lib/db/types'
34

45
const USER_BILLING_IDENTITY_LOCK_TIMEOUT_MS = 5_000
@@ -12,7 +13,5 @@ export async function acquireUserBillingIdentityLock(tx: DbOrTx, userId: string)
1213
await tx.execute(
1314
sql`select set_config('lock_timeout', ${`${USER_BILLING_IDENTITY_LOCK_TIMEOUT_MS}ms`}, true)`
1415
)
15-
await tx.execute(
16-
sql`select pg_advisory_xact_lock(hashtextextended(${`user-billing-identity:${userId}`}, 0))`
17-
)
16+
await acquireAdvisoryXactLock(tx, 'user_billing_identity', `user-billing-identity:${userId}`)
1817
}

‎apps/sim/lib/billing/organizations/membership.ts‎

Lines changed: 6 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -51,6 +51,7 @@ import { isBillingEnabled } from '@/lib/core/config/env-flags'
5151
import { OrchestrationError } from '@/lib/core/orchestration/types'
5252
import { enqueueOutboxEvent } from '@/lib/core/outbox/service'
5353
import { revokeWorkspaceCredentialMembershipsTx } from '@/lib/credentials/access'
54+
import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks'
5455
import { isRetryableTransactionError } from '@/lib/db/transaction'
5556
import type { DbOrTx } from '@/lib/db/types'
5657
import { acquireInvitationMutationLocks } from '@/lib/invitations/locks'
@@ -83,8 +84,10 @@ export async function acquireOrganizationMutationLock(
8384
await tx.execute(
8485
sql`select set_config('lock_timeout', ${`${ORG_MEMBERSHIP_LOCK_TIMEOUT_MS}ms`}, true)`
8586
)
86-
await tx.execute(
87-
sql`select pg_advisory_xact_lock(hashtextextended(${`organization-mutation:${organizationId}`}, 0))`
87+
await acquireAdvisoryXactLock(
88+
tx,
89+
'organization_mutation',
90+
`organization-mutation:${organizationId}`
8891
)
8992
}
9093

@@ -108,9 +111,7 @@ export async function acquireOrgMembershipLock(
108111
await tx.execute(
109112
sql`select set_config('lock_timeout', ${`${ORG_MEMBERSHIP_LOCK_TIMEOUT_MS}ms`}, true)`
110113
)
111-
await tx.execute(
112-
sql`select pg_advisory_xact_lock(hashtextextended(${`${userId}:${organizationId}`}, 0))`
113-
)
114+
await acquireAdvisoryXactLock(tx, 'organization_membership', `${userId}:${organizationId}`)
114115
}
115116

116117
/**

‎apps/sim/lib/billing/webhooks/enterprise.ts‎

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -43,6 +43,7 @@ import {
4343
enqueueOutboxEvents,
4444
patchOutboxEventPayload,
4545
} from '@/lib/core/outbox/service'
46+
import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks'
4647
import { sendEmail } from '@/lib/messaging/email/mailer'
4748
import { getFromEmailAddress } from '@/lib/messaging/email/utils'
4849
import { captureServerEvent } from '@/lib/posthog/server'
@@ -173,8 +174,10 @@ async function reconcileManualEnterpriseSubscription(
173174

174175
const coreResult = await db.transaction(async (tx) => {
175176
await acquireOrganizationMutationLock(tx, referenceId)
176-
await tx.execute(
177-
sql`select pg_advisory_xact_lock(hashtextextended(${`stripe-subscription:${stripeSubscription.id}`}, 0))`
177+
await acquireAdvisoryXactLock(
178+
tx,
179+
'stripe_subscription',
180+
`stripe-subscription:${stripeSubscription.id}`
178181
)
179182
// The authoritative Stripe read happened under a durable subscription
180183
// lease. Fence the write before touching billing state so a crashed holder

‎apps/sim/lib/credential-groups/enrollments.ts‎

Lines changed: 9 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -44,6 +44,7 @@ import type {
4444
CredentialGroupEnrollmentRecord,
4545
InviteCredentialGroupEnrollmentsInput,
4646
} from '@/lib/credential-groups/types'
47+
import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks'
4748
import type { DbOrTx } from '@/lib/db/types'
4849
import { sendEmail } from '@/lib/messaging/email/mailer'
4950
import { getFromEmailAddress } from '@/lib/messaging/email/utils'
@@ -177,8 +178,10 @@ export async function lockCredentialGroupEnrollmentLifecycle(
177178
enrollmentId: string
178179
): Promise<void> {
179180
if (!enrollmentId.trim()) throw new Error('Credential group enrollment ID is required')
180-
await executor.execute(
181-
sql`SELECT pg_advisory_xact_lock(hashtextextended(${`credential-group-enrollment:${enrollmentId}`}, 0))`
181+
await acquireAdvisoryXactLock(
182+
executor,
183+
'credential_group_enrollment',
184+
`credential-group-enrollment:${enrollmentId}`
182185
)
183186
}
184187

@@ -190,8 +193,10 @@ async function lockCredentialGroupInvitationTarget(
190193
): Promise<void> {
191194
if (!groupId.trim()) throw new Error('Credential group ID is required')
192195
if (!email.trim()) throw new Error('Credential group enrollment email is required')
193-
await executor.execute(
194-
sql`SELECT pg_advisory_xact_lock(hashtextextended(${`credential-group-invitation:${groupId}:${email}`}, 0))`
196+
await acquireAdvisoryXactLock(
197+
executor,
198+
'credential_group_invitation',
199+
`credential-group-invitation:${groupId}:${email}`
195200
)
196201
}
197202

0 commit comments

Comments
 (0)