Files
WeHub Mirror 6bf8bebf51
CI / Test and Build (push) Failing after 1s
CI / Migrate Dev DB (push) Has been skipped
CI / Migrate DB (push) Has been skipped
CodeQL / Analyze actions (push) Has been cancelled
CodeQL / Analyze javascript-typescript (push) Has been cancelled
CI / Detect Version (push) Has been cancelled
CI / Detect Desktop Changes (push) Has been cancelled
CI / Build AMD64 (blacksmith-2vcpu-ubuntu-2404, ./docker/cron.Dockerfile, ubuntu-latest, ghcr.io/simstudioai/cron) (push) Has been cancelled
CI / Build AMD64 (blacksmith-2vcpu-ubuntu-2404, ./docker/db.Dockerfile, ECR_MIGRATIONS, ubuntu-latest, ghcr.io/simstudioai/migrations) (push) Has been cancelled
CI / Build AMD64 (blacksmith-4vcpu-ubuntu-2404, ./docker/pii.Dockerfile, ECR_PII, ubuntu-latest, ghcr.io/simstudioai/pii) (push) Has been cancelled
CI / Build AMD64 (blacksmith-4vcpu-ubuntu-2404, ./docker/realtime.Dockerfile, ECR_REALTIME, ubuntu-latest, ghcr.io/simstudioai/realtime) (push) Has been cancelled
CI / Build AMD64 (blacksmith-8vcpu-ubuntu-2404, ./docker/app.Dockerfile, ECR_APP, linux-x64-8-core, ghcr.io/simstudioai/simstudio) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-4vcpu-ubuntu-2404-arm, ./docker/cron.Dockerfile, ubuntu-24.04-arm, ghcr.io/simstudioai/cron) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-4vcpu-ubuntu-2404-arm, ./docker/db.Dockerfile, ubuntu-24.04-arm, ghcr.io/simstudioai/migrations) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-4vcpu-ubuntu-2404-arm, ./docker/pii.Dockerfile, ubuntu-24.04-arm, ghcr.io/simstudioai/pii) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-4vcpu-ubuntu-2404-arm, ./docker/realtime.Dockerfile, ubuntu-24.04-arm, ghcr.io/simstudioai/realtime) (push) Has been cancelled
CI / Build ARM64 (GHCR Only) (blacksmith-8vcpu-ubuntu-2404-arm, ./docker/app.Dockerfile, linux-arm64-8-core, ghcr.io/simstudioai/simstudio) (push) Has been cancelled
CI / Check Docs Changes (push) Has been cancelled
Publish CLI Package / publish-npm (push) Has been cancelled
Publish Python SDK / publish-pypi (push) Has been cancelled
CI / Deploy Trigger.dev (Dev) (push) Has been cancelled
Helm Chart / Lint, test, and validate chart (push) Has been cancelled
Helm Chart / Chart version bumped (push) Has been cancelled
Publish TypeScript SDK / publish-npm (push) Has been cancelled
CI / Build Dev ECR (blacksmith-8vcpu-ubuntu-2404, ./docker/app.Dockerfile, ECR_APP, linux-x64-8-core) (push) Has been cancelled
CI / Promote Images (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/cron) (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/migrations) (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/pii) (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/realtime) (push) Has been cancelled
CI / Build Dev ECR (blacksmith-2vcpu-ubuntu-2404, ./docker/db.Dockerfile, ECR_MIGRATIONS, ubuntu-latest) (push) Has been cancelled
CI / Build Dev ECR (blacksmith-4vcpu-ubuntu-2404, ./docker/pii.Dockerfile, ECR_PII, ubuntu-latest) (push) Has been cancelled
CI / Build Dev ECR (blacksmith-4vcpu-ubuntu-2404, ./docker/realtime.Dockerfile, ECR_REALTIME, ubuntu-latest) (push) Has been cancelled
CI / Create GHCR Manifests (ghcr.io/simstudioai/simstudio) (push) Has been cancelled
CI / Process Docs (push) Has been cancelled
CI / Create GitHub Release (push) Has been cancelled
CI / Check Desktop Signing Secrets (push) Has been cancelled
CI / Desktop Release (push) Has been cancelled
CI / Create Desktop Prerelease (push) Has been cancelled
CI / Desktop Prerelease Build (push) Has been cancelled
CI / Publish Desktop Prerelease (push) Has been cancelled
CI / Prune Desktop Prereleases (push) Has been cancelled
Helm Chart / Install on kind and run helm test (push) Has been cancelled
WeHub snapshot of cb28d14c6f2c081de7a0d8729a8c816c9adef67a
2026-08-10 11:17:50 +08:00

