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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 11 additions & 6 deletions apps/sim/app/api/schedules/execute/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -885,10 +886,12 @@ async function recoverStaleDatabaseScheduleJobs(now: Date): Promise<void> {
const disabledScheduleIds = new Set<string>()

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'
)
Expand Down Expand Up @@ -1065,10 +1068,12 @@ async function tryStartDatabaseScheduleJob(jobId: string): Promise<DatabaseSched
const now = new Date()

return 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) return 'capacity_full'
if (!acquired) return 'capacity_full'

const [row] = await tx
.select({
Expand Down
5 changes: 2 additions & 3 deletions apps/sim/app/api/workspaces/[id]/byok-keys/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'

Expand Down Expand Up @@ -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() })
Expand Down
9 changes: 3 additions & 6 deletions apps/sim/ee/workspace-forking/lib/lineage/lineage.ts
Original file line number Diff line number Diff line change
@@ -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 {
Expand Down Expand Up @@ -110,9 +111,7 @@ export async function setForkLockTimeout(tx: DbOrTx): Promise<void> {
* unnecessary serialization, never a correctness issue.
*/
export async function acquireForkEdgeLock(tx: DbOrTx, childWorkspaceId: string): Promise<void> {
await tx.execute(
sql`select pg_advisory_xact_lock(hashtextextended(${`fork-edge:${childWorkspaceId}`}, 0))`
)
await acquireAdvisoryXactLock(tx, 'fork_edge', `fork-edge:${childWorkspaceId}`)
}

/**
Expand All @@ -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<void> {
await tx.execute(
sql`select pg_advisory_xact_lock(hashtextextended(${`fork-target:${targetWorkspaceId}`}, 0))`
)
await acquireAdvisoryXactLock(tx, 'fork_target', `fork-target:${targetWorkspaceId}`)
}
7 changes: 5 additions & 2 deletions apps/sim/lib/api-key/application/organization-byok-keys.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand Down Expand Up @@ -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
Expand Down
3 changes: 2 additions & 1 deletion apps/sim/lib/billing/core/usage-log.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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')
Expand Down Expand Up @@ -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
Expand Down
7 changes: 5 additions & 2 deletions apps/sim/lib/billing/enterprise-owner-claim.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand Down Expand Up @@ -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 })
Expand Down
5 changes: 2 additions & 3 deletions apps/sim/lib/billing/organizations/billing-identity-lock.ts
Original file line number Diff line number Diff line change
@@ -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
Expand All @@ -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}`)
}
11 changes: 6 additions & 5 deletions apps/sim/lib/billing/organizations/membership.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand Down Expand Up @@ -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}`
)
}

Expand All @@ -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}`)
}

/**
Expand Down
7 changes: 5 additions & 2 deletions apps/sim/lib/billing/webhooks/enterprise.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand Down Expand Up @@ -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
Expand Down
13 changes: 9 additions & 4 deletions apps/sim/lib/credential-groups/enrollments.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand Down Expand Up @@ -177,8 +178,10 @@ export async function lockCredentialGroupEnrollmentLifecycle(
enrollmentId: string
): Promise<void> {
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}`
)
}

Expand All @@ -190,8 +193,10 @@ async function lockCredentialGroupInvitationTarget(
): Promise<void> {
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}`
)
}

Expand Down
9 changes: 6 additions & 3 deletions apps/sim/lib/credential-groups/oauth.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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()
Expand Down Expand Up @@ -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({
Expand Down
9 changes: 6 additions & 3 deletions apps/sim/lib/credential-groups/service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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 }
Expand Down Expand Up @@ -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()
Expand Down
9 changes: 6 additions & 3 deletions apps/sim/lib/credential-groups/slack-managed-users.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand All @@ -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'
Expand Down Expand Up @@ -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()
Expand Down
9 changes: 5 additions & 4 deletions apps/sim/lib/credentials/env-locks.ts
Original file line number Diff line number Diff line change
@@ -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
Expand All @@ -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<void> {
async function lockEnvMap(tx: DbOrTx, tag: string, lockKey: string): Promise<void> {
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<void> {
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<void> {
await lockEnvMap(tx, userId)
await lockEnvMap(tx, 'personal_env_map', userId)
}
7 changes: 3 additions & 4 deletions apps/sim/lib/credentials/managed-oauth.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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')
Expand Down Expand Up @@ -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 {
Expand Down
Loading
Loading