diff --git a/apps/sim/app/api/files/uploads/utils.ts b/apps/sim/app/api/files/uploads/utils.ts index e19ea1a9210..17e29390761 100644 --- a/apps/sim/app/api/files/uploads/utils.ts +++ b/apps/sim/app/api/files/uploads/utils.ts @@ -4,8 +4,8 @@ import { type InternalFileUploadSession, internalFileUploadSessionSchema, } from '@/lib/api/contracts/upload-sessions' +import { orchestrationFailureResponse } from '@/lib/api/server/orchestration-response' import { getSession } from '@/lib/auth' -import { asOrchestrationError, statusForOrchestrationError } from '@/lib/core/orchestration/types' import type { UploadSessionRecord } from '@/lib/uploads/upload-session/service' import type { UploadActor, UploadPurposeResult } from '@/app/api/files/uploads/finalizers' @@ -32,13 +32,7 @@ export async function requireUploadUser(): Promise { null ) - await db.transaction(async (tx) => { - await tx.insert(workflow).values( - await buildNewWorkflowRow(tx, { - id: newWorkflowId, - userId: session.user.id, - workspaceId: targetWorkspaceId, - folderId: null, - name: dedupedName, - description: sourceWorkflow.description, - variables: sourceWorkflow.variables || {}, - }) - ) - }) + await db.transaction((tx) => + insertNewWorkflowRow(tx, { + id: newWorkflowId, + userId: session.user.id, + workspaceId: targetWorkspaceId, + folderId: null, + name: dedupedName, + description: sourceWorkflow.description, + variables: sourceWorkflow.variables || {}, + }) + ) // Save using existing persistence logic const saveResult = await saveWorkflowToNormalizedTables(newWorkflowId, importedData, { @@ -228,7 +226,7 @@ export const POST = withRouteHandler(async (request: NextRequest) => { copilotChatsImported, }) } catch (error) { - if (error instanceof OrchestrationError && error.code === 'not_found') { + if (asOrchestrationError(error)?.code === 'not_found') { return NextResponse.json({ error: 'Target workspace not found' }, { status: 404 }) } logger.error('Error importing workflow', error) diff --git a/apps/sim/app/api/table/utils.ts b/apps/sim/app/api/table/utils.ts index 875b8005d0c..62d27b5e86e 100644 --- a/apps/sim/app/api/table/utils.ts +++ b/apps/sim/app/api/table/utils.ts @@ -2,8 +2,8 @@ import { createLogger } from '@sim/logger' import { permissionSatisfies } from '@sim/platform-authz/workspace' import { toError } from '@sim/utils/errors' import { NextResponse } from 'next/server' +import { orchestrationFailureResponse } from '@/lib/api/server/orchestration-response' import { - asOrchestrationError, messageForOrchestrationError, type OrchestrationErrorCode, statusForOrchestrationError, @@ -137,13 +137,7 @@ export function orchestrationErrorResponse(error: unknown): NextResponse | null const lockResponse = tableLockErrorResponse(error) if (lockResponse) return lockResponse - const classified = asOrchestrationError(error) - if (!classified) return null - - return NextResponse.json( - { error: classified.message }, - { status: statusForOrchestrationError(classified.code) } - ) + return orchestrationFailureResponse(error) } /** diff --git a/apps/sim/app/api/v1/admin/workflows/import/route.ts b/apps/sim/app/api/v1/admin/workflows/import/route.ts index 12a02a6c855..47c31583a1c 100644 --- a/apps/sim/app/api/v1/admin/workflows/import/route.ts +++ b/apps/sim/app/api/v1/admin/workflows/import/route.ts @@ -28,10 +28,10 @@ import { and, eq, isNull } from 'drizzle-orm' import { NextResponse } from 'next/server' import { adminV1ImportWorkflowContract } from '@/lib/api/contracts/v1/admin' import { parseRequest } from '@/lib/api/server' -import { OrchestrationError } from '@/lib/core/orchestration/types' +import { asOrchestrationError } from '@/lib/core/orchestration/types' import { withRouteHandler } from '@/lib/core/utils/with-route-handler' import { parseWorkflowJson } from '@/lib/workflows/operations/import-export' -import { buildNewWorkflowRow } from '@/lib/workflows/persistence/new-workflow-row' +import { insertNewWorkflowRow } from '@/lib/workflows/persistence/new-workflow-row' import { prepareWorkflowStateForPersistence } from '@/lib/workflows/persistence/prepare-state' import { saveWorkflowToNormalizedTables } from '@/lib/workflows/persistence/utils' import { deduplicateWorkflowName } from '@/lib/workflows/utils' @@ -116,18 +116,16 @@ export const POST = withRouteHandler( const workflowId = generateId() const dedupedName = await deduplicateWorkflowName(workflowName, workspaceId, folderId || null) - await db.transaction(async (tx) => { - await tx.insert(workflow).values( - await buildNewWorkflowRow(tx, { - id: workflowId, - userId: workspaceData.ownerId, - workspaceId, - folderId: folderId || null, - name: dedupedName, - description: workflowDescription, - }) - ) - }) + await db.transaction((tx) => + insertNewWorkflowRow(tx, { + id: workflowId, + userId: workspaceData.ownerId, + workspaceId, + folderId: folderId || null, + name: dedupedName, + description: workflowDescription, + }) + ) /** * Same normalization the editor and the v1 import API run, via the one @@ -183,7 +181,7 @@ export const POST = withRouteHandler( return NextResponse.json(response) } catch (error) { - if (error instanceof OrchestrationError && error.code === 'not_found') { + if (asOrchestrationError(error)?.code === 'not_found') { return notFoundResponse('Workspace') } if (error instanceof FolderNotFoundError) { diff --git a/apps/sim/app/api/v1/admin/workspaces/[id]/import/route.ts b/apps/sim/app/api/v1/admin/workspaces/[id]/import/route.ts index 44ee2523b7b..34964a8e287 100644 --- a/apps/sim/app/api/v1/admin/workspaces/[id]/import/route.ts +++ b/apps/sim/app/api/v1/admin/workspaces/[id]/import/route.ts @@ -45,7 +45,7 @@ import { extractWorkflowsFromZip, parseWorkflowJson, } from '@/lib/workflows/operations/import-export' -import { buildNewWorkflowRow } from '@/lib/workflows/persistence/new-workflow-row' +import { insertNewWorkflowRow } from '@/lib/workflows/persistence/new-workflow-row' import { prepareWorkflowStateForPersistence } from '@/lib/workflows/persistence/prepare-state' import { saveWorkflowToNormalizedTables } from '@/lib/workflows/persistence/utils' import { deduplicateWorkflowName } from '@/lib/workflows/utils' @@ -350,18 +350,16 @@ async function importSingleWorkflow( const workflowId = generateId() const dedupedName = await deduplicateWorkflowName(workflowName, workspaceId, targetFolderId) - await db.transaction(async (tx) => { - await tx.insert(workflow).values( - await buildNewWorkflowRow(tx, { - id: workflowId, - userId: ownerId, - workspaceId, - folderId: targetFolderId, - name: dedupedName, - description: workflowData.metadata?.description || 'Imported via Admin API', - }) - ) - }) + await db.transaction((tx) => + insertNewWorkflowRow(tx, { + id: workflowId, + userId: ownerId, + workspaceId, + folderId: targetFolderId, + name: dedupedName, + description: workflowData.metadata?.description || 'Imported via Admin API', + }) + ) /** * Same normalization the editor, the v1 import API and the single-workflow diff --git a/apps/sim/app/api/workspaces/[id]/route.ts b/apps/sim/app/api/workspaces/[id]/route.ts index afc662ad4b7..1ddd4e05cc2 100644 --- a/apps/sim/app/api/workspaces/[id]/route.ts +++ b/apps/sim/app/api/workspaces/[id]/route.ts @@ -5,9 +5,9 @@ import { and, eq, isNull } from 'drizzle-orm' import { type NextRequest, NextResponse } from 'next/server' import { deleteWorkspaceBodySchema, updateWorkspaceContract } from '@/lib/api/contracts' import { parseRequest, validationErrorResponse } from '@/lib/api/server' +import { orchestrationFailureResponse } from '@/lib/api/server/orchestration-response' import { getSession } from '@/lib/auth' import { changeWorkspaceStoragePayerInTx } from '@/lib/billing/storage/payer-transfer' -import { OrchestrationError, statusForOrchestrationError } from '@/lib/core/orchestration/types' import { captureServerEvent } from '@/lib/posthog/server' import { archiveWorkspace } from '@/lib/workspaces/lifecycle' @@ -307,6 +307,20 @@ export const DELETE = withRouteHandler( }, request, }) + if (archiveResult.archivedProject) { + recordAudit({ + workspaceId, + actorId: session.user.id, + actorName: session.user.name, + actorEmail: session.user.email, + action: AuditAction.PROJECT_ARCHIVED, + resourceType: AuditResourceType.PROJECT, + resourceId: archiveResult.archivedProject.id, + resourceName: archiveResult.archivedProject.name, + description: `Archived Project "${archiveResult.archivedProject.name}" with its last active environment`, + request, + }) + } captureServerEvent( session.user.id, @@ -317,12 +331,8 @@ export const DELETE = withRouteHandler( return NextResponse.json({ success: true }) } catch (error) { - if (error instanceof OrchestrationError) { - return NextResponse.json( - { error: error.message }, - { status: statusForOrchestrationError(error.code) } - ) - } + const failure = orchestrationFailureResponse(error, 'Failed to delete workspace') + if (failure) return failure logger.error(`Error deleting workspace ${workspaceId}:`, error) return NextResponse.json({ error: 'Failed to delete workspace' }, { status: 500 }) } diff --git a/apps/sim/ee/access-control/components/group-detail.tsx b/apps/sim/ee/access-control/components/group-detail.tsx index 166c83e3284..ee7190065d5 100644 --- a/apps/sim/ee/access-control/components/group-detail.tsx +++ b/apps/sim/ee/access-control/components/group-detail.tsx @@ -1784,16 +1784,6 @@ export function GroupDetail({ {configTab === 'platform' && (
- - setEditingConfig((previous) => ({ - ...previous, - deniedPartialAccessProjectIssues: value, - })) - } - />
))} + + setEditingConfig((previous) => ({ + ...previous, + deniedPartialAccessProjectIssues: value, + })) + } + />
)} diff --git a/apps/sim/ee/access-control/components/project-issue-restrictions.tsx b/apps/sim/ee/access-control/components/project-issue-restrictions.tsx index d2404c7b548..689ab44738c 100644 --- a/apps/sim/ee/access-control/components/project-issue-restrictions.tsx +++ b/apps/sim/ee/access-control/components/project-issue-restrictions.tsx @@ -1,7 +1,8 @@ 'use client' -import { Checkbox, Chip } from '@sim/emcn' +import { Checkbox, Chip, Info, OverflowText } from '@sim/emcn' import { isApiClientError } from '@/lib/api/client/errors' +import { SettingsSection } from '@/app/workspace/[workspaceId]/settings/components/settings-section/settings-section' import { useProjects } from '@/hooks/queries/projects' interface ProjectIssueRestrictionsProps { @@ -10,7 +11,10 @@ interface ProjectIssueRestrictionsProps { onChange: (value: string[]) => void } -/** Project choices use the authorized inventory; policy remains enforced at the application boundary. */ +/** + * Project choices use the authorized inventory; policy remains enforced at the + * application boundary. + */ export function ProjectIssueRestrictions({ organizationId, value, @@ -20,44 +24,58 @@ export function ProjectIssueRestrictions({ if (projects.isPending || (isApiClientError(projects.error) && projects.error.status === 503)) { return null } + const choices = projects.data?.pages.flatMap((page) => page.projects) ?? [] + if (!projects.error && choices.length === 0) return null const selected = new Set(value) return ( -
-

Restrict Issues for partial-access teammates

-

- For selected Projects, teammates governed by this group need access to every active - environment to use Issues. -

+ + For selected Projects, teammates governed by this group need access to every active + environment to use Issues. + + } + action={ + projects.hasNextPage ? ( + void projects.fetchNextPage()} + > + {projects.isFetchingNextPage ? 'Loading…' : 'Load more'} + + ) : undefined + } + > {projects.error && ( -

{projects.error.message}

+

{projects.error.message}

)} - {projects.data?.pages - .flatMap((page) => page.projects) - .map((project) => ( -
+
+ ) } diff --git a/apps/sim/ee/workspace-forking/application/lineage-details.ts b/apps/sim/ee/workspace-forking/application/lineage-details.ts index 0b640febdc0..951bf9061b9 100644 --- a/apps/sim/ee/workspace-forking/application/lineage-details.ts +++ b/apps/sim/ee/workspace-forking/application/lineage-details.ts @@ -1,7 +1,6 @@ import { db } from '@sim/db' import { workspace } from '@sim/db/schema' -import { eq } from 'drizzle-orm' -import { readForkSyncNewWorkflowsExcluded } from '@/lib/workflows/persistence/new-workflow-row' +import { and, eq, isNull } from 'drizzle-orm' import { getEffectiveWorkspacePermission } from '@/lib/workspaces/permissions/utils' import { getForkChildren, getForkParent } from '@/ee/workspace-forking/lib/lineage/lineage' import { getUndoableRunForTarget } from '@/ee/workspace-forking/lib/promote/promote-run-store' @@ -24,6 +23,21 @@ async function withViewerAccess { + const [row] = await db + .select({ excluded: workspace.forkSyncNewWorkflowsExcluded }) + .from(workspace) + .where(and(eq(workspace.id, workspaceId), isNull(workspace.archivedAt))) + .limit(1) + return row?.excluded ?? false +} + export const getWorkspaceForkLineageDetails = defineForkUseCase({ operation: forkOperations.discover, availability: true, @@ -40,7 +54,7 @@ export const getWorkspaceForkLineageDetails = defineForkUseCase({ getForkChildren(workspaceId), getUndoableRunForTarget(db, workspaceId), // Lineage-uniform, so this workspace's own value is the lineage's value. - readForkSyncNewWorkflowsExcluded(db, workspaceId), + readForkSyncNewWorkflowsExcluded(workspaceId), ]) const [parent, children] = await Promise.all([ diff --git a/apps/sim/ee/workspace-forking/lib/copy/copy-workflows.ts b/apps/sim/ee/workspace-forking/lib/copy/copy-workflows.ts index 21c9c21b83d..847f2ec2984 100644 --- a/apps/sim/ee/workspace-forking/lib/copy/copy-workflows.ts +++ b/apps/sim/ee/workspace-forking/lib/copy/copy-workflows.ts @@ -26,7 +26,6 @@ import { type SubBlockTransform, } from '@/lib/workflows/references/remap-references' import type { CanonicalModeOverrides } from '@/lib/workflows/subblocks/visibility' -import { lockActiveWorkspace } from '@/lib/workspaces/active-workspace' import { deriveForkBlockId, type ForkBlockIdResolver, @@ -503,7 +502,6 @@ export async function copyWorkflowStateIntoTarget( requestId = 'unknown', } = params - await lockActiveWorkspace(tx, targetWorkspaceId) const targetFolderId = sourceMeta.folderId ? (folderIdMap.get(sourceMeta.folderId) ?? null) : null const varIdMapping = new Map() diff --git a/apps/sim/ee/workspace-forking/lib/create-fork.ts b/apps/sim/ee/workspace-forking/lib/create-fork.ts index 995652275e4..69c161ed921 100644 --- a/apps/sim/ee/workspace-forking/lib/create-fork.ts +++ b/apps/sim/ee/workspace-forking/lib/create-fork.ts @@ -1,5 +1,5 @@ import { db } from '@sim/db' -import { permissions, projectWorkspace, workflow, workspace } from '@sim/db/schema' +import { permissions, projectWorkspace, workspace } from '@sim/db/schema' import { createLogger } from '@sim/logger' import type { PermissionType } from '@sim/platform-authz/workspace' import { getErrorMessage } from '@sim/utils/errors' @@ -9,7 +9,7 @@ import type { Workspace } from '@/lib/api/contracts/workspaces' import { enqueueOutboxEvent } from '@/lib/core/outbox/service' import { requireForkProject } from '@/lib/projects/membership' import { buildDefaultWorkflowArtifacts } from '@/lib/workflows/defaults' -import { buildNewWorkflowRow } from '@/lib/workflows/persistence/new-workflow-row' +import { insertNewWorkflowRow } from '@/lib/workflows/persistence/new-workflow-row' import { saveWorkflowToNormalizedTables } from '@/lib/workflows/persistence/utils' import { collectReferencedDocumentIds, @@ -511,17 +511,15 @@ export async function createFork(params: CreateForkParams): Promise { resetDbChainMock() + queueTableRows(workspace, [{ name: 'Target' }]) + queueTableRows(workspace, [{ archivedAt: null, forkSyncNewWorkflowsExcluded: false }]) mockGetUsersWithPermissions.mockResolvedValue([]) mockLoadSourceDeployedStates.mockResolvedValue({ deployedWorkflows: [], diff --git a/apps/sim/ee/workspace-forking/lib/promote/promote.ts b/apps/sim/ee/workspace-forking/lib/promote/promote.ts index 266a3bca4ab..60df6e03a4e 100644 --- a/apps/sim/ee/workspace-forking/lib/promote/promote.ts +++ b/apps/sim/ee/workspace-forking/lib/promote/promote.ts @@ -29,6 +29,7 @@ import { type ForkRemapKind, } from '@/lib/workflows/references/remap-references' import { getMcpServerMetaByIds } from '@/lib/workflows/references/resources' +import { lockActiveWorkspace } from '@/lib/workspaces/active-workspace' import { findWorkspaceOperationReceipt, lockWorkspaceOperationRequest, @@ -629,6 +630,8 @@ export async function promoteFork(params: PromoteForkParams): Promise page.nextCursor, staleTime: PROJECT_LIST_STALE_TIME, retry: (failureCount, error) => - !(isApiClientError(error) && error.status === 503) && failureCount < 3, + failureCount < 1 && + (!isApiClientError(error) || + error.status === 408 || + error.status === 429 || + (error.status >= 500 && error.status !== 503)), enabled: Boolean(organizationId), }) } diff --git a/apps/sim/lib/api/contracts/projects.ts b/apps/sim/lib/api/contracts/projects.ts index f25549855e0..0339cd07738 100644 --- a/apps/sim/lib/api/contracts/projects.ts +++ b/apps/sim/lib/api/contracts/projects.ts @@ -1,15 +1,18 @@ import { z } from 'zod' -import { nonEmptyIdSchema } from '@/lib/api/contracts/primitives' +import { + nonEmptyIdSchema, + organizationIdSchema, + workspaceIdSchema, +} from '@/lib/api/contracts/primitives' import { defineRouteContract } from '@/lib/api/contracts/types' import { createProjectInputSchema } from '@/lib/projects/create-input' -export const projectEnvironmentSchema = z.object({ +const projectEnvironmentSchema = z.object({ id: nonEmptyIdSchema, name: z.string(), forkedFromWorkspaceId: nonEmptyIdSchema.nullable(), }) -export type ProjectEnvironment = z.output -export const projectSchema = z.object({ +const projectSchema = z.object({ id: nonEmptyIdSchema, name: z.string(), organizationId: nonEmptyIdSchema.nullable(), @@ -20,33 +23,27 @@ export const projectSchema = z.object({ environments: z.array(projectEnvironmentSchema), capabilities: z.object({ administer: z.boolean(), issues: z.boolean() }), }) -export type Project = z.output -export const projectParamsSchema = z.object({ id: nonEmptyIdSchema }) -export type ProjectParams = z.input -export const projectQuerySchema = z.object({ - organizationId: nonEmptyIdSchema.optional(), - workspaceId: nonEmptyIdSchema.optional(), +const projectParamsSchema = z.object({ id: nonEmptyIdSchema }) +const projectQuerySchema = z.object({ + organizationId: organizationIdSchema.optional(), + workspaceId: workspaceIdSchema.optional(), }) -export type ProjectQuery = z.input -export const listProjectsQuerySchema = z.object({ - organizationId: nonEmptyIdSchema.optional(), +const listProjectsQuerySchema = z.object({ + organizationId: organizationIdSchema.optional(), cursor: nonEmptyIdSchema.optional(), limit: z.coerce.number().int().min(1).max(100).default(50), }) -export type ListProjectsQuery = z.input -export const listProjectsResponseSchema = z.object({ +const listProjectsResponseSchema = z.object({ projects: z.array(projectSchema), nextCursor: nonEmptyIdSchema.nullable(), }) -export type ListProjectsResponse = z.output export const listProjectsContract = defineRouteContract({ method: 'GET', path: '/api/projects', query: listProjectsQuerySchema, response: { mode: 'json', schema: listProjectsResponseSchema }, }) -export const getProjectResponseSchema = z.object({ project: projectSchema }) -export type GetProjectResponse = z.output +const getProjectResponseSchema = z.object({ project: projectSchema }) export const getProjectContract = defineRouteContract({ method: 'GET', path: '/api/projects/[id]', @@ -54,10 +51,8 @@ export const getProjectContract = defineRouteContract({ query: projectQuerySchema, response: { mode: 'json', schema: getProjectResponseSchema }, }) -export const renameProjectBodySchema = z.object({ name: z.string().trim().min(1).max(100) }) -export type RenameProjectBody = z.input -export const renameProjectResponseSchema = z.object({ id: nonEmptyIdSchema, name: z.string() }) -export type RenameProjectResponse = z.output +const renameProjectBodySchema = z.object({ name: z.string().trim().min(1).max(100) }) +const renameProjectResponseSchema = z.object({ id: nonEmptyIdSchema, name: z.string() }) export const renameProjectContract = defineRouteContract({ method: 'PATCH', path: '/api/projects/[id]', @@ -65,11 +60,10 @@ export const renameProjectContract = defineRouteContract({ body: renameProjectBodySchema, response: { mode: 'json', schema: renameProjectResponseSchema }, }) -export const archiveProjectResponseSchema = z.object({ +const archiveProjectResponseSchema = z.object({ id: nonEmptyIdSchema, archived: z.boolean(), }) -export type ArchiveProjectResponse = z.output export const archiveProjectContract = defineRouteContract({ method: 'DELETE', path: '/api/projects/[id]', @@ -77,8 +71,7 @@ export const archiveProjectContract = defineRouteContract({ response: { mode: 'json', schema: archiveProjectResponseSchema }, }) -export const workspaceProjectParamsSchema = z.object({ workspaceId: nonEmptyIdSchema }) -export type WorkspaceProjectParams = z.input +const workspaceProjectParamsSchema = z.object({ workspaceId: workspaceIdSchema }) export const getWorkspaceProjectContract = defineRouteContract({ method: 'GET', path: '/api/projects/by-workspace/[workspaceId]', @@ -86,13 +79,11 @@ export const getWorkspaceProjectContract = defineRouteContract({ response: { mode: 'json', schema: getProjectResponseSchema }, }) -export const createProjectBodySchema = createProjectInputSchema -export type CreateProjectBody = z.input -export const createProjectResponseSchema = z.object({ +const createProjectBodySchema = createProjectInputSchema +const createProjectResponseSchema = z.object({ project: z.object({ id: nonEmptyIdSchema, name: z.string() }), initialEnvironment: z.object({ id: nonEmptyIdSchema, name: z.string() }), }) -export type CreateProjectResponse = z.output export const createProjectContract = defineRouteContract({ method: 'POST', path: '/api/projects', diff --git a/apps/sim/lib/api/server/orchestration-response.ts b/apps/sim/lib/api/server/orchestration-response.ts new file mode 100644 index 00000000000..684371a26fc --- /dev/null +++ b/apps/sim/lib/api/server/orchestration-response.ts @@ -0,0 +1,28 @@ +import { NextResponse } from 'next/server' +import { + asOrchestrationError, + messageForOrchestrationError, + statusForOrchestrationError, +} from '@/lib/core/orchestration/types' + +/** + * Maps a classified domain failure anywhere in `error`'s cause chain to its status for a + * raw route; `null` when unclassified, so the caller logs it and returns its own 500. An + * `internal` code answers with `fallback` rather than a message that may carry internals. + */ +export function orchestrationFailureResponse( + error: unknown, + fallback = 'Internal server error' +): NextResponse | null { + const classified = asOrchestrationError(error) + if (!classified) return null + return NextResponse.json( + { + error: messageForOrchestrationError( + { error: classified.message, errorCode: classified.code }, + fallback + ), + }, + { status: statusForOrchestrationError(classified.code) } + ) +} diff --git a/apps/sim/lib/billing/organizations/membership.ts b/apps/sim/lib/billing/organizations/membership.ts index 220b0dba82d..e7b25750441 100644 --- a/apps/sim/lib/billing/organizations/membership.ts +++ b/apps/sim/lib/billing/organizations/membership.ts @@ -15,7 +15,6 @@ import { organization, permissionGroupMember, permissions, - project, subscription as subscriptionTable, user, userStats, @@ -61,7 +60,7 @@ import { revokePersonalApiKeysTx, revokeUserSessionsTx, } from '@/lib/organizations/members/revocation' -import { lockProjectBackfillWrites, tryLockProject } from '@/lib/projects/membership' +import { reassignOrganizationProjects } from '@/lib/projects/membership' import { removeWorkspaceSkillMembershipsTx } from '@/lib/skills/access' import { reassignWorkflowOwnershipForWorkspaceMemberRemovalTx, @@ -557,25 +556,12 @@ async function reassignOwnedOrganizationResourcesTx({ const ownerId = ownerMembership?.userId if (!ownerId || ownerId === userId) return 0 - await lockProjectBackfillWrites(tx, workspaceIds) - const ownedProjects = await tx - .select({ id: project.id }) - .from(project) - .where(and(eq(project.organizationId, organizationId), eq(project.ownerId, userId))) - .orderBy(project.id) - for (const row of ownedProjects) { - await tryLockProject(tx, row.id) - await tx - .update(project) - .set({ ownerId, updatedAt: new Date() }) - .where( - and( - eq(project.id, row.id), - eq(project.ownerId, userId), - eq(project.organizationId, organizationId) - ) - ) - } + await reassignOrganizationProjects(tx, { + organizationId, + fromUserId: userId, + toUserId: ownerId, + workspaceIds, + }) /** Creator attribution must survive account deletion without changing document ACLs. */ await tx @@ -1850,6 +1836,12 @@ export async function transferOrganizationOwnership( .returning({ id: workspace.id }) result.workspacesReassigned = ownerUpdate.length + await reassignOrganizationProjects(tx, { + organizationId, + fromUserId: currentOwnerUserId, + toUserId: newOwnerUserId, + workspaceIds: ownerUpdate.map((workspaceRow) => workspaceRow.id), + }) const reassignedWorkspaceIds = Array.from( new Set([...billedWorkspaceIds, ...ownerUpdate.map((workspaceRow) => workspaceRow.id)]) diff --git a/apps/sim/lib/core/config/env.ts b/apps/sim/lib/core/config/env.ts index 9904b0143e3..4c889ae2a71 100644 --- a/apps/sim/lib/core/config/env.ts +++ b/apps/sim/lib/core/config/env.ts @@ -599,6 +599,7 @@ export const env = createEnv({ AGENTMAIL_DOMAIN: z.string().optional(), // Custom domain for AgentMail inboxes (default: agentmail.to) MSHIP_PLAN_MODE: z.boolean().optional(), DASHBOARDS: z.boolean().optional(), + PROJECT_API_ENABLED: z.boolean().optional(), // Fallback for the `projects` feature flag off AppConfig MSHIP_MODEL_SELECTOR: z.boolean().optional(), INBOX_ENABLED: z.boolean().optional(), // Enable inbox (Sim Mailer) on self-hosted (bypasses hosted requirements) SANDBOXES_ENABLED: z.boolean().optional(), // Enable custom sandboxes on self-hosted (bypasses hosted requirements) @@ -669,7 +670,6 @@ export const env = createEnv({ // SSO Configuration (for script-based registration) SSO_ENABLED: z.boolean().optional(), // Enable SSO functionality - PROJECT_API_ENABLED: z.boolean().optional(), // Expose Projects after backfill and contract enforcement SCIM_ENABLED: z.boolean().optional(), // Enable SCIM directory provisioning USAGE_MONITORING_ENABLED: z.boolean().optional(), // Enable organization usage monitoring on self-hosted (bypasses hosted requirements) SSO_PROVIDER_TYPE: z.enum(['oidc', 'saml']).optional(), // [REQUIRED] SSO provider type diff --git a/apps/sim/lib/core/config/feature-flags.ts b/apps/sim/lib/core/config/feature-flags.ts index 83925454c75..566d3481832 100644 --- a/apps/sim/lib/core/config/feature-flags.ts +++ b/apps/sim/lib/core/config/feature-flags.ts @@ -110,6 +110,13 @@ const FEATURE_FLAGS = { 'requires knowledge-member-access. Off-AppConfig falls back to CREDENTIAL_GROUPS.', fallback: 'CREDENTIAL_GROUPS', }, + projects: { + description: + 'Expose the Project APIs once the membership backfill has validated. Global on/off only; ' + + 'workspace creation assigns Projects and lifecycle protections apply either way. ' + + 'Off-AppConfig falls back to PROJECT_API_ENABLED.', + fallback: 'PROJECT_API_ENABLED', + }, 'knowledge-member-access': { description: 'Organization Search (live) and the permission-aware workspace connector modes: members ' + diff --git a/apps/sim/lib/db/advisory-locks.ts b/apps/sim/lib/db/advisory-locks.ts index dacbcac9c54..862e952b36e 100644 --- a/apps/sim/lib/db/advisory-locks.ts +++ b/apps/sim/lib/db/advisory-locks.ts @@ -1,4 +1,5 @@ import { type SQL, sql } from 'drizzle-orm' +import { textArrayLiteral } from '@/lib/db/arrays' import type { DbTransaction } from '@/lib/db/types' const LOCK_TAG_PATTERN = /^[a-z][a-z0-9_]*$/ @@ -45,6 +46,24 @@ export async function tryAcquireAdvisoryXactLock( return Boolean(lock?.acquired) } +/** + * Tries every transaction-scoped advisory lock in `keys` in one round trip, without + * waiting, so their order cannot deadlock. Returns whether all are held; locks taken + * alongside a refusal stay held until the transaction ends, so a caller that sees + * `false` should abort it. + */ +export async function tryAcquireAdvisoryXactLocks( + tx: DbTransaction, + tag: string, + keys: readonly string[] +): Promise { + if (keys.length === 0) return true + const [lock] = await tx.execute<{ acquired: boolean }>(sql` + SELECT bool_and(pg_try_advisory_xact_lock(hashtextextended(key, 0))) AS acquired + FROM unnest(${textArrayLiteral(keys)}) AS key ${lockTag(tag)}`) + return Boolean(lock?.acquired) +} + /** One lock of an {@link acquireAdvisoryXactLocks} set. */ export interface AdvisoryXactLockRequest { key: string diff --git a/apps/sim/lib/db/arrays.ts b/apps/sim/lib/db/arrays.ts new file mode 100644 index 00000000000..9bad5ea9325 --- /dev/null +++ b/apps/sim/lib/db/arrays.ts @@ -0,0 +1,15 @@ +import { type SQL, sql } from 'drizzle-orm' + +/** + * The pool uses fetch_types: false, so arrays must be constructed from scalar + * parameters. A JSON scalar keeps large sets below PostgreSQL's bind limit. + */ +export function textArrayLiteral(values: readonly string[]): SQL { + if (values.length > 1000) { + return sql`ARRAY(SELECT jsonb_array_elements_text(${JSON.stringify(values)}::text::jsonb))` + } + return sql`ARRAY[${sql.join( + values.map((value) => sql`${value}`), + sql`, ` + )}]::text[]` +} diff --git a/apps/sim/lib/knowledge/__integration__/workspace-lifecycle.integration.ts b/apps/sim/lib/knowledge/__integration__/workspace-lifecycle.integration.ts index 817f3903db8..5317a7da3da 100644 --- a/apps/sim/lib/knowledge/__integration__/workspace-lifecycle.integration.ts +++ b/apps/sim/lib/knowledge/__integration__/workspace-lifecycle.integration.ts @@ -16,9 +16,6 @@ vi.mock('@/lib/mcp/pubsub', () => ({ mcpPubSub: null })) vi.mock('@/lib/mcp/service', () => ({ mcpService: { clearCache: vi.fn().mockResolvedValue(undefined) }, })) -vi.mock('@/lib/workflows/lifecycle', () => ({ - archiveWorkflowsForWorkspace: vi.fn().mockResolvedValue(0), -})) import { createKnowledgeAclFixtureIds, diff --git a/apps/sim/lib/knowledge/access/predicate.ts b/apps/sim/lib/knowledge/access/predicate.ts index a316783db85..a7da9e6a427 100644 --- a/apps/sim/lib/knowledge/access/predicate.ts +++ b/apps/sim/lib/knowledge/access/predicate.ts @@ -11,6 +11,7 @@ import { user, } from '@sim/db/schema' import { type SQL, sql } from 'drizzle-orm' +import { textArrayLiteral } from '@/lib/db/arrays' import { EXTERNAL_GROUP_STALE_AFTER_MS } from '@/lib/knowledge/access/external-groups' import { SOURCE_ACL_MAX_AGE_MS } from '@/lib/knowledge/access/freshness' import { confluenceReaderGroupCondition } from '@/lib/knowledge/access/group-membership' @@ -279,17 +280,3 @@ function storedKnowledgeAccessCondition( export function aclOverlap(tokens: SQL): SQL { return sql`${document.acl} && ${tokens}` } - -/** - * The pool uses fetch_types: false, so arrays must be constructed from scalar - * parameters. A JSON scalar keeps large sets below PostgreSQL's bind limit. - */ -export function textArrayLiteral(values: readonly string[]): SQL { - if (values.length > 1000) { - return sql`ARRAY(SELECT jsonb_array_elements_text(${JSON.stringify(values)}::text::jsonb))` - } - return sql`ARRAY[${sql.join( - values.map((value) => sql`${value}`), - sql`, ` - )}]::text[]` -} diff --git a/apps/sim/lib/knowledge/connectors/member-observations.ts b/apps/sim/lib/knowledge/connectors/member-observations.ts index 2146bcac2f3..6cbc90c4121 100644 --- a/apps/sim/lib/knowledge/connectors/member-observations.ts +++ b/apps/sim/lib/knowledge/connectors/member-observations.ts @@ -26,8 +26,8 @@ import { type SQL, sql, } from 'drizzle-orm' +import { textArrayLiteral } from '@/lib/db/arrays' import type { DbOrTx } from '@/lib/db/types' -import { textArrayLiteral } from '@/lib/knowledge/access/predicate' import { walkReconciliationWindows } from '@/lib/knowledge/connectors/reconciliation-window' import { ACL_CHANGE_BATCH_SIZE, diff --git a/apps/sim/lib/knowledge/connectors/sync-persistence.ts b/apps/sim/lib/knowledge/connectors/sync-persistence.ts index b12ecba5a24..0759377789d 100644 --- a/apps/sim/lib/knowledge/connectors/sync-persistence.ts +++ b/apps/sim/lib/knowledge/connectors/sync-persistence.ts @@ -6,7 +6,7 @@ import { generateId } from '@sim/utils/id' import { truncateAtCodePoint } from '@sim/utils/string' import { and, eq, exists, inArray, isNull, lt, not, or, type SQL, sql } from 'drizzle-orm' import { getInternalApiBaseUrl } from '@/lib/core/utils/urls' -import { textArrayLiteral } from '@/lib/knowledge/access/predicate' +import { textArrayLiteral } from '@/lib/db/arrays' import { EMPTY_ACL, validateMirroredDocumentAcl, diff --git a/apps/sim/lib/knowledge/search/vector-leg.ts b/apps/sim/lib/knowledge/search/vector-leg.ts index 2113961b7f9..19155646dc7 100644 --- a/apps/sim/lib/knowledge/search/vector-leg.ts +++ b/apps/sim/lib/knowledge/search/vector-leg.ts @@ -3,7 +3,7 @@ import { document, embedding, embeddingSearch } from '@sim/db/schema' import { createLogger } from '@sim/logger' import { getErrorMessage, getPostgresErrorCode } from '@sim/utils/errors' import { and, eq, inArray, type SQL, sql } from 'drizzle-orm' -import { textArrayLiteral } from '@/lib/knowledge/access/predicate' +import { textArrayLiteral } from '@/lib/db/arrays' import { runSearchQuery, type SearchBudget, diff --git a/apps/sim/lib/knowledge/service.ts b/apps/sim/lib/knowledge/service.ts index 423de6a68bc..7e34ad16f89 100644 --- a/apps/sim/lib/knowledge/service.ts +++ b/apps/sim/lib/knowledge/service.ts @@ -31,9 +31,10 @@ import { } from '@/lib/core/resource-scope' import { resourceScopeCondition } from '@/lib/core/resource-scope.server' import { generateRestoreName } from '@/lib/core/utils/restore-name' +import { textArrayLiteral } from '@/lib/db/arrays' import { findActiveFolder, resolveRestoredFolderId } from '@/lib/folders/queries' import { isKnowledgeMemberAccessAvailable } from '@/lib/knowledge/access/availability' -import { knowledgeAccessCondition, textArrayLiteral } from '@/lib/knowledge/access/predicate' +import { knowledgeAccessCondition } from '@/lib/knowledge/access/predicate' import type { KnowledgeAccessProvider } from '@/lib/knowledge/access/types' import { mirrorsSourceAcls } from '@/lib/knowledge/connectors/access-modes' import { diff --git a/apps/sim/lib/projects/README.md b/apps/sim/lib/projects/README.md index 88c3a330a9f..dc3deb8a198 100644 --- a/apps/sim/lib/projects/README.md +++ b/apps/sim/lib/projects/README.md @@ -2,7 +2,7 @@ A Project groups environments. An environment is an existing `workspace` record; there is no separate environment table. Every newly created Project starts with an environment. -Project APIs return HTTP 503 until the deployment enables `PROJECT_API_ENABLED`. Workspace creation always assigns a Project atomically. The API control defaults off and does not disable assignment. Existing assigned Projects always retain their lifecycle protections, including fork inheritance and disconnect behavior, even if activation is disabled. +Project APIs return HTTP 503 until the `projects` feature flag is on (AppConfig on hosted deployments; the `PROJECT_API_ENABLED` secret elsewhere). The flag defaults off and gates only the Project APIs: workspace creation always assigns a Project atomically, and assigned Projects keep their lifecycle protections (fork inheritance, disconnect) either way. ## Choose the creation flow @@ -12,9 +12,9 @@ Project APIs return HTTP 503 until the deployment enables `PROJECT_API_ENABLED`. | Existing workspace creation UI or caller | `POST /api/workspaces` | Creates a workspace and automatically creates its Project, preserving the existing workspace response. | | Create another environment by forking | Existing workspace fork operation | Inherits the source workspace's Project; unassigned legacy families remain unassigned until backfilled. | -Both POST endpoints are internal, session-authenticated APIs. `POST /api/projects` is not a public `/api/v2` endpoint and does not accept API-key principals. This foundation does not remove or deprecate existing workspace creation endpoints. +Both POST endpoints are internal, session-authenticated APIs. `POST /api/projects` is not a public `/api/v2` endpoint and does not accept API-key principals. Existing workspace creation endpoints remain supported. -Do not call both creation endpoints for one onboarding flow: each creates a new workspace and a new Project. Neither endpoint attaches a workspace to an existing Project. A general Project environment-creation endpoint is follow-up work. +Do not call both creation endpoints for one onboarding flow: each creates a new workspace and a new Project. Neither endpoint attaches a workspace to an existing Project. ## Create a Project and its first environment @@ -71,8 +71,12 @@ Existing callers can continue to use `createWorkspaceContract` and `POST /api/wo The workspace and its Project are created atomically. The generated Project name is `Support workspace - Project`; long names are bounded to 100 characters while retaining the suffix. The response remains `{ "workspace": ... }` with HTTP 200, without a new Project response wrapper. Call `GET /api/projects/by-workspace/[workspaceId]` when an existing workspace caller needs its authorized Project details. +## Archiving + +Archiving a workspace through the existing workspace deletion flow never strands an active Project: when the workspace is its Project's last active environment, the Project is archived in the same transaction. Account deletion applies the same rule to a Project whose surviving environments are all archived. `DELETE /api/projects/[id]` archives a Project and every environment together. + ## Server implementation The Project route calls `createProject` in `application/create-project.ts`. Existing workspace callers continue through their current creation paths. Both use the shared transaction primitive in `lib/workspaces/create.ts`; surface adapters must not independently commit Project and workspace creation. -Project descriptions and Project-scoped files are not part of this creation contract. Project-scoped files and a designated Project brief are follow-up work. +Project descriptions and Project-scoped files are not part of this creation contract. diff --git a/apps/sim/lib/projects/__integration__/foundation.integration.ts b/apps/sim/lib/projects/__integration__/foundation.integration.ts index 876079dc5d7..7e219bc9d80 100644 --- a/apps/sim/lib/projects/__integration__/foundation.integration.ts +++ b/apps/sim/lib/projects/__integration__/foundation.integration.ts @@ -23,13 +23,16 @@ import { createWorkspaceApiKeyPrincipal, } from '@sim/testing/factories/principal.factory' import { createDeferred } from '@sim/testing/helpers/deferred' +import { featureFlagsMock, featureFlagsMockFns } from '@sim/testing/mocks/feature-flags.mock' import { getErrorMessage, getPostgresErrorCode } from '@sim/utils/errors' -import { sleep } from '@sim/utils/helpers' import { generateId } from '@sim/utils/id' -import { and, eq, inArray, sql } from 'drizzle-orm' +import { and, eq, inArray, or, sql } from 'drizzle-orm' import { NextRequest } from 'next/server' import { afterAll, beforeEach, describe, expect, it, vi } from 'vitest' -import { removeUserFromOrganization } from '@/lib/billing/organizations/membership' +import { + removeUserFromOrganization, + transferOrganizationOwnership, +} from '@/lib/billing/organizations/membership' import { prepareProjectsForAccountDeletion } from '@/lib/projects/account-deletion' import { archiveProject, @@ -45,6 +48,7 @@ import { createProjectForWorkspace, lockProject, lockWorkspaceProject, + projectBackfillLockKey, splitForkProject, transferWorkspaceProjects, } from '@/lib/projects/membership' @@ -64,18 +68,25 @@ vi.hoisted(() => { process.env.ADMIN_API_KEY = 'project-fixture-admin-key' }) +vi.mock('@/lib/core/config/feature-flags', () => featureFlagsMock) + +function setProjectsEnabled(enabled: boolean) { + featureFlagsMockFns.mockIsFeatureEnabled.mockImplementation( + async (flag) => enabled && flag === 'projects' + ) +} + beforeEach(() => { - vi.stubEnv('PROJECT_API_ENABLED', 'true') + setProjectsEnabled(true) }) const users: string[] = [] const organizations: string[] = [] -const environments: string[] = [] const request = { requestId: 'project-foundation-integration', headers: new Headers() } const checks: { name: string; status: 'passed' | 'failed'; durationMs: number; error?: string }[] = [] -/** Exercises durable auth and lifecycle invariants against real Postgres, including concurrent writers. */ +/** Registers a test and records its status and duration in the suite report. */ function check(name: string, run: () => Promise) { it(name, async () => { const started = performance.now() @@ -125,7 +136,6 @@ async function fixture(org = true, count = 2) { let projectId = '' await db.transaction(async (tx) => { for (const [index, id] of ids.entries()) { - environments.push(id) await tx.insert(workspace).values({ id, name: `Environment ${index}`, @@ -164,6 +174,48 @@ async function fixture(org = true, count = 2) { return { ownerId, teammateId, outsiderId, organizationId, projectId, ids, owner, teammate } } +/** + * Waits until the operation under test is blocked behind `blockerPid`: a session whose + * current statement contains `waitingIn` (an advisory lock tag or row-lock clause), so an + * unrelated waiter cannot release the barrier early. + */ +async function waitUntilBlockedBy(blockerPid: number, waitingIn: string) { + await expect + .poll( + async () => + ( + await db.execute(sql` + SELECT 1 FROM pg_stat_activity + WHERE ${blockerPid} = ANY(pg_blocking_pids(pid)) + AND position(${waitingIn.toLowerCase()} in lower(query)) > 0 + `) + ).length, + { timeout: 2000, interval: 10 } + ) + .toBeGreaterThan(0) +} + +async function addOrganizationProject(organizationId: string, ownerId: string) { + const workspaceId = generateId() + const projectId = await db.transaction(async (tx) => { + await tx.insert(workspace).values({ + id: workspaceId, + name: 'Sibling environment', + ownerId, + billedAccountUserId: ownerId, + organizationId, + workspaceMode: 'organization', + }) + return createProjectForWorkspace(tx, { + workspaceId, + name: 'Sibling environment', + organizationId, + ownerId, + }) + }) + return { workspaceId, projectId } +} + async function addWorkflow(workspaceId: string, userId: string) { const id = generateId() const now = new Date() @@ -186,16 +238,11 @@ afterAll(async () => { process.env.PROJECT_FOUNDATION_REPORT_PATH ?? resolve('test-results/project-foundation.json') await mkdir(dirname(reportPath), { recursive: true }) await writeFile(reportPath, JSON.stringify({ checks }, null, 2)) - if (environments.length) { - const memberships = await db - .select({ id: projectWorkspace.projectId }) - .from(projectWorkspace) - .where(inArray(projectWorkspace.workspaceId, environments)) - const ids = [...new Set(memberships.map((row) => row.id))] - await db.delete(projectWorkspace).where(inArray(projectWorkspace.workspaceId, environments)) - if (ids.length) await db.delete(project).where(inArray(project.id, ids)) - await db.delete(workspace).where(inArray(workspace.id, environments)) - } + if (users.length) + await db + .delete(workspace) + .where(or(inArray(workspace.ownerId, users), inArray(workspace.billedAccountUserId, users))) + if (users.length) await db.delete(project).where(inArray(project.ownerId, users)) if (organizations.length) await db.delete(organization).where(inArray(organization.id, organizations)) if (users.length) await db.delete(user).where(inArray(user.id, users)) @@ -205,7 +252,7 @@ describe('Project foundation at the database and application boundary', () => { check( 'workspace creation and fork/disconnect assign Projects while APIs remain disabled', async () => { - vi.stubEnv('PROJECT_API_ENABLED', 'false') + setProjectsEnabled(false) const f = await fixture(false, 1) const source = await db.transaction((tx) => createWorkspaceInTransaction(tx, { @@ -219,7 +266,6 @@ describe('Project foundation at the database and application boundary', () => { skipDefaultWorkflow: true, }) ) - environments.push(source.id) expect( await db.select().from(projectWorkspace).where(eq(projectWorkspace.workspaceId, source.id)) ).toHaveLength(1) @@ -231,7 +277,6 @@ describe('Project foundation at the database and application boundary', () => { userId: f.ownerId, name: 'Legacy child', }) - environments.push(fork.workspace.id) expect( await db .select() @@ -249,8 +294,9 @@ describe('Project foundation at the database and application boundary', () => { 'Project operations remain unavailable until API activation with no partial creation', async () => { const f = await fixture(false, 1) - vi.stubEnv('PROJECT_API_ENABLED', 'false') + setProjectsEnabled(false) const input = { projectId: f.projectId } + const [before] = await db.select().from(project).where(eq(project.id, f.projectId)) const calls = [ () => createProject.execute({ @@ -283,70 +329,46 @@ describe('Project foundation at the database and application boundary', () => { expect( await db.select().from(workspace).where(eq(workspace.ownerId, f.ownerId)) ).toHaveLength(1) - const [record] = await db.select().from(project).where(eq(project.id, f.projectId)) - expect(record.name).toBe('Environment 0 - Project') - expect(record.archivedAt).toBeNull() + expect(await db.select().from(project).where(eq(project.id, f.projectId))).toEqual([before]) } ) - check( - 'disabling activation preserves assigned fork membership and lifecycle protections', - async () => { - const f = await fixture(false, 1) - vi.stubEnv('PROJECT_API_ENABLED', 'false') - const parent = await getWorkspaceWithOwner(f.ids[0]) - if (!parent) throw new Error('Missing source fixture') - const fork = await createFork({ - source: parent, - policy: await getWorkspaceCreationPolicy({ userId: f.ownerId }), - userId: f.ownerId, - name: 'Assigned child', - }) - environments.push(fork.workspace.id) - const [membership] = await db - .select() - .from(projectWorkspace) - .where(eq(projectWorkspace.workspaceId, fork.workspace.id)) - expect(membership.projectId).toBe(f.projectId) - await unlinkForkEdge({ parentWorkspaceId: f.ids[0], childWorkspaceId: fork.workspace.id }) - const [detached] = await db - .select() - .from(projectWorkspace) - .where(eq(projectWorkspace.workspaceId, fork.workspace.id)) - expect(detached.projectId).not.toBe(f.projectId) - await expect( - archiveWorkspace(fork.workspace.id, { requestId: 'disabled-project-rollout' }) - ).rejects.toMatchObject({ code: 'conflict' }) - } - ) - - check('new workspaces receive Projects while Project APIs remain disabled', async () => { - vi.stubEnv('PROJECT_API_ENABLED', 'false') + check('a detached fork keeps its own Project, archived with its only environment', async () => { const f = await fixture(false, 1) - const created = await db.transaction((tx) => - createWorkspaceInTransaction(tx, { - userId: f.ownerId, - name: 'Writer activation', - organizationId: null, - observedOrganizationId: null, - governingPermissionGroupOrganizationId: null, - workspaceMode: 'personal', - billedAccountUserId: f.ownerId, - skipDefaultWorkflow: true, - }) - ) - environments.push(created.id) - expect( - await db.select().from(projectWorkspace).where(eq(projectWorkspace.workspaceId, created.id)) - ).toHaveLength(1) + const parent = await getWorkspaceWithOwner(f.ids[0]) + if (!parent) throw new Error('Missing source fixture') + const fork = await createFork({ + source: parent, + policy: await getWorkspaceCreationPolicy({ userId: f.ownerId }), + userId: f.ownerId, + name: 'Assigned child', + }) + const [membership] = await db + .select() + .from(projectWorkspace) + .where(eq(projectWorkspace.workspaceId, fork.workspace.id)) + expect(membership.projectId).toBe(f.projectId) + await unlinkForkEdge({ parentWorkspaceId: f.ids[0], childWorkspaceId: fork.workspace.id }) + const [detached] = await db + .select() + .from(projectWorkspace) + .where(eq(projectWorkspace.workspaceId, fork.workspace.id)) + expect(detached.projectId).not.toBe(f.projectId) + await expect( + archiveWorkspace(fork.workspace.id, { requestId: 'detached-fork-archive' }) + ).resolves.toMatchObject({ archived: true }) + const [detachedProject] = await db + .select() + .from(project) + .where(eq(project.id, detached.projectId)) + expect(detachedProject.archivedAt).not.toBeNull() }) - check('legacy fork and disconnect refuse a partially assigned subtree', async () => { + check('fork refuses a partially assigned lineage', async () => { const f = await fixture(false, 3) await db .delete(projectWorkspace) .where(inArray(projectWorkspace.workspaceId, f.ids.slice(0, 2))) - vi.stubEnv('PROJECT_API_ENABLED', 'false') const parent = await getWorkspaceWithOwner(f.ids[1]) if (!parent) throw new Error('Missing source fixture') await expect( @@ -357,9 +379,6 @@ describe('Project foundation at the database and application boundary', () => { name: 'Invalid child', }) ).rejects.toMatchObject({ code: 'conflict' }) - await expect( - unlinkForkEdge({ parentWorkspaceId: f.ids[0], childWorkspaceId: f.ids[1] }) - ).rejects.toMatchObject({ code: 'conflict' }) const [child] = await db.select().from(workspace).where(eq(workspace.id, f.ids[1])) expect(child.forkedFromWorkspaceId).toBe(f.ids[0]) expect(await db.select().from(workspace).where(eq(workspace.ownerId, f.ownerId))).toHaveLength( @@ -388,13 +407,9 @@ describe('Project foundation at the database and application boundary', () => { await Promise.race([read.promise, writer]) await db.transaction(async (tx) => { const [lock] = await tx.execute<{ acquired: boolean }>(sql` - SELECT pg_try_advisory_xact_lock(hashtextextended(${`project-backfill:${f.ids[0]}`}, 0)) AS acquired + SELECT pg_try_advisory_xact_lock(hashtextextended(${projectBackfillLockKey(f.ids[0])}, 0)) AS acquired `) expect(lock.acquired).toBe(false) - const [unrelated] = await tx.execute<{ acquired: boolean }>(sql` - SELECT pg_try_advisory_xact_lock(hashtextextended('project-backfill:unrelated', 0)) AS acquired - `) - expect(unrelated.acquired).toBe(true) }) } finally { release.resolve() @@ -402,7 +417,7 @@ describe('Project foundation at the database and application boundary', () => { } await db.transaction(async (tx) => { await tx.execute( - sql`SELECT pg_advisory_xact_lock(hashtextextended(${`project-backfill:${f.ids[0]}`}, 0))` + sql`SELECT pg_advisory_xact_lock(hashtextextended(${projectBackfillLockKey(f.ids[0])}, 0))` ) await createProjectForWorkspace(tx, { workspaceId: f.ids[0], @@ -424,13 +439,14 @@ describe('Project foundation at the database and application boundary', () => { const parent = await getWorkspaceWithOwner(f.ids[0]) if (!parent) throw new Error('Missing source fixture') const policy = await getWorkspaceCreationPolicy({ userId: f.ownerId }) - const locked = createDeferred() + const locked = createDeferred() const release = createDeferred() const backfill = db.transaction(async (tx) => { await tx.execute( - sql`SELECT pg_advisory_xact_lock(hashtextextended(${`project-backfill:${f.ids[0]}`}, 0))` + sql`SELECT pg_advisory_xact_lock(hashtextextended(${projectBackfillLockKey(f.ids[0])}, 0))` ) - locked.resolve() + const [connection] = await tx.execute<{ pid: number }>(sql`SELECT pg_backend_pid() AS pid`) + locked.resolve(connection.pid) await release.promise return createProjectForWorkspace(tx, { workspaceId: f.ids[0], @@ -445,27 +461,13 @@ describe('Project foundation at the database and application boundary', () => { policy, userId: f.ownerId, name: 'Concurrent child', - }).then((result) => { - environments.push(result.workspace.id) - return result }) - let blocked = false try { - for (let attempt = 0; attempt < 100; attempt++) { - const rows = await db.execute<{ waiting: boolean }>(sql`SELECT EXISTS ( - SELECT 1 FROM pg_locks WHERE locktype = 'advisory' AND mode = 'ShareLock' AND NOT granted - ) AS waiting`) - if (rows[0]?.waiting) { - blocked = true - break - } - await sleep(20) - } + await waitUntilBlockedBy(await locked.promise, "lock='project_backfill'") } finally { release.resolve() } const [projectId, result] = await Promise.all([backfill, fork]) - expect(blocked).toBe(true) const [membership] = await db .select() .from(projectWorkspace) @@ -496,7 +498,6 @@ describe('Project foundation at the database and application boundary', () => { }, request, }) - environments.push(result.initialEnvironment.id) const [created] = await db.select().from(project).where(eq(project.id, result.project.id)) expect(created).toMatchObject({ name: 'Customer support', @@ -639,44 +640,6 @@ describe('Project foundation at the database and application boundary', () => { } ) - check('workspace creation and forks commit exactly one Project membership', async () => { - const f = await fixture(false, 1) - const source = await db.transaction((tx) => - createWorkspaceInTransaction(tx, { - userId: f.ownerId, - name: 'New environment', - organizationId: null, - observedOrganizationId: null, - governingPermissionGroupOrganizationId: null, - workspaceMode: 'personal', - billedAccountUserId: f.ownerId, - skipDefaultWorkflow: true, - }) - ) - environments.push(source.id) - const before = await db - .select() - .from(projectWorkspace) - .where(eq(projectWorkspace.workspaceId, source.id)) - expect(before).toHaveLength(1) - const policy = await getWorkspaceCreationPolicy({ userId: f.ownerId }) - const parent = await getWorkspaceWithOwner(source.id) - if (!parent) throw new Error('Missing source fixture') - const fork = await createFork({ - source: parent, - policy, - userId: f.ownerId, - name: 'Child environment', - }) - environments.push(fork.workspace.id) - const child = await db - .select() - .from(projectWorkspace) - .where(eq(projectWorkspace.workspaceId, fork.workspace.id)) - expect(child).toHaveLength(1) - expect(child[0].projectId).toBe(before[0].projectId) - }) - check( 'a departing organization member transfers Project lifecycle ownership to the org owner', async () => { @@ -742,7 +705,7 @@ describe('Project foundation at the database and application boundary', () => { input, request, }) - ).rejects.toThrow() + ).rejects.toThrow('cannot perform operation') await expect( renameProject.execute({ principal: f.teammate, @@ -825,13 +788,21 @@ describe('Project foundation at the database and application boundary', () => { } ) - check('concurrent individual removals preserve the last active environment', async () => { + check('concurrent removal of every environment archives the Project', async () => { const f = await fixture(false) const results = await Promise.allSettled(f.ids.map((id) => archiveWorkspace(id, request))) - expect(results.filter((result) => result.status === 'fulfilled')).toHaveLength(1) - expect(results.filter((result) => result.status === 'rejected')).toHaveLength(1) + expect(results).toMatchObject(f.ids.map(() => ({ status: 'fulfilled' }))) const rows = await db.select().from(workspace).where(inArray(workspace.id, f.ids)) - expect(rows.filter((row) => !row.archivedAt)).toHaveLength(1) + expect(rows.map((row) => row.archivedAt)).not.toContain(null) + const [record] = await db.select().from(project).where(eq(project.id, f.projectId)) + expect(record.archivedAt).not.toBeNull() + }) + + check('removing one of several environments keeps the Project active', async () => { + const f = await fixture(false) + await archiveWorkspace(f.ids[0], request) + const [record] = await db.select().from(project).where(eq(project.id, f.projectId)) + expect(record.archivedAt).toBeNull() }) check( @@ -876,18 +847,7 @@ describe('Project foundation at the database and application boundary', () => { {} ) try { - let waiting = false - for (let attempt = 0; attempt < 100; attempt++) { - const rows = await db.execute( - sql`SELECT 1 FROM pg_stat_activity WHERE ${blocker} = ANY(pg_blocking_pids(pid))` - ) - if (rows.length) { - waiting = true - break - } - await sleep(10) - } - expect(waiting).toBe(true) + await waitUntilBlockedBy(blocker, 'for share') } finally { release.resolve() await archive @@ -916,18 +876,7 @@ describe('Project foundation at the database and application boundary', () => { (error: unknown) => error ) try { - let waiting = false - for (let attempt = 0; attempt < 100; attempt++) { - const rows = await db.execute(sql` - SELECT 1 FROM pg_stat_activity WHERE ${blocker} = ANY(pg_blocking_pids(pid)) - `) - if (rows.length) { - waiting = true - break - } - await sleep(10) - } - expect(waiting).toBe(true) + await waitUntilBlockedBy(blocker, 'for share') } finally { release.resolve() await archive @@ -948,8 +897,15 @@ describe('Project foundation at the database and application boundary', () => { throw new Error('Abort compound archive') }) ).rejects.toThrow('Abort compound archive') - const before = await db.select().from(workspace).where(inArray(workspace.id, f.ids)) - expect(before.every((row) => row.archivedAt === null)).toBe(true) + const untouched = await db.select().from(workspace).where(inArray(workspace.id, f.ids)) + expect(untouched.map((row) => row.archivedAt)).toEqual([null, null]) + const activeWorkflows = await db + .select() + .from(workflow) + .where(inArray(workflow.id, workflowIds)) + expect(activeWorkflows).toMatchObject( + workflowIds.map(() => ({ archivedAt: null, isDeployed: true })) + ) const args = { principal: f.owner, input: { projectId: f.projectId }, request } await archiveProject.execute(args) await archiveProject.execute(args) @@ -1091,6 +1047,8 @@ describe('Project foundation at the database and application boundary', () => { const organizationId = source.organizationId if (!organizationId) throw new Error('Missing organization fixture') const destination = await fixture(true, 1) + const sibling = await addOrganizationProject(organizationId, source.ownerId) + const staying = await addOrganizationProject(organizationId, source.ownerId) const groupId = generateId() await db.insert(permissionGroup).values({ id: groupId, @@ -1098,17 +1056,28 @@ describe('Project foundation at the database and application boundary', () => { name: 'Source policy', createdBy: source.ownerId, isDefault: true, - config: { deniedPartialAccessProjectIssues: [source.projectId], hideTablesTab: true }, + config: { + deniedPartialAccessProjectIssues: [ + source.projectId, + staying.projectId, + sibling.projectId, + ], + hideTablesTab: true, + }, }) + const moving = [...source.ids, sibling.workspaceId] await db.transaction(async (tx) => { - await transferWorkspaceProjects(tx, source.ids, destination.organizationId) + await transferWorkspaceProjects(tx, moving, destination.organizationId) await tx .update(workspace) .set({ organizationId: destination.organizationId }) - .where(eq(workspace.id, source.ids[0])) + .where(inArray(workspace.id, moving)) }) const [group] = await db.select().from(permissionGroup).where(eq(permissionGroup.id, groupId)) - expect(group.config).toEqual({ deniedPartialAccessProjectIssues: [], hideTablesTab: true }) + expect(group.config).toEqual({ + deniedPartialAccessProjectIssues: [staying.projectId], + hideTablesTab: true, + }) const [moved] = await db.select().from(project).where(eq(project.id, source.projectId)) expect(moved).toMatchObject({ organizationId: destination.organizationId, @@ -1117,6 +1086,26 @@ describe('Project foundation at the database and application boundary', () => { } ) + check('organization ownership transfer moves the previous owner’s Projects', async () => { + const f = await fixture(true, 1) + if (!f.organizationId) throw new Error('Missing organization fixture') + await db.insert(member).values({ + id: generateId(), + organizationId: f.organizationId, + userId: f.teammateId, + role: 'member', + createdAt: new Date(), + }) + const result = await transferOrganizationOwnership({ + organizationId: f.organizationId, + currentOwnerUserId: f.ownerId, + newOwnerUserId: f.teammateId, + }) + expect(result).toMatchObject({ success: true }) + const [record] = await db.select().from(project).where(eq(project.id, f.projectId)) + expect(record.ownerId).toBe(f.teammateId) + }) + check( 'organization deletion preserves Project identity and assigns its former owner explicitly', async () => { @@ -1142,7 +1131,7 @@ describe('Project foundation at the database and application boundary', () => { ) check( - 'account deletion preview reports a surviving Project losing its last active environment', + 'account deletion archives a surviving Project losing its last active environment', async () => { const f = await fixture(false) await db @@ -1164,14 +1153,11 @@ describe('Project foundation at the database and application boundary', () => { .where(eq(workspace.id, f.ids[1])) const plan = await getAccountDeletionPlan(f.ownerId) expect(plan.workspacesToDelete.map((row) => row.id)).toEqual([f.ids[0]]) - expect(plan.blockers).toEqual([ - { code: 'project_lifecycle', message: expect.stringContaining('Archive') }, - ]) - await expect( - db.transaction((tx) => prepareProjectsForAccountDeletion(tx, f.ownerId, [f.ids[0]])) - ).rejects.toMatchObject({ code: 'conflict' }) - await db.transaction((tx) => archiveProjectInTransaction(tx, f.projectId)) - expect((await getAccountDeletionPlan(f.ownerId)).blockers).toEqual([]) + expect(plan.blockers).toEqual([]) + await db.transaction((tx) => prepareProjectsForAccountDeletion(tx, f.ownerId, [f.ids[0]])) + const [record] = await db.select().from(project).where(eq(project.id, f.projectId)) + expect(record.archivedAt).not.toBeNull() + expect(record.ownerId).toBe(f.teammateId) } ) @@ -1216,18 +1202,7 @@ describe('Project foundation at the database and application boundary', () => { prepareProjectsForAccountDeletion(tx, f.ownerId, [f.ids[1]]) ) try { - let waiting = false - for (let attempt = 0; attempt < 100; attempt++) { - const rows = await db.execute( - sql`SELECT 1 FROM pg_stat_activity WHERE ${blocker} = ANY(pg_blocking_pids(pid))` - ) - if (rows.length) { - waiting = true - break - } - await sleep(10) - } - expect(waiting).toBe(true) + await waitUntilBlockedBy(blocker, "lock='project'") } finally { release.resolve() await unlink @@ -1257,18 +1232,7 @@ describe('Project foundation at the database and application boundary', () => { const blocker = await held.promise const ban = disableUserResources(f.ownerId) try { - let waiting = false - for (let attempt = 0; attempt < 100; attempt++) { - const rows = await db.execute( - sql`SELECT 1 FROM pg_stat_activity WHERE ${blocker} = ANY(pg_blocking_pids(pid))` - ) - if (rows.length) { - waiting = true - break - } - await sleep(10) - } - expect(waiting).toBe(true) + await waitUntilBlockedBy(blocker, "lock='project'") } finally { release.resolve() await transfer diff --git a/apps/sim/lib/projects/account-deletion.ts b/apps/sim/lib/projects/account-deletion.ts index 1831df0480d..7b0370d82f5 100644 --- a/apps/sim/lib/projects/account-deletion.ts +++ b/apps/sim/lib/projects/account-deletion.ts @@ -1,73 +1,73 @@ import { db } from '@sim/db' import { member, permissions, project, projectWorkspace, workspace } from '@sim/db/schema' +import { ORG_ADMIN_ROLES } from '@sim/platform-authz/workspace' import { and, asc, eq, inArray, ne, or, sql } from 'drizzle-orm' -import { OrchestrationError } from '@/lib/core/orchestration/types' import type { DbOrTx, DbTransaction } from '@/lib/db/types' -import { lockProject, lockProjectBackfillWrites } from '@/lib/projects/membership' +import { + lockProjectBackfillWrites, + lockProjects, + ProjectConflictError, +} from '@/lib/projects/membership' +/** Two indexed lookups; an `OR` around a membership subquery would scan every Project. */ async function loadRelatedProjects(executor: DbOrTx, userId: string, doomedWorkspaceIds: string[]) { + const doomedMemberships = doomedWorkspaceIds.length + ? await executor + .select({ projectId: projectWorkspace.projectId }) + .from(projectWorkspace) + .where(inArray(projectWorkspace.workspaceId, doomedWorkspaceIds)) + : [] return executor - .select({ id: project.id }) + .select() .from(project) .where( or( eq(project.ownerId, userId), - doomedWorkspaceIds.length - ? sql`${project.id} in ( - select ${projectWorkspace.projectId} from ${projectWorkspace} - where ${inArray(projectWorkspace.workspaceId, doomedWorkspaceIds)} - )` + doomedMemberships.length + ? inArray( + project.id, + doomedMemberships.map((row) => row.projectId) + ) : undefined ) ) .orderBy(asc(project.id)) } -interface ProjectDeletionDecision { - blocker?: string - remove?: boolean - ownerId?: string -} +type ProjectDeletionDecision = + | { blocker: string } + | { remove: true } + | { archive: boolean; ownerId?: string } -async function planProjectDeletion( +/** + * An org admin, else a teammate who administers every surviving environment. With `hold`, + * the successor's membership or grants stay share-locked until commit, so the handoff + * cannot land on someone demoted concurrently. + */ +async function findProjectSuccessor( executor: DbOrTx, record: typeof project.$inferSelect, userId: string, - doomed: Set -): Promise { - const members = await executor - .select({ id: workspace.id, archivedAt: workspace.archivedAt }) - .from(projectWorkspace) - .innerJoin(workspace, eq(workspace.id, projectWorkspace.workspaceId)) - .where(eq(projectWorkspace.projectId, record.id)) - const survivors = members.filter((row) => !doomed.has(row.id)) - if (!survivors.length) { - return record.organizationId || record.ownerId !== userId - ? { blocker: 'Account deletion cannot remove another owner’s Project' } - : { remove: true } - } - if (!record.archivedAt && survivors.every((row) => row.archivedAt)) { - return { - blocker: 'Archive the Project before deleting its last active environment with your account', - } - } - if (record.ownerId !== userId) return {} + survivorIds: string[], + hold: boolean +): Promise { if (record.organizationId) { - const [successor] = await executor + const adminQuery = executor .select({ userId: member.userId }) .from(member) .where( and( eq(member.organizationId, record.organizationId), ne(member.userId, userId), - inArray(member.role, ['owner', 'admin']) + inArray(member.role, ORG_ADMIN_ROLES) ) ) .orderBy(asc(member.userId)) .limit(1) - if (successor) return { ownerId: successor.userId } + const [admin] = await (hold ? adminQuery.for('share') : adminQuery) + if (admin) return admin.userId } - const [successor] = await executor + const [teammate] = await executor .select({ userId: permissions.userId }) .from(permissions) .where( @@ -75,18 +75,81 @@ async function planProjectDeletion( eq(permissions.entityType, 'workspace'), eq(permissions.permissionType, 'admin'), ne(permissions.userId, userId), - inArray( - permissions.entityId, - survivors.map((row) => row.id) - ) + inArray(permissions.entityId, survivorIds) ) ) .groupBy(permissions.userId) - .having(sql`count(*) = ${survivors.length}`) + .having(sql`count(*) = ${survivorIds.length}`) .orderBy(asc(permissions.userId)) .limit(1) - return successor - ? { ownerId: successor.userId } + if (!teammate || !hold) return teammate?.userId ?? null + const held = await executor + .select({ id: permissions.id }) + .from(permissions) + .where( + and( + eq(permissions.entityType, 'workspace'), + eq(permissions.permissionType, 'admin'), + eq(permissions.userId, teammate.userId), + inArray(permissions.entityId, survivorIds) + ) + ) + .for('share') + return held.length === survivorIds.length ? teammate.userId : null +} + +interface ProjectEnvironment { + id: string + archivedAt: Date | null +} + +/** Every environment of `projectIds`, in one query, keyed by Project. */ +async function loadProjectEnvironments(executor: DbOrTx, projectIds: string[]) { + const rows = projectIds.length + ? await executor + .select({ + projectId: projectWorkspace.projectId, + id: workspace.id, + archivedAt: workspace.archivedAt, + }) + .from(projectWorkspace) + .innerJoin(workspace, eq(workspace.id, projectWorkspace.workspaceId)) + .where(inArray(projectWorkspace.projectId, projectIds)) + : [] + const byProject = new Map() + for (const { projectId, ...environment } of rows) { + const environments = byProject.get(projectId) + if (environments) environments.push(environment) + else byProject.set(projectId, [environment]) + } + return byProject +} + +async function planProjectDeletion( + executor: DbOrTx, + record: typeof project.$inferSelect, + members: ProjectEnvironment[], + userId: string, + doomed: Set, + hold: boolean +): Promise { + const survivors = members.filter((row) => !doomed.has(row.id)) + if (!survivors.length) { + return record.organizationId || record.ownerId !== userId + ? { blocker: 'Account deletion cannot remove another owner’s Project' } + : { remove: true } + } + const archive = !record.archivedAt && survivors.every((row) => row.archivedAt) + if (record.ownerId !== userId) return { archive } + const ownerId = await findProjectSuccessor( + executor, + record, + userId, + survivors.map((row) => row.id), + hold + ) + return ownerId + ? { archive, ownerId } : { blocker: 'Give a teammate admin access to every environment before deleting the Project owner’s account', @@ -99,13 +162,22 @@ export async function getProjectAccountDeletionBlockers( doomedWorkspaceIds: string[] ): Promise { const records = await loadRelatedProjects(db, userId, doomedWorkspaceIds) + const environments = await loadProjectEnvironments( + db, + records.map((record) => record.id) + ) const doomed = new Set(doomedWorkspaceIds) const blockers: string[] = [] - for (const { id } of records) { - const [record] = await db.select().from(project).where(eq(project.id, id)) - if (!record) continue - const decision = await planProjectDeletion(db, record, userId, doomed) - if (decision.blocker) blockers.push(decision.blocker) + for (const record of records) { + const decision = await planProjectDeletion( + db, + record, + environments.get(record.id) ?? [], + userId, + doomed, + false + ) + if ('blocker' in decision) blockers.push(decision.blocker) } return blockers } @@ -125,31 +197,46 @@ export async function prepareProjectsForAccountDeletion( ...ownedEnvironments.map((row) => row.id), ]) const locked = new Set() + let records: (typeof project.$inferSelect)[] for (;;) { - const records = await loadRelatedProjects(tx, userId, doomedWorkspaceIds) + records = await loadRelatedProjects(tx, userId, doomedWorkspaceIds) const pending = records.filter((row) => !locked.has(row.id)) if (!pending.length) break - for (const { id } of pending) { - await lockProject(tx, id) - locked.add(id) - } + await lockProjects( + tx, + pending.map((row) => row.id) + ) + for (const { id } of pending) locked.add(id) } - const records = await loadRelatedProjects(tx, userId, doomedWorkspaceIds) + const environments = await loadProjectEnvironments( + tx, + records.map((record) => record.id) + ) const doomed = new Set(doomedWorkspaceIds) - for (const { id } of records) { - await lockProject(tx, id) - const [record] = await tx.select().from(project).where(eq(project.id, id)) - if (!record) continue - const decision = await planProjectDeletion(tx, record, userId, doomed) - if (decision.blocker) throw new OrchestrationError('conflict', decision.blocker) - if (decision.remove) { - await tx.delete(projectWorkspace).where(eq(projectWorkspace.projectId, id)) - await tx.delete(project).where(eq(project.id, id)) + const now = new Date() + for (const record of records) { + const decision = await planProjectDeletion( + tx, + record, + environments.get(record.id) ?? [], + userId, + doomed, + true + ) + if ('blocker' in decision) throw new ProjectConflictError(decision.blocker) + if ('remove' in decision) { + await tx.delete(projectWorkspace).where(eq(projectWorkspace.projectId, record.id)) + await tx.delete(project).where(eq(project.id, record.id)) + continue } - if (!decision.ownerId) continue + if (!decision.archive && !decision.ownerId) continue await tx .update(project) - .set({ ownerId: decision.ownerId, updatedAt: new Date() }) - .where(eq(project.id, id)) + .set({ + ...(decision.archive ? { archivedAt: now } : {}), + ...(decision.ownerId ? { ownerId: decision.ownerId } : {}), + updatedAt: now, + }) + .where(eq(project.id, record.id)) } } diff --git a/apps/sim/lib/projects/application/authorization.ts b/apps/sim/lib/projects/application/authorization.ts index 47ebc4f7b43..eab358ab472 100644 --- a/apps/sim/lib/projects/application/authorization.ts +++ b/apps/sim/lib/projects/application/authorization.ts @@ -1,52 +1,82 @@ import type { Principal, SessionPrincipal } from '@sim/auth/principal' -import { member, permissions, project, projectWorkspace, workspace } from '@sim/db/schema' +import { + member, + permissionGroup, + permissions, + project, + projectWorkspace, + workspace, +} from '@sim/db/schema' +import { createLogger } from '@sim/logger' import { isOrgAdminRole } from '@sim/platform-authz/workspace' -import { and, asc, eq, inArray } from 'drizzle-orm' +import { and, asc, eq, inArray, sql } from 'drizzle-orm' import { PrincipalKindAuthorizationError } from '@/lib/core/application/workspace-authorization' import { OrchestrationError } from '@/lib/core/orchestration/types' +import { textArrayLiteral } from '@/lib/db/arrays' import type { DbTransaction } from '@/lib/db/types' import { CAPABILITY_RULES, refuseCapability } from '@/lib/permission-groups/capabilities' import { acquirePermissionGroupOrgLock } from '@/lib/permission-groups/locks' import { resolveVerifiedUserAccessControlContext } from '@/lib/permission-groups/resolve.server' -import type { ProjectOperation } from '@/lib/projects/application/operations' +import { type ProjectOperation, projectOperations } from '@/lib/projects/application/operations' import { lockProject } from '@/lib/projects/membership' +const logger = createLogger('ProjectAuthorization') + export function requireProjectPrincipal( principal: Principal, - operation: Pick + operation: Pick ): asserts principal is SessionPrincipal { if (principal.kind !== 'session') throw new PrincipalKindAuthorizationError(principal.kind, operation.id) } -/** The complete environment set is loaded server-side; hidden environments never enter the result. */ -export async function authorizeProject( +type ProjectRecord = typeof project.$inferSelect + +interface ProjectEnvironmentAccess { + id: string + name: string + organizationId: string | null + archivedAt: Date | null + parentId: string | null + permission: string | null +} + +export interface ProjectAuthorizationInput { + organizationId?: string + workspaceId?: string +} + +/** + * `hold` locks the Project and the rows the decision reads until commit, for callers that + * act on it. `snapshot` takes no locks and relies on the caller's read-only snapshot. + */ +type ProjectAccessMode = 'hold' | 'snapshot' + +type ProjectAccess = Awaited> + +/** + * Loads the caller's org role and every environment with its grant for `records` in three + * queries, whatever their count. A snapshot read also learns, in one more query, which + * records any permission group restricts, so unrestricted ones skip per-environment policy. + */ +async function loadProjectAccess( tx: DbTransaction, - principal: SessionPrincipal, - operation: ProjectOperation, - input: { projectId: string; organizationId?: string; workspaceId?: string } + userId: string, + records: ProjectRecord[], + mode: ProjectAccessMode ) { - await lockProject(tx, input.projectId) - const [record] = await tx.select().from(project).where(eq(project.id, input.projectId)).limit(1) - if ( - !record || - (input.organizationId !== undefined && input.organizationId !== record.organizationId) - ) { - throw new OrchestrationError('not_found', 'Project not found') - } - const [orgMember] = record.organizationId - ? await tx - .select({ role: member.role }) - .from(member) - .where( - and(eq(member.userId, principal.userId), eq(member.organizationId, record.organizationId)) - ) - .limit(1) - .for('share') + const lock = mode === 'hold' + const memberQuery = tx + .select({ organizationId: member.organizationId, role: member.role }) + .from(member) + .where(eq(member.userId, userId)) + .limit(1) + const [membership] = records.some((record) => record.organizationId) + ? await (lock ? memberQuery.for('share') : memberQuery) : [] - const orgAdmin = isOrgAdminRole(orgMember?.role) const environments = await tx .select({ + projectId: projectWorkspace.projectId, id: workspace.id, name: workspace.name, organizationId: workspace.organizationId, @@ -55,30 +85,84 @@ export async function authorizeProject( }) .from(projectWorkspace) .innerJoin(workspace, eq(workspace.id, projectWorkspace.workspaceId)) - .where(eq(projectWorkspace.projectId, record.id)) + .where( + inArray( + projectWorkspace.projectId, + records.map((record) => record.id) + ) + ) .orderBy(asc(workspace.id)) - if (environments.some((row) => row.organizationId !== record.organizationId)) - throw new OrchestrationError('conflict', 'Project ownership needs reconciliation') - const grants = environments.length - ? await tx - .select({ id: permissions.entityId, permission: permissions.permissionType }) - .from(permissions) - .where( - and( - eq(permissions.entityType, 'workspace'), - eq(permissions.userId, principal.userId), - inArray( - permissions.entityId, - environments.map((row) => row.id) - ) - ) + const grantQuery = tx + .select({ id: permissions.entityId, permission: permissions.permissionType }) + .from(permissions) + .where( + and( + eq(permissions.entityType, 'workspace'), + eq(permissions.userId, userId), + inArray( + permissions.entityId, + environments.map((row) => row.id) ) - .orderBy(asc(permissions.entityId)) - .for('share') - : [] + ) + ) + .orderBy(asc(permissions.entityId)) + const grants = environments.length ? await (lock ? grantQuery.for('share') : grantQuery) : [] const grantsById = new Map(grants.map((row) => [row.id, row.permission])) - const rows = environments.map((row) => ({ ...row, permission: grantsById.get(row.id) ?? null })) - if (operation.access === 'issues' && record.organizationId) + const environmentsByProject = new Map() + for (const { projectId, ...row } of environments) { + const access = { ...row, permission: grantsById.get(row.id) ?? null } + const rows = environmentsByProject.get(projectId) + if (rows) rows.push(access) + else environmentsByProject.set(projectId, [access]) + } + const organizationIds = [ + ...new Set(records.flatMap((record) => (record.organizationId ? [record.organizationId] : []))), + ] + const restricted = + mode === 'snapshot' && organizationIds.length + ? new Set( + ( + await tx.execute<{ id: string }>(sql` + SELECT DISTINCT denied.id FROM ${permissionGroup}, + jsonb_array_elements_text( + CASE + WHEN jsonb_typeof(${permissionGroup.config}->'deniedPartialAccessProjectIssues') = 'array' + THEN ${permissionGroup.config}->'deniedPartialAccessProjectIssues' + ELSE '[]'::jsonb + END + ) AS denied(id) + WHERE ${inArray(permissionGroup.organizationId, organizationIds)} + AND denied.id = ANY(${textArrayLiteral(records.map((record) => record.id))}) + `) + ).map((row) => row.id) + ) + : null + return { + /** Whether a permission group might restrict Issues for the Project; held reads check all. */ + mayRestrictIssues: (projectId: string) => restricted === null || restricted.has(projectId), + isOrgAdmin: (organizationId: string | null) => + organizationId !== null && + membership?.organizationId === organizationId && + isOrgAdminRole(membership.role), + environmentsFor: (projectId: string) => environmentsByProject.get(projectId) ?? [], + } +} + +/** Applies the access rules to one loaded Project; hidden environments never enter the result. */ +async function evaluateProjectAccess( + tx: DbTransaction, + principal: SessionPrincipal, + operation: ProjectOperation, + record: ProjectRecord, + access: ProjectAccess, + input: ProjectAuthorizationInput, + mode: ProjectAccessMode +) { + const orgAdmin = access.isOrgAdmin(record.organizationId) + const rows = access.environmentsFor(record.id) + if (rows.some((row) => row.organizationId !== record.organizationId)) + throw new OrchestrationError('conflict', 'Project ownership needs reconciliation') + if (mode === 'hold' && operation.access === 'issues' && record.organizationId) await acquirePermissionGroupOrgLock(tx, record.organizationId) const active = rows.filter((row) => !row.archivedAt) const visible = active.filter((row) => orgAdmin || row.permission !== null) @@ -94,7 +178,12 @@ export async function authorizeProject( 'Organization admin or admin access to every environment is required' ) let canUseIssues = !record.archivedAt && visible.length > 0 - if (canUseIssues && visible.length < active.length && record.organizationId) { + if ( + canUseIssues && + visible.length < active.length && + record.organizationId && + access.mayRestrictIssues(record.id) + ) { for (const environment of visible) { const { config } = await resolveVerifiedUserAccessControlContext( principal.userId, @@ -116,7 +205,6 @@ export async function authorizeProject( const visibleIds = new Set(visible.map((row) => row.id)) return { record, - environmentIds: rows.map((row) => row.id), canAdminister, canUseIssues, environments: visible.map((row) => ({ @@ -126,3 +214,64 @@ export async function authorizeProject( })), } } + +export type AuthorizedProject = Awaited> + +/** Authorizes one Project for `operation`; see {@link ProjectAccessMode} for locking. */ +export async function authorizeProject( + tx: DbTransaction, + principal: SessionPrincipal, + operation: ProjectOperation, + input: ProjectAuthorizationInput & { projectId: string }, + mode: ProjectAccessMode +): Promise { + if (mode === 'hold') await lockProject(tx, input.projectId) + const [record] = await tx.select().from(project).where(eq(project.id, input.projectId)).limit(1) + if ( + !record || + (input.organizationId !== undefined && input.organizationId !== record.organizationId) + ) { + throw new OrchestrationError('not_found', 'Project not found') + } + const access = await loadProjectAccess(tx, principal.userId, [record], mode) + return evaluateProjectAccess(tx, principal, operation, record, access, input, mode) +} + +/** Authorizes a listed page of Projects in id order inside the caller's read-only snapshot. */ +export async function authorizeProjectsForRead( + tx: DbTransaction, + principal: SessionPrincipal, + projectIds: string[] +): Promise { + if (!projectIds.length) return [] + const records = await tx + .select() + .from(project) + .where(inArray(project.id, projectIds)) + .orderBy(asc(project.id)) + const access = await loadProjectAccess(tx, principal.userId, records, 'snapshot') + const authorized: AuthorizedProject[] = [] + for (const record of records) { + /** One inconsistent Project must not hide the rest of the caller's page. */ + if ( + access.environmentsFor(record.id).some((row) => row.organizationId !== record.organizationId) + ) { + logger.warn('Skipping a Project whose environments need ownership reconciliation', { + projectId: record.id, + }) + continue + } + authorized.push( + await evaluateProjectAccess( + tx, + principal, + projectOperations.list, + record, + access, + {}, + 'snapshot' + ) + ) + } + return authorized +} diff --git a/apps/sim/lib/projects/application/create-project.ts b/apps/sim/lib/projects/application/create-project.ts index 95baeca9441..162f6adf6e3 100644 --- a/apps/sim/lib/projects/application/create-project.ts +++ b/apps/sim/lib/projects/application/create-project.ts @@ -7,11 +7,12 @@ import { OrchestrationError } from '@/lib/core/orchestration/types' import { refuseCapability } from '@/lib/permission-groups/capabilities' import { requireProjectPrincipal } from '@/lib/projects/application/authorization' import { projectOperations } from '@/lib/projects/application/operations' -import { type CreateProjectInput, createProjectInputSchema } from '@/lib/projects/create-input' +import type { CreateProjectInput } from '@/lib/projects/create-input' import { requireProjectApiEnabled } from '@/lib/projects/rollout.server' import { createWorkspaceWithProjectInTransaction, emitWorkspaceCreatedPlatformEvent, + WORKSPACE_USER_FK_CONSTRAINTS, } from '@/lib/workspaces/create' import { getWorkspaceCreationPolicy, @@ -32,14 +33,8 @@ export const createProject: OperationUseCase< operation: projectOperations.create, async execute({ principal, input, request }) { requireProjectPrincipal(principal, projectOperations.create) - requireProjectApiEnabled() - const parsed = createProjectInputSchema.safeParse(input) - if (!parsed.success) - throw new OrchestrationError( - 'validation', - 'A scope, Project name and initial environment name are required' - ) - const { organizationId, name, initialEnvironment } = parsed.data + await requireProjectApiEnabled() + const { organizationId, name, initialEnvironment } = input const policy = await getWorkspaceCreationPolicy({ userId: principal.userId, activeOrganizationId: organizationId, @@ -80,21 +75,16 @@ export const createProject: OperationUseCase< ) if (getPostgresErrorCode(error) === '55P03') throw new OrchestrationError( - 'locked', + 'conflict', 'This organization is being updated; retry Project creation' ) if ( getPostgresErrorCode(error) === '23503' && - getPostgresConstraintName(error) === 'workspace_owner_id_user_id_fk' - ) - throw new OrchestrationError('unauthorized', 'Unauthorized') - if ( - getPostgresErrorCode(error) === '23503' && - getPostgresConstraintName(error) === 'workspace_billed_account_user_id_user_id_fk' + WORKSPACE_USER_FK_CONSTRAINTS.has(getPostgresConstraintName(error) ?? '') ) throw new OrchestrationError( 'conflict', - 'The billing account changed; retry Project creation' + 'The owning or billing account changed; retry Project creation' ) throw error } diff --git a/apps/sim/lib/projects/application/use-cases.ts b/apps/sim/lib/projects/application/use-cases.ts index eb7484471d8..98cfb3495fc 100644 --- a/apps/sim/lib/projects/application/use-cases.ts +++ b/apps/sim/lib/projects/application/use-cases.ts @@ -1,22 +1,27 @@ import { AuditAction, AuditResourceType } from '@sim/audit' import { db } from '@sim/db' import { member, permissions, project, projectWorkspace, workspace } from '@sim/db/schema' -import { and, asc, eq, gt, isNull, sql } from 'drizzle-orm' +import { ORG_ADMIN_ROLES } from '@sim/platform-authz/workspace' +import { eq, inArray, sql } from 'drizzle-orm' import { recordProjectedUseCaseAuditEntries } from '@/lib/core/application/authorized-workspace-use-case' import type { OperationUseCase } from '@/lib/core/application/operation' import { OrchestrationError } from '@/lib/core/orchestration/types' -import { authorizeProject, requireProjectPrincipal } from '@/lib/projects/application/authorization' +import { + type AuthorizedProject, + authorizeProject, + authorizeProjectsForRead, + type ProjectAuthorizationInput, + requireProjectPrincipal, +} from '@/lib/projects/application/authorization' import { projectOperations } from '@/lib/projects/application/operations' import { archiveProjectInTransaction, finishProjectArchive } from '@/lib/projects/lifecycle' import { requireProjectApiEnabled } from '@/lib/projects/rollout.server' -interface ProjectInput { - projectId: string - organizationId?: string - workspaceId?: string -} -type ProjectContext = Awaited> -function presentProject(context: ProjectContext) { +type ProjectInput = ProjectAuthorizationInput & { projectId: string } +/** Reads see one consistent snapshot without locking the rows writers need. */ +const READ_SNAPSHOT = { isolationLevel: 'repeatable read', accessMode: 'read only' } as const + +function presentProject(context: AuthorizedProject) { return { ...context.record, environments: context.environments, @@ -32,14 +37,22 @@ export const getProject: OperationUseCase< operation: projectOperations.get, async execute({ principal, input }) { requireProjectPrincipal(principal, projectOperations.get) - requireProjectApiEnabled() - return db.transaction(async (tx) => ({ - project: presentProject(await authorizeProject(tx, principal, projectOperations.get, input)), - })) + await requireProjectApiEnabled() + return db.transaction( + async (tx) => ({ + project: presentProject( + await authorizeProject(tx, principal, projectOperations.get, input, 'snapshot') + ), + }), + READ_SNAPSHOT + ) }, } -/** Read-only capability probe. Issue mutations call authorizeProject inside their own transaction. */ +/** + * Read-only capability probe. Issue mutations call authorizeProject in `hold` mode + * inside their own transaction. + */ export const getProjectIssueAccess: OperationUseCase< typeof projectOperations.issues, ProjectInput, @@ -48,11 +61,17 @@ export const getProjectIssueAccess: OperationUseCase< operation: projectOperations.issues, async execute({ principal, input }) { requireProjectPrincipal(principal, projectOperations.issues) - requireProjectApiEnabled() + await requireProjectApiEnabled() return db.transaction(async (tx) => { - const context = await authorizeProject(tx, principal, projectOperations.issues, input) + const context = await authorizeProject( + tx, + principal, + projectOperations.issues, + input, + 'snapshot' + ) return { projectId: context.record.id } - }) + }, READ_SNAPSHOT) }, } @@ -64,43 +83,39 @@ export const listProjects: OperationUseCase< operation: projectOperations.list, async execute({ principal, input }) { requireProjectPrincipal(principal, projectOperations.list) - requireProjectApiEnabled() - if (!Number.isInteger(input.limit) || input.limit < 1 || input.limit > 100) - throw new OrchestrationError('validation', 'Limit must be between 1 and 100') + await requireProjectApiEnabled() return db.transaction(async (tx) => { - const candidates = await tx - .select({ id: project.id }) - .from(project) - .where( - and( - isNull(project.archivedAt), - input.organizationId ? eq(project.organizationId, input.organizationId) : undefined, - input.cursor ? gt(project.id, input.cursor) : undefined, - sql`exists (select 1 from ${projectWorkspace} pw join ${workspace} w on w.id = pw.workspace_id - where pw.project_id = ${project.id} and w.archived_at is null and ( - exists (select 1 from ${permissions} pe where pe.entity_type = 'workspace' and pe.entity_id = w.id and pe.user_id = ${principal.userId}) - or exists (select 1 from ${member} m where m.organization_id = w.organization_id and m.user_id = ${principal.userId} and m.role in ('owner', 'admin')) - ))` - ) - ) - .orderBy(asc(project.id)) - .limit(input.limit + 1) - const page = candidates.slice(0, input.limit) - const projects = [] - for (const row of page) - projects.push( - presentProject( - await authorizeProject(tx, principal, projectOperations.list, { - projectId: row.id, - organizationId: input.organizationId, - }) - ) + /** Driven from the caller's grants and admin organization, so cost tracks their reach. */ + const candidates = await tx.execute<{ id: string }>(sql` + WITH accessible AS ( + SELECT ${permissions.entityId} AS workspace_id FROM ${permissions} + WHERE ${permissions.userId} = ${principal.userId} + AND ${permissions.entityType} = 'workspace' + UNION + SELECT ${workspace.id} FROM ${member} + JOIN ${workspace} ON ${workspace.organizationId} = ${member.organizationId} + WHERE ${member.userId} = ${principal.userId} AND ${inArray(member.role, ORG_ADMIN_ROLES)} ) + SELECT DISTINCT ${projectWorkspace.projectId} AS id + FROM accessible + JOIN ${projectWorkspace} ON ${projectWorkspace.workspaceId} = accessible.workspace_id + JOIN ${workspace} + ON ${workspace.id} = accessible.workspace_id AND ${workspace.archivedAt} IS NULL + JOIN ${project} + ON ${project.id} = ${projectWorkspace.projectId} AND ${project.archivedAt} IS NULL + WHERE TRUE + ${input.organizationId ? sql`AND ${project.organizationId} = ${input.organizationId}` : sql``} + ${input.cursor ? sql`AND ${projectWorkspace.projectId} > ${input.cursor}` : sql``} + ORDER BY 1 + LIMIT ${input.limit + 1} + `) + const page = candidates.slice(0, input.limit).map((row) => row.id) + const projects = await authorizeProjectsForRead(tx, principal, page) return { - projects, - nextCursor: candidates.length > input.limit ? (page.at(-1)?.id ?? null) : null, + projects: projects.map(presentProject), + nextCursor: candidates.length > input.limit ? (page.at(-1) ?? null) : null, } - }) + }, READ_SNAPSHOT) }, } @@ -112,12 +127,10 @@ export const renameProject: OperationUseCase< operation: projectOperations.rename, async execute({ principal, input, request }) { requireProjectPrincipal(principal, projectOperations.rename) - requireProjectApiEnabled() - const name = input.name.trim() - if (!name || name.length > 100) - throw new OrchestrationError('validation', 'Project name must contain 1–100 characters') + await requireProjectApiEnabled() + const { name } = input const result = await db.transaction(async (tx) => { - const context = await authorizeProject(tx, principal, projectOperations.rename, input) + const context = await authorizeProject(tx, principal, projectOperations.rename, input, 'hold') if (context.record.archivedAt) throw new OrchestrationError('conflict', 'Project is archived') if (context.record.name === name) return { context, changed: false } await tx @@ -154,9 +167,15 @@ export const archiveProject: OperationUseCase< operation: projectOperations.archive, async execute({ principal, input, request }) { requireProjectPrincipal(principal, projectOperations.archive) - requireProjectApiEnabled() + await requireProjectApiEnabled() const result = await db.transaction(async (tx) => { - const context = await authorizeProject(tx, principal, projectOperations.archive, input) + const context = await authorizeProject( + tx, + principal, + projectOperations.archive, + input, + 'hold' + ) const effects = await archiveProjectInTransaction(tx, context.record.id) return { context, effects } }) @@ -190,7 +209,7 @@ export const getWorkspaceProject: OperationUseCase< operation: projectOperations.get, async execute({ principal, input }) { requireProjectPrincipal(principal, projectOperations.get) - requireProjectApiEnabled() + await requireProjectApiEnabled() return db.transaction(async (tx) => { const [membership] = await tx .select({ projectId: projectWorkspace.projectId }) @@ -200,12 +219,18 @@ export const getWorkspaceProject: OperationUseCase< if (!membership) throw new OrchestrationError('not_found', 'Project not found') return { project: presentProject( - await authorizeProject(tx, principal, projectOperations.get, { - projectId: membership.projectId, - workspaceId: input.workspaceId, - }) + await authorizeProject( + tx, + principal, + projectOperations.get, + { + projectId: membership.projectId, + workspaceId: input.workspaceId, + }, + 'snapshot' + ) ), } - }) + }, READ_SNAPSHOT) }, } diff --git a/apps/sim/lib/projects/create-input.ts b/apps/sim/lib/projects/create-input.ts index 4eb91f72df2..c342feb4d5c 100644 --- a/apps/sim/lib/projects/create-input.ts +++ b/apps/sim/lib/projects/create-input.ts @@ -1,7 +1,8 @@ import { z } from 'zod' +import { organizationIdSchema } from '@/lib/api/contracts/primitives' export const createProjectInputSchema = z.object({ - organizationId: z.string().trim().min(1).nullable(), + organizationId: z.string().trim().pipe(organizationIdSchema).nullable(), name: z.string().trim().min(1).max(100), initialEnvironment: z.object({ name: z.string().trim().min(1).max(100) }), }) diff --git a/apps/sim/lib/projects/lifecycle.ts b/apps/sim/lib/projects/lifecycle.ts index d089f7c2e38..2fe08df4890 100644 --- a/apps/sim/lib/projects/lifecycle.ts +++ b/apps/sim/lib/projects/lifecycle.ts @@ -1,66 +1,40 @@ -import { - project, - projectWorkspace, - workflow, - workflowMcpServer, - workflowMcpTool, - workspace, -} from '@sim/db/schema' -import { and, asc, eq, isNull } from 'drizzle-orm' +import { project, projectWorkspace, workspace } from '@sim/db/schema' +import { asc, eq } from 'drizzle-orm' import { OrchestrationError } from '@/lib/core/orchestration/types' import type { DbTransaction } from '@/lib/db/types' import { lockProject } from '@/lib/projects/membership' -import { archiveWorkflowInTransaction, finishWorkflowArchive } from '@/lib/workflows/lifecycle' -import { archiveWorkspaceInTransaction, finishWorkspaceArchive } from '@/lib/workspaces/lifecycle' +import { + archiveEnvironmentInTransaction, + type EnvironmentArchiveEffects, + finishEnvironmentArchive, +} from '@/lib/workspaces/lifecycle' /** All durable archive state commits together; external notifications follow the commit. */ -export async function archiveProjectInTransaction(tx: DbTransaction, projectId: string) { +export async function archiveProjectInTransaction( + tx: DbTransaction, + projectId: string +): Promise { await lockProject(tx, projectId) const [record] = await tx.select().from(project).where(eq(project.id, projectId)) if (!record) throw new OrchestrationError('not_found', 'Project not found') const now = record.archivedAt ?? new Date() - const members = await tx - .select({ id: projectWorkspace.workspaceId }) + const environments = await tx + .select({ id: workspace.id }) .from(projectWorkspace) + .innerJoin(workspace, eq(workspace.id, projectWorkspace.workspaceId)) .where(eq(projectWorkspace.projectId, projectId)) - .orderBy(asc(projectWorkspace.workspaceId)) - const workflows: { id: string; workspaceId: string; serverIds: string[] }[] = [] - const environments: { id: string; serverIds: string[] }[] = [] - for (const { id: workspaceId } of members) { - await tx - .select({ id: workspace.id }) - .from(workspace) - .where(eq(workspace.id, workspaceId)) - .for('update') - const rows = await tx - .select({ id: workflow.id }) - .from(workflow) - .where(and(eq(workflow.workspaceId, workspaceId), isNull(workflow.archivedAt))) - .orderBy(asc(workflow.id)) - for (const row of rows) { - const servers = await tx - .select({ id: workflowMcpTool.serverId }) - .from(workflowMcpTool) - .where(eq(workflowMcpTool.workflowId, row.id)) - await archiveWorkflowInTransaction(tx, row.id, now) - workflows.push({ id: row.id, workspaceId, serverIds: servers.map((server) => server.id) }) - } - const servers = await tx - .select({ id: workflowMcpServer.id }) - .from(workflowMcpServer) - .where(eq(workflowMcpServer.workspaceId, workspaceId)) - await archiveWorkspaceInTransaction(tx, workspaceId, now) - environments.push({ id: workspaceId, serverIds: servers.map((server) => server.id) }) - } + .orderBy(asc(workspace.id)) + .for('no key update', { of: workspace }) + const effects: EnvironmentArchiveEffects[] = [] + for (const { id } of environments) + effects.push(await archiveEnvironmentInTransaction(tx, id, now)) await tx.update(project).set({ archivedAt: now, updatedAt: now }).where(eq(project.id, projectId)) - return { workflows, environments } + return effects } export async function finishProjectArchive( - effects: Awaited>, + effects: EnvironmentArchiveEffects[], requestId: string ): Promise { - for (const row of effects.workflows) - await finishWorkflowArchive(row.id, row.workspaceId, row.serverIds, { requestId }) - for (const row of effects.environments) await finishWorkspaceArchive(row.id, row.serverIds) + for (const environment of effects) await finishEnvironmentArchive(environment, requestId) } diff --git a/apps/sim/lib/projects/membership.ts b/apps/sim/lib/projects/membership.ts index ed93c8059f0..536f95133a6 100644 --- a/apps/sim/lib/projects/membership.ts +++ b/apps/sim/lib/projects/membership.ts @@ -1,9 +1,72 @@ import { permissionGroup, project, projectWorkspace, workspace } from '@sim/db/schema' import { getPostgresErrorCode } from '@sim/utils/errors' import { generateId } from '@sim/utils/id' -import { and, asc, eq, inArray, isNull, ne, notInArray, sql } from 'drizzle-orm' +import { compareStrings, truncateAtCodePoint } from '@sim/utils/string' +import { and, asc, eq, inArray, isNull, notInArray, type SQL, sql } from 'drizzle-orm' import { OrchestrationError } from '@/lib/core/orchestration/types' -import type { DbOrTx, DbTransaction } from '@/lib/db/types' +import { + acquireAdvisoryXactLock, + acquireAdvisoryXactLocks, + tryAcquireAdvisoryXactLocks, +} from '@/lib/db/advisory-locks' +import { textArrayLiteral } from '@/lib/db/arrays' +import type { DbTransaction } from '@/lib/db/types' +import { acquirePermissionGroupOrgLock } from '@/lib/permission-groups/locks' + +const PROJECT_LOCK_TIMEOUT_MS = 5_000 + +/** A Project lifecycle rule or lock refused the change; callers may map it to their own error. */ +export class ProjectConflictError extends OrchestrationError { + constructor(message: string) { + super('conflict', message) + this.name = 'ProjectConflictError' + } +} + +/** + * Waits in `acquire` are bounded by {@link PROJECT_LOCK_TIMEOUT_MS} and a timeout or + * deadlock becomes a retryable Project conflict. The caller's own `lock_timeout` is + * restored afterwards, so the bound never leaks into the work done under the locks. + */ +async function withProjectLockTimeout( + tx: DbTransaction, + message: string, + acquire: () => Promise +): Promise { + const [setting] = await tx.execute<{ previous: string }>( + sql`SELECT current_setting('lock_timeout') AS previous` + ) + await tx.execute(sql`SELECT set_config('lock_timeout', ${`${PROJECT_LOCK_TIMEOUT_MS}ms`}, true)`) + let result: T + try { + result = await acquire() + } catch (error) { + const code = getPostgresErrorCode(error) + if (code === '55P03' || code === '40P01') throw new ProjectConflictError(message) + throw error + } + await tx.execute(sql`SELECT set_config('lock_timeout', ${setting?.previous ?? '0'}, true)`) + return result +} + +/** + * The advisory key the Project membership backfill holds exclusively per workspace while it + * assigns it; writers take it shared. Every holder locks in code-unit order of workspace id + * (`ORDER BY id COLLATE "C"` in SQL) so the two sides cannot deadlock. + */ +export function projectBackfillLockKey(workspaceId: string): string { + return `project-backfill:${workspaceId}` +} + +const BACKFILL_RUNNING = 'Project backfill is running; retry the operation' +const PROJECT_CHANGING = 'Project is changing; retry the operation' + +function acquireBackfillWriteLocks(tx: DbTransaction, workspaceIds: string[]) { + const locks = [...new Set(workspaceIds)] + .sort(compareStrings) + .map((id) => ({ key: projectBackfillLockKey(id), shared: true })) + return acquireAdvisoryXactLocks(tx, 'project_backfill', locks) +} /** Shared per-environment gate keeps membership absence reads stable during SQL backfill. */ export async function lockProjectBackfillWrites( @@ -11,39 +74,47 @@ export async function lockProjectBackfillWrites( workspaceIds: string[] ): Promise { if (!workspaceIds.length) return - await tx.execute(sql`SET LOCAL lock_timeout = '5s'`) - try { - await tx.execute(sql` - SELECT pg_advisory_xact_lock_shared(hashtextextended('project-backfill:' || id, 0)) - FROM (SELECT DISTINCT unnest(ARRAY[${sql.join( - workspaceIds.map((id) => sql`${id}`), - sql`, ` - )}]::text[]) AS id ORDER BY id) ids - `) - } catch (error) { - if (getPostgresErrorCode(error) === '55P03') - throw new OrchestrationError('conflict', 'Project backfill is running; retry the operation') - throw error - } + await withProjectLockTimeout(tx, BACKFILL_RUNNING, () => + acquireBackfillWriteLocks(tx, workspaceIds) + ) } /** Canonical Project mutex; membership and lifecycle writers hold it until commit. */ export async function lockProject(tx: DbTransaction, projectId: string): Promise { - await tx.execute(sql`SELECT set_config('lock_timeout', '5000ms', true)`) - try { - await tx.execute( - sql`SELECT pg_advisory_xact_lock(hashtextextended(${`project:${projectId}`}, 0))` - ) - } catch (error) { - if (getPostgresErrorCode(error) === '55P03') - throw new OrchestrationError('conflict', 'Project is changing; retry the operation') - throw error - } + await withProjectLockTimeout(tx, PROJECT_CHANGING, () => + acquireAdvisoryXactLock(tx, 'project', projectLockKey(projectId)) + ) +} + +/** Takes several Project mutexes in code-unit id order, the order every multi-lock holder uses. */ +export async function lockProjects(tx: DbTransaction, projectIds: string[]): Promise { + if (!projectIds.length) return + const locks = [...new Set(projectIds)] + .sort(compareStrings) + .map((id) => ({ key: projectLockKey(id), shared: false })) + await withProjectLockTimeout(tx, PROJECT_CHANGING, () => + acquireAdvisoryXactLocks(tx, 'project', locks) + ) +} + +function projectLockKey(projectId: string): string { + return `project:${projectId}` } function generatedProjectName(workspaceName: string): string { const suffix = ' - Project' - return `${(workspaceName.trim() || 'Untitled').slice(0, 100 - suffix.length)}${suffix}` + return `${truncateAtCodePoint(workspaceName.trim() || 'Untitled', 100 - suffix.length, '')}${suffix}` +} + +/** Selects `workspaceId` and every fork descendant, archived ones included, as `descendants`. */ +function forkSubtree(workspaceId: string): SQL { + return sql` + WITH RECURSIVE descendants AS ( + SELECT id, name, owner_id, archived_at FROM workspace WHERE id = ${workspaceId} + UNION + SELECT w.id, w.name, w.owner_id, w.archived_at + FROM workspace w JOIN descendants d ON w.forked_from_workspace_id = d.id + )` } export async function createProjectForWorkspace( @@ -53,7 +124,6 @@ export async function createProjectForWorkspace( name: string organizationId: string | null ownerId: string - archivedAt?: Date | null projectName?: string } ): Promise { @@ -63,7 +133,6 @@ export async function createProjectForWorkspace( name: input.projectName ?? generatedProjectName(input.name), organizationId: input.organizationId, ownerId: input.ownerId, - archivedAt: input.archivedAt ?? null, }) await tx.insert(projectWorkspace).values({ projectId: id, workspaceId: input.workspaceId }) return id @@ -71,95 +140,93 @@ export async function createProjectForWorkspace( /** Returns null only for a legacy workspace awaiting the SQL backfill. */ export async function lockWorkspaceProject(tx: DbTransaction, workspaceId: string) { - await lockProjectBackfillWrites(tx, [workspaceId]) - const [membership] = await tx - .select() - .from(projectWorkspace) - .where(eq(projectWorkspace.workspaceId, workspaceId)) - .limit(1) - if (!membership) return null - await lockProject(tx, membership.projectId) - const [current] = await tx - .select({ project }) - .from(projectWorkspace) - .innerJoin(project, eq(project.id, projectWorkspace.projectId)) - .where(eq(projectWorkspace.workspaceId, workspaceId)) - .limit(1) - if (!current || current.project.id !== membership.projectId) { - throw new OrchestrationError('conflict', 'Project membership changed; retry the operation') - } - return current.project + return withProjectLockTimeout(tx, PROJECT_CHANGING, async () => { + await acquireBackfillWriteLocks(tx, [workspaceId]) + const [membership] = await tx + .select() + .from(projectWorkspace) + .where(eq(projectWorkspace.workspaceId, workspaceId)) + .limit(1) + if (!membership) return null + await acquireAdvisoryXactLock(tx, 'project', projectLockKey(membership.projectId)) + const [current] = await tx + .select({ project }) + .from(projectWorkspace) + .innerJoin(project, eq(project.id, projectWorkspace.projectId)) + .where(eq(projectWorkspace.workspaceId, workspaceId)) + .limit(1) + if (!current || current.project.id !== membership.projectId) { + throw new ProjectConflictError('Project membership changed; retry the operation') + } + return current.project + }) } +/** + * Locks the parent's Project for a new fork; null means a legacy parent awaiting + * backfill whose subtree must still be unassigned. Refuses an archived Project. + */ export async function requireForkProject(tx: DbTransaction, parentWorkspaceId: string) { const parent = await lockWorkspaceProject(tx, parentWorkspaceId) if (!parent) { await requireUnassignedForkSubtree(tx, parentWorkspaceId) return null } - if (parent.archivedAt) throw new OrchestrationError('conflict', 'Cannot fork an archived Project') + if (parent.archivedAt) throw new ProjectConflictError('Cannot fork an archived Project') return parent } /** Legacy fallback must not hide partially assigned descendants. Caller holds the lineage lock. */ async function requireUnassignedForkSubtree(tx: DbTransaction, workspaceId: string): Promise { - const descendants = await tx.execute<{ id: string }>(sql` - WITH RECURSIVE descendants AS ( - SELECT id FROM workspace WHERE id = ${workspaceId} - UNION - SELECT w.id FROM workspace w JOIN descendants d ON w.forked_from_workspace_id = d.id - ) SELECT id FROM descendants - `) - if (!descendants.length) return - await lockProjectBackfillWrites( - tx, - descendants.map((row) => row.id) + const descendants = await tx.execute<{ id: string }>( + sql`${forkSubtree(workspaceId)} SELECT id FROM descendants` ) + if (!descendants.length) return + const ids = descendants.map((row) => row.id) + await lockProjectBackfillWrites(tx, ids) const rows = await tx .select({ id: projectWorkspace.workspaceId }) .from(projectWorkspace) - .where( - inArray( - projectWorkspace.workspaceId, - descendants.map((row) => row.id) - ) - ) + .where(inArray(projectWorkspace.workspaceId, ids)) .limit(1) if (rows.length) - throw new OrchestrationError( - 'conflict', - 'Fork descendants need Project membership reconciliation' - ) + throw new ProjectConflictError('Fork descendants need Project membership reconciliation') } -/** Individual removal cannot leave an active Project without an active environment. */ -export async function requireRemainingProjectEnvironment( +/** + * Keeps an active Project from outliving its environments: archiving its last active + * environment archives it too. Returns whether it did. Caller holds the Project lock. + */ +export async function archiveProjectWithLastEnvironment( tx: DbTransaction, - workspaceId: string -): Promise { - const owner = await lockWorkspaceProject(tx, workspaceId) - if (!owner) return - const [remaining] = await tx - .select({ id: workspace.id }) - .from(projectWorkspace) - .innerJoin(workspace, eq(workspace.id, projectWorkspace.workspaceId)) + projectId: string, + workspaceId: string, + now: Date +): Promise { + const archived = await tx + .update(project) + .set({ archivedAt: now, updatedAt: now }) .where( and( - eq(projectWorkspace.projectId, owner.id), - ne(workspace.id, workspaceId), - isNull(workspace.archivedAt) + eq(project.id, projectId), + isNull(project.archivedAt), + sql`NOT EXISTS ( + SELECT 1 FROM ${projectWorkspace} + JOIN ${workspace} ON ${workspace.id} = ${projectWorkspace.workspaceId} + WHERE ${projectWorkspace.projectId} = ${projectId} + AND ${workspace.id} <> ${workspaceId} + AND ${workspace.archivedAt} IS NULL + )` ) ) - .limit(1) - if (!remaining && !owner.archivedAt) { - throw new OrchestrationError( - 'conflict', - 'The last active environment cannot be removed. Archive the Project instead.' - ) - } + .returning({ id: project.id }) + return archived.length > 0 } -/** Called before clearing the edge, under the existing lineage lock. */ +/** + * Moves the detached subtree into a new Project and returns its id; null for a legacy + * unassigned subtree. Called before clearing the edge, under the existing lineage lock. + */ export async function splitForkProject( tx: DbTransaction, workspaceId: string @@ -169,29 +236,23 @@ export async function splitForkProject( await requireUnassignedForkSubtree(tx, workspaceId) return null } - if (owner.archivedAt) - throw new OrchestrationError('conflict', 'Cannot disconnect an archived Project') + if (owner.archivedAt) throw new ProjectConflictError('Cannot disconnect an archived Project') const rows = await tx.execute<{ id: string name: string owner_id: string archived_at: Date | null project_id: string | null - }>(sql` - WITH RECURSIVE descendants AS ( - SELECT id, name, owner_id, archived_at FROM workspace WHERE id = ${workspaceId} - UNION - SELECT w.id, w.name, w.owner_id, w.archived_at FROM workspace w JOIN descendants d ON w.forked_from_workspace_id = d.id - ) SELECT d.*, pw.project_id FROM descendants d LEFT JOIN project_workspace pw ON pw.workspace_id = d.id + }>(sql`${forkSubtree(workspaceId)} + SELECT d.*, pw.project_id FROM descendants d LEFT JOIN project_workspace pw ON pw.workspace_id = d.id `) if (rows.some((row) => row.project_id !== owner.id)) - throw new OrchestrationError( - 'conflict', + throw new ProjectConflictError( 'Fork descendants need Project membership reconciliation before disconnecting' ) const root = rows.find((row) => row.id === workspaceId) if (!root || rows.every((row) => row.archived_at)) - throw new OrchestrationError('conflict', 'A new Project needs an active environment') + throw new ProjectConflictError('A new Project needs an active environment') const ids = rows.map((row) => row.id) const [remaining] = await tx .select({ id: workspace.id }) @@ -206,8 +267,7 @@ export async function splitForkProject( ) .limit(1) if (!remaining) - throw new OrchestrationError( - 'conflict', + throw new ProjectConflictError( 'Disconnecting would remove the last active environment from this Project' ) const id = generateId() @@ -224,6 +284,7 @@ export async function splitForkProject( and(eq(projectWorkspace.projectId, owner.id), inArray(projectWorkspace.workspaceId, ids)) ) if (owner.organizationId) { + await acquirePermissionGroupOrgLock(tx, owner.organizationId) await tx.execute(sql` UPDATE ${permissionGroup} SET config = jsonb_set(config, '{deniedPartialAccessProjectIssues}', @@ -235,7 +296,11 @@ export async function splitForkProject( return id } -/** Ownership changes include the complete Project; never silently split Project-wide resources. */ +/** + * Ownership changes include the complete Project; never silently split Project-wide + * resources. Callers hold the organization mutation lock of every organization the + * Projects leave, which serializes the permission-group edit with group mutations. + */ export async function transferWorkspaceProjects( tx: DbTransaction, workspaceIds: string[], @@ -249,48 +314,83 @@ export async function transferWorkspaceProjects( .from(projectWorkspace) .where(inArray(projectWorkspace.workspaceId, workspaceIds)) .orderBy(asc(projectWorkspace.projectId)) + if (!owners.length) return + const projectIds = owners.map((row) => row.id) + await tryLockProjects(tx, projectIds) const selected = new Set(workspaceIds) - for (const owner of owners) { - await tryLockProject(tx, owner.id) - const members = await tx - .select({ id: projectWorkspace.workspaceId }) - .from(projectWorkspace) - .where(eq(projectWorkspace.projectId, owner.id)) - if (members.some((row) => !selected.has(row.id))) { - throw new OrchestrationError( - 'conflict', - 'Move all environments in the Project together, or disconnect the fork first' - ) - } - const [current] = await tx - .select({ organizationId: project.organizationId }) - .from(project) - .where(eq(project.id, owner.id)) - if (current?.organizationId && current.organizationId !== organizationId) { - await tx.execute(sql` - UPDATE ${permissionGroup} - SET config = jsonb_set(config, '{deniedPartialAccessProjectIssues}', - (config->'deniedPartialAccessProjectIssues') - ${owner.id}), updated_at = now() - WHERE organization_id = ${current.organizationId} - AND config->'deniedPartialAccessProjectIssues' ? ${owner.id} - `) - } - await tx - .update(project) - .set({ - organizationId, - ownerId, - updatedAt: new Date(), - }) - .where(eq(project.id, owner.id)) + const members = await tx + .select({ id: projectWorkspace.workspaceId }) + .from(projectWorkspace) + .where(inArray(projectWorkspace.projectId, projectIds)) + if (members.some((row) => !selected.has(row.id))) { + throw new ProjectConflictError( + 'Move all environments in the Project together, or disconnect the fork first' + ) + } + const current = await tx + .select({ id: project.id, organizationId: project.organizationId }) + .from(project) + .where(inArray(project.id, projectIds)) + const leavingByOrganization = new Map() + for (const row of current) { + if (!row.organizationId || row.organizationId === organizationId) continue + const leaving = leavingByOrganization.get(row.organizationId) + if (leaving) leaving.push(row.id) + else leavingByOrganization.set(row.organizationId, [row.id]) } + for (const [previousOrganizationId, leaving] of leavingByOrganization) { + const ids = textArrayLiteral(leaving) + await tx.execute(sql` + UPDATE ${permissionGroup} + SET config = jsonb_set(config, '{deniedPartialAccessProjectIssues}', + (config->'deniedPartialAccessProjectIssues') - ${ids}), updated_at = now() + WHERE organization_id = ${previousOrganizationId} + AND config->'deniedPartialAccessProjectIssues' ?| ${ids} + `) + } + await tx + .update(project) + .set({ organizationId, ownerId, updatedAt: new Date() }) + .where(inArray(project.id, projectIds)) } -/** Existing ownership paths can hold workspace rows first; refuse contention instead of inverting locks. */ -export async function tryLockProject(tx: DbOrTx, projectId: string): Promise { - const [lock] = await tx.execute<{ acquired: boolean }>( - sql`SELECT pg_try_advisory_xact_lock(hashtextextended(${`project:${projectId}`}, 0)) AS acquired` - ) - if (!lock?.acquired) - throw new OrchestrationError('conflict', 'Project is changing; retry the ownership change') +/** + * Moves the organization Projects `fromUserId` owns to `toUserId`, alongside the + * workspace ownership change that `workspaceIds` names. + */ +export async function reassignOrganizationProjects( + tx: DbTransaction, + input: { organizationId: string; fromUserId: string; toUserId: string; workspaceIds: string[] } +): Promise { + await lockProjectBackfillWrites(tx, input.workspaceIds) + const owned = await tx + .select({ id: project.id }) + .from(project) + .where( + and(eq(project.organizationId, input.organizationId), eq(project.ownerId, input.fromUserId)) + ) + .orderBy(asc(project.id)) + if (!owned.length) return + const ids = owned.map((row) => row.id) + await tryLockProjects(tx, ids) + await tx + .update(project) + .set({ ownerId: input.toUserId, updatedAt: new Date() }) + .where( + and( + inArray(project.id, ids), + eq(project.organizationId, input.organizationId), + eq(project.ownerId, input.fromUserId) + ) + ) +} + +/** + * Existing ownership paths can hold workspace rows first; refuse contention instead + * of inverting locks. + */ +async function tryLockProjects(tx: DbTransaction, projectIds: string[]): Promise { + const keys = projectIds.map(projectLockKey) + if (!(await tryAcquireAdvisoryXactLocks(tx, 'project', keys))) + throw new ProjectConflictError('Project is changing; retry the ownership change') } diff --git a/apps/sim/lib/projects/rollout.server.ts b/apps/sim/lib/projects/rollout.server.ts index 1f6e24962e4..758f5c0b963 100644 --- a/apps/sim/lib/projects/rollout.server.ts +++ b/apps/sim/lib/projects/rollout.server.ts @@ -1,4 +1,4 @@ -import { envBoolean, getEnv } from '@/lib/core/config/env' +import { isFeatureEnabled } from '@/lib/core/config/feature-flags' import { HttpError } from '@/lib/core/utils/http-error' class ProjectUnavailableError extends HttpError { @@ -8,6 +8,6 @@ class ProjectUnavailableError extends HttpError { } } -export function requireProjectApiEnabled(): void { - if (!(envBoolean(getEnv('PROJECT_API_ENABLED')) ?? false)) throw new ProjectUnavailableError() +export async function requireProjectApiEnabled(): Promise { + if (!(await isFeatureEnabled('projects'))) throw new ProjectUnavailableError() } diff --git a/apps/sim/lib/users/account-deletion.ts b/apps/sim/lib/users/account-deletion.ts index ce66288f5d4..f748152778f 100644 --- a/apps/sim/lib/users/account-deletion.ts +++ b/apps/sim/lib/users/account-deletion.ts @@ -15,6 +15,7 @@ import { workspace as workspaceTable, } from '@sim/db/schema' import { createLogger } from '@sim/logger' +import { getPostgresConstraintName, getPostgresErrorCode } from '@sim/utils/errors' import { formatQuotedNameList } from '@sim/utils/string' import { and, eq, gt, inArray, isNotNull, isNull, lte, ne, notExists, or, sql } from 'drizzle-orm' import type { @@ -817,7 +818,23 @@ export async function deleteUserAccount(userId: string): Promise { - await tx.delete(webhookPathClaim).where(eq(webhookPathClaim.workflowId, workflowId)) +export async function releaseWebhookPathClaims( + tx: DbOrTx, + workflowIds: readonly string[] +): Promise { + if (workflowIds.length === 0) return + await tx.delete(webhookPathClaim).where(inArray(webhookPathClaim.workflowId, workflowIds)) } /** diff --git a/apps/sim/lib/workflows/lifecycle.ts b/apps/sim/lib/workflows/lifecycle.ts index 67d51ff0095..0c78a3c3ed6 100644 --- a/apps/sim/lib/workflows/lifecycle.ts +++ b/apps/sim/lib/workflows/lifecycle.ts @@ -3,7 +3,6 @@ import { apiKey, chat, folder as folderTable, - projectWorkspace, webhook, workflow, workflowDeploymentVersion, @@ -109,7 +108,7 @@ export async function archiveWorkflow( .from(workflowMcpTool) .where(and(eq(workflowMcpTool.workflowId, workflowId), isNull(workflowMcpTool.archivedAt))) - await db.transaction((tx) => archiveWorkflowInTransaction(tx, workflowId, now)) + await db.transaction((tx) => archiveWorkflowsInTransaction(tx, [workflowId], now)) await finishWorkflowArchive( workflowId, @@ -300,33 +299,8 @@ export async function disableUserResources(userId: string): Promise { .from(workspace) .where(and(eq(workspace.ownerId, userId), isNull(workspace.archivedAt))) - const { archiveProjectInTransaction, finishProjectArchive } = await import( - '@/lib/projects/lifecycle' - ) - const { lockWorkspaceProject } = await import('@/lib/projects/membership') - const processed = new Set() for (const row of ownedWorkspaces) { - if (processed.has(row.id)) continue - const archived = await db.transaction(async (tx) => { - const record = await lockWorkspaceProject(tx, row.id) - if (!record) return null - const active = await tx - .select({ id: workspace.id, ownerId: workspace.ownerId }) - .from(projectWorkspace) - .innerJoin(workspace, eq(workspace.id, projectWorkspace.workspaceId)) - .where(and(eq(projectWorkspace.projectId, record.id), isNull(workspace.archivedAt))) - .orderBy(workspace.id) - .for('no key update', { of: workspace }) - if (!active.length || !active.every((entry) => entry.ownerId === userId)) return null - return archiveProjectInTransaction(tx, record.id) - }) - if (archived) { - for (const entry of archived.environments) processed.add(entry.id) - await finishProjectArchive(archived, requestId) - } else { - await archiveWorkspace(row.id, { requestId, expectedOwnerId: userId }) - processed.add(row.id) - } + await archiveWorkspace(row.id, { requestId, expectedOwnerId: userId }) } await db.delete(apiKey).where(eq(apiKey.userId, userId)) @@ -335,14 +309,18 @@ export async function disableUserResources(userId: string): Promise { ) } -/** Durable archive state shared by single-workflow and compound Project archival. */ -export async function archiveWorkflowInTransaction( +/** + * Durable archive state shared by single-workflow and compound Project archival. Each + * statement covers the whole batch, so the round trips do not grow with its size. + */ +export async function archiveWorkflowsInTransaction( tx: DbTransaction, - workflowId: string, + workflowIds: readonly string[], now: Date ): Promise { - await supersedeInFlightDeploymentOperations(tx, workflowId) - await releaseWebhookPathClaims(tx, workflowId) + if (workflowIds.length === 0) return + await supersedeInFlightDeploymentOperations(tx, workflowIds) + await releaseWebhookPathClaims(tx, workflowIds) await tx .update(workflowSchedule) @@ -353,7 +331,9 @@ export async function archiveWorkflowInTransaction( nextRunAt: null, lastQueuedAt: null, }) - .where(and(eq(workflowSchedule.workflowId, workflowId), isNull(workflowSchedule.archivedAt))) + .where( + and(inArray(workflowSchedule.workflowId, workflowIds), isNull(workflowSchedule.archivedAt)) + ) await tx .update(webhook) @@ -362,7 +342,7 @@ export async function archiveWorkflowInTransaction( updatedAt: now, isActive: false, }) - .where(and(eq(webhook.workflowId, workflowId), isNull(webhook.archivedAt))) + .where(and(inArray(webhook.workflowId, workflowIds), isNull(webhook.archivedAt))) await tx .update(chat) @@ -371,7 +351,7 @@ export async function archiveWorkflowInTransaction( updatedAt: now, isActive: false, }) - .where(and(eq(chat.workflowId, workflowId), isNull(chat.archivedAt))) + .where(and(inArray(chat.workflowId, workflowIds), isNull(chat.archivedAt))) await tx .update(workflowMcpTool) @@ -379,14 +359,21 @@ export async function archiveWorkflowInTransaction( archivedAt: now, updatedAt: now, }) - .where(and(eq(workflowMcpTool.workflowId, workflowId), isNull(workflowMcpTool.archivedAt))) + .where( + and(inArray(workflowMcpTool.workflowId, workflowIds), isNull(workflowMcpTool.archivedAt)) + ) await tx .update(workflowDeploymentVersion) .set({ isActive: false, }) - .where(eq(workflowDeploymentVersion.workflowId, workflowId)) + .where( + and( + inArray(workflowDeploymentVersion.workflowId, workflowIds), + eq(workflowDeploymentVersion.isActive, true) + ) + ) await tx .update(workflow) @@ -396,7 +383,7 @@ export async function archiveWorkflowInTransaction( isDeployed: false, isPublicApi: false, }) - .where(and(eq(workflow.id, workflowId), isNull(workflow.archivedAt))) + .where(and(inArray(workflow.id, workflowIds), isNull(workflow.archivedAt))) } /** Best-effort external notifications run only after durable archive state commits. */ @@ -419,13 +406,9 @@ export async function finishWorkflowArchive( await cleanupExternalWebhooksForWorkflow(workflowId, options.requestId) - if (workspaceId && mcpPubSub && serverIds.length > 0) { - const uniqueServerIds = [...new Set(serverIds)] - for (const serverId of uniqueServerIds) { - mcpPubSub.publishWorkflowToolsChanged({ - serverId, - workspaceId: workspaceId, - }) + if (workspaceId && mcpPubSub) { + for (const serverId of new Set(serverIds)) { + mcpPubSub.publishWorkflowToolsChanged({ serverId, workspaceId }) } } } diff --git a/apps/sim/lib/workflows/persistence/deployment-operations.ts b/apps/sim/lib/workflows/persistence/deployment-operations.ts index 6bac6ef19fa..4d1ffca15c1 100644 --- a/apps/sim/lib/workflows/persistence/deployment-operations.ts +++ b/apps/sim/lib/workflows/persistence/deployment-operations.ts @@ -716,14 +716,15 @@ export async function recordDeploymentOperationRetry( } /** - * Supersedes every in-flight operation for a workflow. Must run inside the + * Supersedes every in-flight operation for `workflowIds`. Must run inside the * undeploy/archive transaction so a queued preparation cannot activate a * version after the user explicitly took the workflow offline. */ export async function supersedeInFlightDeploymentOperations( executor: DbOrTx, - workflowId: string + workflowIds: readonly string[] ): Promise { + if (workflowIds.length === 0) return const now = new Date() await executor .update(workflowDeploymentOperation) @@ -734,7 +735,7 @@ export async function supersedeInFlightDeploymentOperations( }) .where( and( - eq(workflowDeploymentOperation.workflowId, workflowId), + inArray(workflowDeploymentOperation.workflowId, workflowIds), inArray(workflowDeploymentOperation.status, IN_FLIGHT_STATUSES) ) ) diff --git a/apps/sim/lib/workflows/persistence/duplicate.ts b/apps/sim/lib/workflows/persistence/duplicate.ts index 09d486204a6..f05c08fc4fb 100644 --- a/apps/sim/lib/workflows/persistence/duplicate.ts +++ b/apps/sim/lib/workflows/persistence/duplicate.ts @@ -20,7 +20,7 @@ import { import { and, eq } from 'drizzle-orm' import type { DbOrTx, DbTransaction } from '@/lib/db/types' import { remapConditionEdgeHandle } from '@/lib/workflows/condition-ids' -import { buildNewWorkflowRow } from '@/lib/workflows/persistence/new-workflow-row' +import { insertNewWorkflowRow } from '@/lib/workflows/persistence/new-workflow-row' import { remapConditionIdsInSubBlocks, remapVariableIdsInSubBlocks, @@ -211,19 +211,17 @@ export async function duplicateWorkflow( // A duplicate is a new workflow, so it takes the workspace's fork-sync policy rather // than inheriting the source's participation, and starts unlocked like any new one. - await tx.insert(workflow).values( - await buildNewWorkflowRow(tx, { - id: newWorkflowId, - userId, - workspaceId: targetWorkspaceId, - folderId: targetFolderId, - sortOrder, - name: deduplicatedName, - description: description || source.description, - now, - variables, - }) - ) + await insertNewWorkflowRow(tx, { + id: newWorkflowId, + userId, + workspaceId: targetWorkspaceId, + folderId: targetFolderId, + sortOrder, + name: deduplicatedName, + description: description || source.description, + now, + variables, + }) // Copy all blocks from source workflow with new IDs const sourceBlocks = await tx diff --git a/apps/sim/lib/workflows/persistence/new-workflow-row.ts b/apps/sim/lib/workflows/persistence/new-workflow-row.ts index 4b7324a8611..6c820d2eb5c 100644 --- a/apps/sim/lib/workflows/persistence/new-workflow-row.ts +++ b/apps/sim/lib/workflows/persistence/new-workflow-row.ts @@ -1,6 +1,5 @@ -import { type workflow, workspace } from '@sim/db/schema' -import { and, eq, isNull } from 'drizzle-orm' -import type { DbOrTx, DbTransaction } from '@/lib/db/types' +import { workflow } from '@sim/db/schema' +import type { DbTransaction } from '@/lib/db/types' import { lockActiveWorkspace } from '@/lib/workspaces/active-workspace' interface NewWorkflowRowInput { @@ -15,33 +14,14 @@ interface NewWorkflowRowInput { now?: Date } -/** - * The workspace's `forkSyncNewWorkflowsExcluded` policy: whether a workflow created now - * starts outside fork sync. - * - * `false` for an archived or missing workspace: a wrongly-synced workflow is visible and - * fixable in the Forks list, while a wrongly-excluded one silently stops syncing. - */ -export async function readForkSyncNewWorkflowsExcluded( - executor: DbOrTx, - workspaceId: string -): Promise { - const [row] = await executor - .select({ excluded: workspace.forkSyncNewWorkflowsExcluded }) - .from(workspace) - .where(and(eq(workspace.id, workspaceId), isNull(workspace.archivedAt))) - .limit(1) - return row?.excluded ?? false -} - /** * The insert row for a genuinely new workflow - created, duplicated, imported, or seeded as * a starter - so every such path takes the workspace's fork-sync policy rather than the * column default. A fork or promote copy is not new and is written by - * `copyWorkflowStateIntoTarget` instead. + * `copyWorkflowStateIntoTarget` instead. Share-locks the workspace and refuses an archived one. */ export async function buildNewWorkflowRow(executor: DbTransaction, input: NewWorkflowRowInput) { - const workspace = await lockActiveWorkspace(executor, input.workspaceId) + const target = await lockActiveWorkspace(executor, input.workspaceId) const now = input.now ?? new Date() return { id: input.id, @@ -57,6 +37,10 @@ export async function buildNewWorkflowRow(executor: DbTransaction, input: NewWor isDeployed: false, runCount: 0, variables: input.variables ?? {}, - forkSyncExcluded: workspace.forkSyncNewWorkflowsExcluded, + forkSyncExcluded: target.forkSyncNewWorkflowsExcluded, } satisfies typeof workflow.$inferInsert } + +export async function insertNewWorkflowRow(tx: DbTransaction, input: NewWorkflowRowInput) { + await tx.insert(workflow).values(await buildNewWorkflowRow(tx, input)) +} diff --git a/apps/sim/lib/workflows/persistence/utils.ts b/apps/sim/lib/workflows/persistence/utils.ts index 3d9c6a14ee5..f2be891b240 100644 --- a/apps/sim/lib/workflows/persistence/utils.ts +++ b/apps/sim/lib/workflows/persistence/utils.ts @@ -972,10 +972,10 @@ export async function undeployWorkflow(params: { .where(eq(workflowDeploymentVersion.workflowId, workflowId)) const deploymentVersionIds = deploymentVersions.map((version) => version.id) - await supersedeInFlightDeploymentOperations(dbCtx, workflowId) + await supersedeInFlightDeploymentOperations(dbCtx, [workflowId]) const { deleteSchedulesForWorkflow } = await import('@/lib/workflows/schedules/deploy') await deleteSchedulesForWorkflow(workflowId, dbCtx) - await releaseWebhookPathClaims(dbCtx, workflowId) + await releaseWebhookPathClaims(dbCtx, [workflowId]) await dbCtx .update(workflowDeploymentVersion) diff --git a/apps/sim/lib/workspaces/active-workspace.ts b/apps/sim/lib/workspaces/active-workspace.ts index f506aa02ae8..8edb9856444 100644 --- a/apps/sim/lib/workspaces/active-workspace.ts +++ b/apps/sim/lib/workspaces/active-workspace.ts @@ -1,11 +1,14 @@ import { workspace } from '@sim/db/schema' import { eq } from 'drizzle-orm' import { OrchestrationError } from '@/lib/core/orchestration/types' -import type { DbOrTx } from '@/lib/db/types' +import type { DbTransaction } from '@/lib/db/types' -/** Transactional resource creation holds this row through insertion so archival cannot overtake it. */ -export async function lockActiveWorkspace(executor: DbOrTx, workspaceId: string) { - const [record] = await executor +/** + * Transactional resource creation and restore hold this row through the write so + * archival cannot overtake it. + */ +export async function lockActiveWorkspace(tx: DbTransaction, workspaceId: string) { + const [record] = await tx .select({ archivedAt: workspace.archivedAt, forkSyncNewWorkflowsExcluded: workspace.forkSyncNewWorkflowsExcluded, diff --git a/apps/sim/lib/workspaces/admin-move.ts b/apps/sim/lib/workspaces/admin-move.ts index 9ad81035fd1..f0a98b5773d 100644 --- a/apps/sim/lib/workspaces/admin-move.ts +++ b/apps/sim/lib/workspaces/admin-move.ts @@ -30,7 +30,6 @@ import { planHasFixedSeatCap, resolveSeatCapacity, } from '@/lib/billing/validation/seat-management' -import { OrchestrationError } from '@/lib/core/orchestration/types' import { addOutboxEventSourceOperationId, enqueueOrReschedulePendingOutboxEvent, @@ -42,7 +41,7 @@ import type { DbOrTx } from '@/lib/db/types' import { getInvitationById, isInvitationExpired } from '@/lib/invitations/core' import { acquireInvitationMutationLocks } from '@/lib/invitations/locks' import { PENDING_INVITATION_UNIQUE_INDEX, sendInvitationEmail } from '@/lib/invitations/send' -import { transferWorkspaceProjects } from '@/lib/projects/membership' +import { ProjectConflictError, transferWorkspaceProjects } from '@/lib/projects/membership' import { invalidateWorkspaceTableLimitsCache } from '@/lib/table/billing' import { deleteCustomBlock } from '@/lib/workflows/custom-blocks/operations' import { @@ -1501,7 +1500,7 @@ export async function moveWorkspaceToOrganization(params: { }) break } catch (error) { - if (error instanceof OrchestrationError && error.code === 'conflict') { + if (error instanceof ProjectConflictError) { throw new WorkspaceMoveError(error.message, 'project-conflict') } if (error instanceof InvitationSetChangedError) { diff --git a/apps/sim/lib/workspaces/create.ts b/apps/sim/lib/workspaces/create.ts index 4b1549bd0ae..b764d7982bd 100644 --- a/apps/sim/lib/workspaces/create.ts +++ b/apps/sim/lib/workspaces/create.ts @@ -1,14 +1,13 @@ import { db } from '@sim/db' -import { permissions, type WorkspaceMode, workflow, workspace } from '@sim/db/schema' +import { permissions, type WorkspaceMode, workspace } from '@sim/db/schema' 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 { DbTransaction } from '@/lib/db/types' import { createProjectForWorkspace } from '@/lib/projects/membership' -import { requireProjectApiEnabled } from '@/lib/projects/rollout.server' import { buildDefaultWorkflowArtifacts } from '@/lib/workflows/defaults' -import { buildNewWorkflowRow } from '@/lib/workflows/persistence/new-workflow-row' +import { insertNewWorkflowRow } from '@/lib/workflows/persistence/new-workflow-row' import { saveWorkflowToNormalizedTables } from '@/lib/workflows/persistence/utils' import { getWorkspaceInvitePolicy, @@ -20,7 +19,7 @@ import { } from '@/lib/workspaces/policy' /** Foreign keys from `workspace` to `user`; a violation means the acting user's row is gone. */ -const WORKSPACE_USER_FK_CONSTRAINTS = new Set([ +export const WORKSPACE_USER_FK_CONSTRAINTS = new Set([ 'workspace_owner_id_user_id_fk', 'workspace_billed_account_user_id_user_id_fk', ]) @@ -88,9 +87,11 @@ export interface TransactionalCreateWorkspaceParams extends CreateWorkspaceParam * The caller supplies the creation-policy snapshot. This function revalidates * that snapshot — including the `workspace.create` capability under the * permission-group advisory lock — before inserting the workspace, owner - * permission and optional starter workflow atomically. + * permission, its Project and optional starter workflow atomically. A caller + * creating an explicit Project names it; otherwise the name derives from the + * workspace. */ -async function createWorkspaceRecordsInTransaction( +export async function createWorkspaceWithProjectInTransaction( tx: DbTransaction, { projectName, @@ -165,17 +166,15 @@ async function createWorkspaceRecordsInTransaction( await tx.insert(permissions).values(permissionRows) if (defaultWorkflowArtifacts) { - await tx.insert(workflow).values( - await buildNewWorkflowRow(tx, { - id: workflowId, - userId, - workspaceId, - folderId: null, - name: 'default-agent', - description: 'Your first workflow - start building here!', - now, - }) - ) + await insertNewWorkflowRow(tx, { + id: workflowId, + userId, + workspaceId, + folderId: null, + name: 'default-agent', + description: 'Your first workflow - start building here!', + now, + }) await saveWorkflowToNormalizedTables( workflowId, defaultWorkflowArtifacts.workflowState, @@ -204,21 +203,11 @@ async function createWorkspaceRecordsInTransaction( } } -/** Explicit Project creation always commits its first environment in the same transaction. */ -export async function createWorkspaceWithProjectInTransaction( - tx: DbTransaction, - params: TransactionalCreateWorkspaceParams & { projectName?: string } -): Promise<{ projectId: string; workspace: CreatedWorkspace }> { - requireProjectApiEnabled() - return createWorkspaceRecordsInTransaction(tx, params) -} - -/** Preserves the workspace-only result for existing creation callers. */ export async function createWorkspaceInTransaction( tx: DbTransaction, params: TransactionalCreateWorkspaceParams ): Promise { - return (await createWorkspaceRecordsInTransaction(tx, params)).workspace + return (await createWorkspaceWithProjectInTransaction(tx, params)).workspace } /** Creates a workspace through the canonical lock-and-insert transaction. */ diff --git a/apps/sim/lib/workspaces/lifecycle.ts b/apps/sim/lib/workspaces/lifecycle.ts index 08fa60e21c9..d6c88cb377c 100644 --- a/apps/sim/lib/workspaces/lifecycle.ts +++ b/apps/sim/lib/workspaces/lifecycle.ts @@ -8,30 +8,53 @@ import { knowledgeConnector, mcpServers, userTableDefinitions, + workflow, workflowMcpServer, + workflowMcpTool, workspace, workspaceFiles, } from '@sim/db/schema' import { createLogger } from '@sim/logger' -import { and, eq, inArray, isNull, sql } from 'drizzle-orm' +import { chunkArray } from '@sim/utils/helpers' +import { and, asc, eq, inArray, isNull, sql } from 'drizzle-orm' +import { mapWithConcurrency } from '@/lib/core/utils/concurrency' import type { DbTransaction } from '@/lib/db/types' import { mcpPubSub } from '@/lib/mcp/pubsub' import { mcpService } from '@/lib/mcp/service' -import { lockWorkspaceProject, requireRemainingProjectEnvironment } from '@/lib/projects/membership' -import { archiveWorkflowsForWorkspace } from '@/lib/workflows/lifecycle' +import { archiveProjectWithLastEnvironment, lockWorkspaceProject } from '@/lib/projects/membership' +import { archiveWorkflowsInTransaction, finishWorkflowArchive } from '@/lib/workflows/lifecycle' import { getWorkspaceWithOwner } from '@/lib/workspaces/permissions/utils' const logger = createLogger('WorkspaceLifecycle') +/** Bounds each batched workflow archive statement's parameter list. */ +const WORKFLOW_ARCHIVE_BATCH_SIZE = 1_000 +/** Bounds concurrent post-commit notifications; each is best-effort and catches its own errors. */ +const ARCHIVE_NOTIFICATION_CONCURRENCY = 8 + +/** What an environment archive must announce once its transaction commits. */ +export interface EnvironmentArchiveEffects { + workspaceId: string + workflows: { id: string; serverIds: string[] }[] + serverIds: string[] +} + interface ArchiveWorkspaceOptions { requestId: string expectedOwnerId?: string } +interface ArchiveWorkspaceResult { + archived: boolean + workspaceName?: string + /** The Project archived with its last active environment, for the caller's audit. */ + archivedProject?: { id: string; name: string } +} + export async function archiveWorkspace( workspaceId: string, options: ArchiveWorkspaceOptions -): Promise<{ archived: boolean; workspaceName?: string }> { +): Promise { const workspaceRecord = await getWorkspaceWithOwner(workspaceId, { includeArchived: true }) if (!workspaceRecord) { @@ -40,48 +63,112 @@ export async function archiveWorkspace( /** Retrying deletion also archives children left active by an older or incomplete deletion. */ const now = workspaceRecord.archivedAt ?? new Date() - const workflowMcpServerIds = await db - .select({ id: workflowMcpServer.id }) - .from(workflowMcpServer) - .where(eq(workflowMcpServer.workspaceId, workspaceId)) - const archived = await db.transaction(async (tx) => { - if (options.expectedOwnerId) { - await lockWorkspaceProject(tx, workspaceId) - const [current] = await tx - .select({ ownerId: workspace.ownerId }) - .from(workspace) - .where(eq(workspace.id, workspaceId)) - .for('no key update') - if (!current || current.ownerId !== options.expectedOwnerId) return false + const outcome = await db.transaction(async (tx) => { + const owningProject = await lockWorkspaceProject(tx, workspaceId) + /** Waits out in-flight workflow creation and restore, so their rows are archived too. */ + const [current] = await tx + .select({ ownerId: workspace.ownerId }) + .from(workspace) + .where(eq(workspace.id, workspaceId)) + .for('no key update') + if (!current) return null + if (options.expectedOwnerId && current.ownerId !== options.expectedOwnerId) return null + const projectArchived = + owningProject !== null && + (await archiveProjectWithLastEnvironment(tx, owningProject.id, workspaceId, now)) + return { + effects: await archiveEnvironmentInTransaction(tx, workspaceId, now), + archivedProject: projectArchived + ? { id: owningProject.id, name: owningProject.name } + : undefined, } - await requireRemainingProjectEnvironment(tx, workspaceId) - await archiveWorkspaceInTransaction(tx, workspaceId, now) - return true }) - if (!archived) return { archived: false } - - await archiveWorkflowsForWorkspace(workspaceId, options) + if (!outcome) return { archived: false } logger.info(`[${options.requestId}] Archived workspace ${workspaceId}`) - await finishWorkspaceArchive( - workspaceId, - workflowMcpServerIds.map((server) => server.id) - ) + await finishEnvironmentArchive(outcome.effects, options.requestId) return { archived: !workspaceRecord.archivedAt, workspaceName: workspaceRecord.name, + archivedProject: outcome.archivedProject, } } -/** Durable environment archive changes; callers own Project minimum-environment checks. */ -export async function archiveWorkspaceInTransaction( +/** + * Archives an environment and its active workflows in the caller's transaction, batching + * the workflow statements; callers own the Project lifecycle. Announce the returned + * effects with {@link finishEnvironmentArchive} after commit. + */ +export async function archiveEnvironmentInTransaction( tx: DbTransaction, workspaceId: string, now: Date +): Promise { + const workflows: EnvironmentArchiveEffects['workflows'] = [] + const active = await tx + .select({ id: workflow.id }) + .from(workflow) + .where(and(eq(workflow.workspaceId, workspaceId), isNull(workflow.archivedAt))) + .orderBy(asc(workflow.id)) + for (const batch of chunkArray( + active.map((row) => row.id), + WORKFLOW_ARCHIVE_BATCH_SIZE + )) { + const tools = await tx + .select({ workflowId: workflowMcpTool.workflowId, serverId: workflowMcpTool.serverId }) + .from(workflowMcpTool) + .where(and(inArray(workflowMcpTool.workflowId, batch), isNull(workflowMcpTool.archivedAt))) + const serverIdsByWorkflow = new Map() + for (const tool of tools) { + const serverIds = serverIdsByWorkflow.get(tool.workflowId) + if (serverIds) serverIds.push(tool.serverId) + else serverIdsByWorkflow.set(tool.workflowId, [tool.serverId]) + } + await archiveWorkflowsInTransaction(tx, batch, now) + for (const id of batch) workflows.push({ id, serverIds: serverIdsByWorkflow.get(id) ?? [] }) + } + const serverIds = await archiveWorkspaceRecordsInTransaction(tx, workspaceId, now) + return { workspaceId, workflows, serverIds } +} + +/** + * Announces a committed environment archive. Every step is best-effort and isolated, so + * one failed notification never skips the rest. + */ +export async function finishEnvironmentArchive( + effects: EnvironmentArchiveEffects, + requestId: string ): Promise { + const { workspaceId } = effects + await mapWithConcurrency(effects.workflows, ARCHIVE_NOTIFICATION_CONCURRENCY, (row) => + finishWorkflowArchive(row.id, workspaceId, row.serverIds, { requestId }).catch((error) => + logger.warn(`[${requestId}] Post-archive notification failed for workflow ${row.id}`, { + error, + }) + ) + ) + await mcpService.clearCache(workspaceId).catch(() => undefined) + if (!mcpPubSub) return + for (const serverId of effects.serverIds) { + try { + mcpPubSub.publishWorkflowToolsChanged({ serverId, workspaceId }) + } catch (error) { + logger.warn(`[${requestId}] MCP tools-changed publish failed for server ${serverId}`, { + error, + }) + } + } +} + +/** The workspace row and its non-workflow resources; returns the deployed MCP server ids. */ +async function archiveWorkspaceRecordsInTransaction( + tx: DbTransaction, + workspaceId: string, + now: Date +): Promise { await tx .update(knowledgeBase) .set({ @@ -161,6 +248,11 @@ export async function archiveWorkspaceInTransaction( .delete(apiKey) .where(and(eq(apiKey.workspaceId, workspaceId), eq(apiKey.type, 'workspace'))) + /** Every server is announced, so a retry still invalidates; only live ones are stamped. */ + const servers = await tx + .select({ id: workflowMcpServer.id }) + .from(workflowMcpServer) + .where(eq(workflowMcpServer.workspaceId, workspaceId)) await tx .update(workflowMcpServer) .set({ @@ -168,7 +260,7 @@ export async function archiveWorkspaceInTransaction( isPublic: false, updatedAt: now, }) - .where(eq(workflowMcpServer.workspaceId, workspaceId)) + .where(and(eq(workflowMcpServer.workspaceId, workspaceId), isNull(workflowMcpServer.deletedAt))) await tx .update(mcpServers) @@ -186,16 +278,5 @@ export async function archiveWorkspaceInTransaction( updatedAt: now, }) .where(and(eq(workspace.id, workspaceId), isNull(workspace.archivedAt))) -} - -/** Refreshes derived MCP state after the archive transaction commits. */ -export async function finishWorkspaceArchive( - workspaceId: string, - serverIds: string[] -): Promise { - await mcpService.clearCache(workspaceId).catch(() => undefined) - if (mcpPubSub) { - for (const serverId of serverIds) - mcpPubSub.publishWorkflowToolsChanged({ serverId, workspaceId }) - } + return servers.map((server) => server.id) } diff --git a/knip.jsonc b/knip.jsonc index e68838874c9..68c63368a0e 100644 --- a/knip.jsonc +++ b/knip.jsonc @@ -45,8 +45,6 @@ // Generated contracts mirror their source of truth; regenerating them must not // trip the unused-export ratchet, and hand edits would be overwritten. "lib/mothership/generated/**": ["exports", "types", "duplicates"], - // API conventions require exported named wire schemas and aliases, including before client adoption. - "lib/api/contracts/projects.ts": ["exports", "types"], "sandbox-tasks/index.ts": ["files"], "components/mcp/index.ts": ["files"], "triggers/quickbooks/index.ts": ["files"], diff --git a/packages/audit/src/types.ts b/packages/audit/src/types.ts index 06547294cf5..1bbeeade26d 100644 --- a/packages/audit/src/types.ts +++ b/packages/audit/src/types.ts @@ -177,6 +177,10 @@ export const AuditAction = { PERMISSION_ACCESS_REQUEST_CLOSED: 'permission_access_request.closed', PERMISSION_ACCESS_REQUEST_SETTINGS_CHANGED: 'permission_access_request.settings_changed', + PROJECT_CREATED: 'project.created', + PROJECT_UPDATED: 'project.updated', + PROJECT_ARCHIVED: 'project.archived', + SANDBOX_CREATED: 'sandbox.created', SANDBOX_UPDATED: 'sandbox.updated', SANDBOX_DELETED: 'sandbox.deleted', @@ -218,9 +222,6 @@ export const AuditAction = { WORKFLOW_EXPORTED: 'workflow.exported', WORKSPACE_CREATED: 'workspace.created', - PROJECT_CREATED: 'project.created', - PROJECT_UPDATED: 'project.updated', - PROJECT_ARCHIVED: 'project.archived', WORKSPACE_UPDATED: 'workspace.updated', WORKSPACE_DELETED: 'workspace.deleted', WORKSPACE_DUPLICATED: 'workspace.duplicated', @@ -278,6 +279,7 @@ export const AuditResourceType = { PASSWORD: 'password', PERMISSION_GROUP: 'permission_group', PERMISSION_ACCESS_REQUEST: 'permission_access_request', + PROJECT: 'project', SANDBOX: 'sandbox', SCHEDULE: 'schedule', SCIM_CONNECTION: 'scim_connection', @@ -290,7 +292,6 @@ export const AuditResourceType = { USER: 'user', WEBHOOK: 'webhook', WORKFLOW: 'workflow', - PROJECT: 'project', WORKSPACE: 'workspace', } as const diff --git a/packages/testing/src/mocks/audit.mock.ts b/packages/testing/src/mocks/audit.mock.ts index 8a89a318ebc..2df4fc55132 100644 --- a/packages/testing/src/mocks/audit.mock.ts +++ b/packages/testing/src/mocks/audit.mock.ts @@ -149,6 +149,10 @@ const AuditAction = { PERMISSION_ACCESS_REQUEST_CANCELLED: 'permission_access_request.cancelled', PERMISSION_ACCESS_REQUEST_CLOSED: 'permission_access_request.closed', PERMISSION_ACCESS_REQUEST_SETTINGS_CHANGED: 'permission_access_request.settings_changed', + + PROJECT_CREATED: 'project.created', + PROJECT_UPDATED: 'project.updated', + PROJECT_ARCHIVED: 'project.archived', SANDBOX_CREATED: 'sandbox.created', SANDBOX_UPDATED: 'sandbox.updated', SANDBOX_DELETED: 'sandbox.deleted', @@ -238,6 +242,7 @@ const AuditResourceType = { PASSWORD: 'password', PERMISSION_GROUP: 'permission_group', PERMISSION_ACCESS_REQUEST: 'permission_access_request', + PROJECT: 'project', SANDBOX: 'sandbox', SCHEDULE: 'schedule', SCIM_CONNECTION: 'scim_connection',