425 lines
13 KiB
TypeScript

import { AuditAction, AuditResourceType, recordAudit } from '@sim/audit'
import { db } from '@sim/db'
import { permissions, type WorkspaceMode, workflow, workspace } from '@sim/db/schema'
import { createLogger } from '@sim/logger'
import { generateId } from '@sim/utils/id'
import { and, eq, isNull } from 'drizzle-orm'
import { type NextRequest, NextResponse } from 'next/server'
import { listWorkspacesQuerySchema } from '@/lib/api/contracts'
import { createWorkspaceContract } from '@/lib/api/contracts/workspaces'
import { parseRequest } from '@/lib/api/server'
import { getSession } from '@/lib/auth'
import { getActiveOrganizationId } from '@/lib/auth/session-response'
import { PlatformEvents } from '@/lib/core/telemetry'
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
import { captureServerEvent } from '@/lib/posthog/server'
import { buildDefaultWorkflowArtifacts } from '@/lib/workflows/defaults'
import { saveWorkflowToNormalizedTables } from '@/lib/workflows/persistence/utils'
import { getRandomWorkspaceColor } from '@/lib/workspaces/colors'
import { listWorkspacesForViewer } from '@/lib/workspaces/list'
import {
getWorkspaceCreationPolicy,
getWorkspaceInvitePolicy,
lockWorkspaceCreationContext,
resolveInviteFlags,
WORKSPACE_MODE,
WorkspaceCreationContextChangedError,
} from '@/lib/workspaces/policy'
const logger = createLogger('Workspaces')
// Get all workspaces for the current user
export const GET = withRouteHandler(async (request: Request) => {
const session = await getSession()
if (!session?.user?.id) {
return NextResponse.json({ error: 'Unauthorized' }, { status: 401 })
}
const scopeResult = listWorkspacesQuerySchema.safeParse(
Object.fromEntries(new URL(request.url).searchParams.entries())
)
if (!scopeResult.success) {
return NextResponse.json(
{ error: 'Invalid query parameters', details: scopeResult.error.issues },
{ status: 400 }
)
}
const { scope } = scopeResult.data
const activeOrganizationId = getActiveOrganizationId(session)
const payload = await listWorkspacesForViewer({
userId: session.user.id,
activeOrganizationId,
scope,
})
const { lastActiveWorkspaceId, pinnedWorkspaceIds, creationPolicy } = payload
if (scope === 'active' && payload.workspaces.length === 0) {
if (!creationPolicy.canCreate) {
return NextResponse.json({
workspaces: [],
lastActiveWorkspaceId,
pinnedWorkspaceIds,
creationPolicy,
})
}
let defaultWorkspace: Awaited<ReturnType<typeof createDefaultWorkspace>>
try {
defaultWorkspace = await createDefaultWorkspace(
session.user.id,
session.user.name,
creationPolicy
)
} catch (error) {
/**
* The user joined an organization between the empty list read and the
* default-workspace insert. Their workspaces (the join sweep's output)
* exist now — re-list and return that instead of failing the load.
*/
if (error instanceof WorkspaceCreationContextChangedError) {
logger.info(
'Default workspace creation raced an organization membership change; re-listing',
{
userId: session.user.id,
}
)
const refreshedPayload = await listWorkspacesForViewer({
userId: session.user.id,
activeOrganizationId,
scope,
})
return NextResponse.json(refreshedPayload)
}
throw error
}
await migrateExistingWorkflows(session.user.id, defaultWorkspace.id)
const refreshedCreationPolicy = await getWorkspaceCreationPolicy({
userId: session.user.id,
activeOrganizationId,
})
return NextResponse.json({
workspaces: [defaultWorkspace],
lastActiveWorkspaceId,
pinnedWorkspaceIds,
creationPolicy: refreshedCreationPolicy,
})
}
if (scope === 'active') {
await ensureWorkflowsHaveWorkspace(session.user.id, payload.workspaces[0].id)
}
return NextResponse.json(payload)
})
// POST /api/workspaces - Create a new workspace
export const POST = withRouteHandler(async (req: NextRequest) => {
const session = await getSession()
if (!session?.user?.id) {
return NextResponse.json({ error: 'Unauthorized' }, { status: 401 })
}
try {
const parsed = await parseRequest(createWorkspaceContract, req, {})
if (!parsed.success) return parsed.response
const { name, color, skipDefaultWorkflow } = parsed.data.body
const activeOrganizationId = getActiveOrganizationId(session)
const creationPolicy = await getWorkspaceCreationPolicy({
userId: session.user.id,
activeOrganizationId,
})
if (!creationPolicy.canCreate) {
return NextResponse.json(
{ error: creationPolicy.reason || 'Workspace creation is not available.' },
{ status: creationPolicy.status }
)
}
const newWorkspace = await createWorkspace({
userId: session.user.id,
name,
skipDefaultWorkflow,
explicitColor: color,
organizationId: creationPolicy.organizationId,
workspaceMode: creationPolicy.workspaceMode,
billedAccountUserId: creationPolicy.billedAccountUserId,
observedOrganizationId: creationPolicy.observedOrganizationId,
})
captureServerEvent(
session.user.id,
'workspace_created',
{
workspace_id: newWorkspace.id,
name: newWorkspace.name,
workspace_mode: newWorkspace.workspaceMode,
organization_id: newWorkspace.organizationId,
},
{
groups: { workspace: newWorkspace.id },
setOnce: { first_workspace_created_at: new Date().toISOString() },
}
)
recordAudit({
workspaceId: newWorkspace.id,
actorId: session.user.id,
actorName: session.user.name,
actorEmail: session.user.email,
action: AuditAction.WORKSPACE_CREATED,
resourceType: AuditResourceType.WORKSPACE,
resourceId: newWorkspace.id,
resourceName: newWorkspace.name,
description: `Created workspace "${newWorkspace.name}"`,
metadata: {
name: newWorkspace.name,
color: newWorkspace.color,
workspaceMode: newWorkspace.workspaceMode,
organizationId: newWorkspace.organizationId,
},
request: req,
})
return NextResponse.json({ workspace: newWorkspace })
} catch (error) {
if (error instanceof WorkspaceCreationContextChangedError) {
return NextResponse.json(
{
error:
'Your organization membership changed while this workspace was being created. Please try again.',
},
{ status: 409 }
)
}
logger.error('Error creating workspace:', error)
return NextResponse.json({ error: 'Failed to create workspace' }, { status: 500 })
}
})
async function createDefaultWorkspace(
userId: string,
userName: string | null | undefined,
creationPolicy: {
organizationId: string | null
workspaceMode: WorkspaceMode
billedAccountUserId: string
observedOrganizationId: string | null
}
) {
const firstName = userName?.split(' ')[0] || null
const workspaceName = firstName ? `${firstName}'s Workspace` : 'My Workspace'
return createWorkspace({
userId,
name: workspaceName,
organizationId: creationPolicy.organizationId,
workspaceMode: creationPolicy.workspaceMode,
billedAccountUserId: creationPolicy.billedAccountUserId,
observedOrganizationId: creationPolicy.observedOrganizationId,
})
}
interface CreateWorkspaceParams {
userId: string
/** Membership the creation policy observed; see WorkspaceCreationPolicy. */
observedOrganizationId: string | null
name: string
skipDefaultWorkflow?: boolean
explicitColor?: string
organizationId: string | null
workspaceMode: WorkspaceMode
billedAccountUserId: string
}
async function createWorkspace({
userId,
observedOrganizationId,
name,
skipDefaultWorkflow = false,
explicitColor,
organizationId,
workspaceMode,
billedAccountUserId,
}: CreateWorkspaceParams) {
const workspaceId = generateId()
const workflowId = generateId()
const now = new Date()
const color = explicitColor || getRandomWorkspaceColor()
let committedBilledAccountUserId = billedAccountUserId
try {
await db.transaction(async (tx) => {
/**
* Creation takes the same organization → user → membership fence as
* source access removal and transfer. If creation commits first, their
* post-lock workspace-set re-read sees this row and cleans it up. If the
* membership mutation commits first, this re-read rejects the stale
* creation policy before inserting anything.
*/
const lockedCreationContext = await lockWorkspaceCreationContext(tx, {
userId,
organizationId,
observedOrganizationId,
})
const currentBilledAccountUserId =
workspaceMode === WORKSPACE_MODE.ORGANIZATION
? lockedCreationContext.billedAccountUserId
: billedAccountUserId
committedBilledAccountUserId = currentBilledAccountUserId
await tx.insert(workspace).values({
id: workspaceId,
name,
color,
ownerId: userId,
organizationId,
workspaceMode,
billedAccountUserId: currentBilledAccountUserId,
allowPersonalApiKeys: true,
createdAt: now,
updatedAt: now,
})
const permissionRows = [
{
id: generateId(),
entityType: 'workspace' as const,
entityId: workspaceId,
userId,
permissionType: 'admin' as const,
createdAt: now,
updatedAt: now,
},
]
if (
workspaceMode === WORKSPACE_MODE.ORGANIZATION &&
currentBilledAccountUserId &&
currentBilledAccountUserId !== userId
) {
permissionRows.push({
id: generateId(),
entityType: 'workspace' as const,
entityId: workspaceId,
userId: currentBilledAccountUserId,
permissionType: 'admin' as const,
createdAt: now,
updatedAt: now,
})
}
await tx.insert(permissions).values(permissionRows)
if (!skipDefaultWorkflow) {
await tx.insert(workflow).values({
id: workflowId,
userId,
workspaceId,
folderId: null,
name: 'default-agent',
description: 'Your first workflow - start building here!',
lastSynced: now,
createdAt: now,
updatedAt: now,
isDeployed: false,
runCount: 0,
variables: {},
})
const { workflowState } = buildDefaultWorkflowArtifacts()
await saveWorkflowToNormalizedTables(workflowId, workflowState, tx)
}
logger.info(
skipDefaultWorkflow
? `Created ${workspaceMode} workspace ${workspaceId} for user ${userId}`
: `Created ${workspaceMode} workspace ${workspaceId} with initial workflow ${workflowId} for user ${userId}`
)
})
} catch (error) {
logger.error(`Failed to create workspace ${workspaceId}:`, error)
throw error
}
try {
PlatformEvents.workspaceCreated({
workspaceId,
userId,
name,
})
} catch {
// Telemetry should not fail the operation
}
const invitePolicy = await getWorkspaceInvitePolicy({
organizationId,
workspaceMode,
billedAccountUserId: committedBilledAccountUserId,
ownerId: userId,
})
return {
id: workspaceId,
name,
color,
ownerId: userId,
organizationId,
workspaceMode,
billedAccountUserId: committedBilledAccountUserId,
allowPersonalApiKeys: true,
createdAt: now,
updatedAt: now,
role: 'owner',
permissions: 'admin',
...resolveInviteFlags(invitePolicy, committedBilledAccountUserId === userId),
}
}
async function migrateExistingWorkflows(userId: string, workspaceId: string) {
const orphanedWorkflows = await db
.select({ id: workflow.id })
.from(workflow)
.where(and(eq(workflow.userId, userId), isNull(workflow.workspaceId)))
if (orphanedWorkflows.length === 0) {
return // No orphaned workflows to migrate
}
logger.info(
`Migrating ${orphanedWorkflows.length} workflows to workspace ${workspaceId} for user ${userId}`
)
await db
.update(workflow)
.set({
workspaceId: workspaceId,
updatedAt: new Date(),
})
.where(and(eq(workflow.userId, userId), isNull(workflow.workspaceId)))
}
async function ensureWorkflowsHaveWorkspace(userId: string, defaultWorkspaceId: string) {
const orphanedWorkflows = await db
.select()
.from(workflow)
.where(and(eq(workflow.userId, userId), isNull(workflow.workspaceId)))
if (orphanedWorkflows.length > 0) {
await db
.update(workflow)
.set({
workspaceId: defaultWorkspaceId,
updatedAt: new Date(),
})
.where(and(eq(workflow.userId, userId), isNull(workflow.workspaceId)))
logger.info(`Fixed ${orphanedWorkflows.length} orphaned workflows for user ${userId}`)
}
}