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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
81 changes: 81 additions & 0 deletions .github/scripts/check-project-rollout.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,81 @@
#!/usr/bin/env python3
"""Read-only, fail-closed ECS retirement check for the Project contract migration.

The expected digest is the operator's release-scoped acknowledgment that Project
compatible writers are deployed and relevant old worker jobs are drained.
AWS checks below independently verify ECS retirement, not worker drainage.
"""
import argparse
import json
import re
import subprocess
import sys


def aws(region, *args):
result = subprocess.run(
['aws', '--region', region, '--no-cli-pager', '--cli-connect-timeout', '10', '--cli-read-timeout', '30', *args, '--output', 'json'],
capture_output=True, text=True, timeout=90, check=False,
)
if result.returncode:
raise RuntimeError('AWS preflight read failed; check deployment-read permissions')
return json.loads(result.stdout)


def verify(environment, region, digest):
if not re.fullmatch(r'sha256:[0-9a-f]{64}', digest):
raise RuntimeError('Set the environment-specific PROJECT_ENFORCEMENT_READY_IMAGE_DIGEST after reviewing compatible rollout and worker-drain evidence')
pipeline = f'sim-{environment}-{region}-app-deployment'
executions = aws(region, 'codepipeline', 'list-pipeline-executions', '--pipeline-name', pipeline).get('pipelineExecutionSummaries', [])
if not executions or executions[0].get('status') != 'Succeeded':
raise RuntimeError('The latest app pipeline has not completed; traffic cutover alone is insufficient')
execution_id = executions[0]['pipelineExecutionId']
group = aws(region, 'deploy', 'get-deployment-group', '--application-name', f'sim-{environment}-{region}-ecs-app',
'--deployment-group-name', f'sim-{environment}-{region}-app-dg')['deploymentGroupInfo']
services = group.get('ecsServices', [])
if len(services) != 1:
raise RuntimeError('Expected exactly one application ECS service')
cluster, service = services[0]['clusterName'], services[0]['serviceName']
description = aws(region, 'ecs', 'describe-services', '--cluster', cluster, '--services', service)
if description.get('failures') or len(description.get('services', [])) != 1:
raise RuntimeError('Cannot inspect the application ECS service')
record = description['services'][0]
if record.get('desiredCount', 0) < 1 or record.get('runningCount') != record['desiredCount'] or record.get('pendingCount') != 0:
raise RuntimeError('Application ECS service is not stable')
arns = set()
for status in ('RUNNING', 'STOPPED'):
arns.update(aws(region, 'ecs', 'list-tasks', '--cluster', cluster, '--service-name', service,
'--desired-status', status).get('taskArns', []))
live = []
ordered = sorted(arns)
for start in range(0, len(ordered), 100):
response = aws(region, 'ecs', 'describe-tasks', '--cluster', cluster, '--tasks', *ordered[start:start + 100])
if response.get('failures'):
raise RuntimeError('Cannot account for every ECS task')
if len(response.get('tasks', [])) != len(ordered[start:start + 100]):
raise RuntimeError('Incomplete ECS task response')
live.extend(task for task in response['tasks'] if task.get('lastStatus') != 'STOPPED')
if len(live) != record['desiredCount']:
raise RuntimeError('Old, stopping, or pending ECS tasks remain')
for task in live:
app = [container for container in task.get('containers', []) if container.get('name') == 'app']
if task.get('lastStatus') != 'RUNNING' or task.get('desiredStatus') != 'RUNNING' or len(app) != 1 or app[0].get('imageDigest') != digest:
raise RuntimeError('A live ECS task does not match the acknowledged compatible release')
latest = aws(region, 'codepipeline', 'list-pipeline-executions', '--pipeline-name', pipeline).get('pipelineExecutionSummaries', [])
if not latest or latest[0].get('pipelineExecutionId') != execution_id or latest[0].get('status') != 'Succeeded':
raise RuntimeError('Application deployment changed during preflight')
print(json.dumps({'ecsRetired': True, 'expectedImageDigest': digest, 'pipelineExecutionId': execution_id,
'operatorAcknowledgedCompatibleWritersAndWorkers': True}))


if __name__ == '__main__':
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument('--environment', required=True, choices=['production', 'staging'])
parser.add_argument('--region', required=True)
parser.add_argument('--expected-image-digest', required=True)
args = parser.parse_args()
try:
verify(args.environment, args.region, args.expected_image_digest)
except (RuntimeError, ValueError, KeyError, TypeError, subprocess.TimeoutExpired) as error:
print(f'Project rollout preflight refused: {error}', file=sys.stderr)
sys.exit(1)
6 changes: 6 additions & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -116,6 +116,9 @@ jobs:
# in place before the new app version deploys (replaces the removed ECS
# migration sidecar)
migrate:
permissions:
contents: read
id-token: write
name: Migrate DB
needs: [test-build]
# Explicit need results instead of the implicit success(): a skipped job
Expand All @@ -134,6 +137,9 @@ jobs:

# Same ordering for dev (schema push before the dev image lands in ECR)
migrate-dev:
permissions:
contents: read
id-token: write
name: Migrate Dev DB
if: github.event_name == 'push' && github.ref == 'refs/heads/dev'
uses: ./.github/workflows/migrations.yml
Expand Down
27 changes: 27 additions & 0 deletions .github/workflows/migrations.yml
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ on:

permissions:
contents: read
id-token: write

jobs:
migrate:
Expand Down Expand Up @@ -51,6 +52,32 @@ jobs:
- name: Install dependencies
run: bun install --frozen-lockfile --ignore-scripts

- name: Check whether Project enforcement needs rollout evidence
id: project-contract
if: inputs.environment != 'dev'
working-directory: ./packages/db
env:
DATABASE_URL: ${{ inputs.environment == 'production' && secrets.DATABASE_URL || inputs.environment == 'staging' && secrets.STAGING_DATABASE_URL || '' }}
MIGRATION_DATABASE_URL: ${{ inputs.environment == 'production' && secrets.MIGRATION_DATABASE_URL || inputs.environment == 'staging' && secrets.STAGING_MIGRATION_DATABASE_URL || '' }}
run: bun --no-env-file scripts/project-contract-required.ts >> "$GITHUB_OUTPUT"

- name: Configure AWS for Project retirement verification
if: steps.project-contract.outputs.required == 'true'
uses: aws-actions/configure-aws-credentials@e7f100cf4c008499ea8adda475de1042d6975c7b
with:
role-to-assume: ${{ inputs.environment == 'production' && secrets.AWS_ROLE_TO_ASSUME || secrets.STAGING_AWS_ROLE_TO_ASSUME }}
aws-region: ${{ inputs.environment == 'production' && secrets.AWS_REGION || secrets.STAGING_AWS_REGION }}

# Set this release-scoped acknowledgment only after compatible writers are deployed
# everywhere and relevant old Trigger.dev runs are drained. SQL performs the backfill.
# The preflight independently checks ECS retirement; it does not inspect worker runs.
- name: Require completed compatible rollout before Project enforcement
if: steps.project-contract.outputs.required == 'true'
env:
ENVIRONMENT: ${{ inputs.environment }}
EXPECTED_IMAGE_DIGEST: ${{ inputs.environment == 'production' && vars.PROJECT_ENFORCEMENT_READY_IMAGE_DIGEST_PRODUCTION || inputs.environment == 'staging' && vars.PROJECT_ENFORCEMENT_READY_IMAGE_DIGEST_STAGING || '' }}
run: python3 .github/scripts/check-project-rollout.py --environment "$ENVIRONMENT" --region "$AWS_REGION" --expected-image-digest "$EXPECTED_IMAGE_DIGEST"

# The expression maps the explicit environment input to exactly one repo
# secret, so the job never holds another environment's database URL. An
# unknown environment resolves to empty and the guard below fails the job.
Expand Down
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import { deleteWorkspaceFixture } from '@sim/db/testing/workspace-fixtures'
/**
* The member sync scheduler's reclaim against real PostgreSQL: a members-mode connector whose
* member lease went stale is put back on the failure ladder, and a connector in any other access
Expand Down Expand Up @@ -51,7 +52,7 @@ describe('member sync reclaim in PostgreSQL', () => {
})

afterAll(async () => {
await db.delete(workspace).where(eq(workspace.id, ids.workspaceId))
await deleteWorkspaceFixture(db, eq(workspace.id, ids.workspaceId))
await db.delete(organization).where(eq(organization.id, ids.organizationId))
await db.delete(user).where(inArray(user.id, [ids.aliceId, ids.bobId]))
await db.$client.end()
Expand Down
3 changes: 2 additions & 1 deletion apps/sim/app/api/v1/knowledge/route.integration.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
import type { Principal } from '@sim/auth/principal'
import { db } from '@sim/db'
import { document, organization, user, workspace } from '@sim/db/schema'
import { deleteWorkspaceFixture } from '@sim/db/testing/workspace-fixtures'
import { authMock, authMockFns, createMockRequest } from '@sim/testing'
import { generateId } from '@sim/utils/id'
import { eq, inArray } from 'drizzle-orm'
Expand Down Expand Up @@ -85,7 +86,7 @@ describe('knowledge-base document totals in PostgreSQL', () => {
})

afterAll(async () => {
await db.delete(workspace).where(eq(workspace.id, ids.workspaceId))
await deleteWorkspaceFixture(db, eq(workspace.id, ids.workspaceId))
await db.delete(organization).where(eq(organization.id, ids.organizationId))
await db.delete(user).where(inArray(user.id, [ids.aliceId, ids.bobId]))
vi.unstubAllGlobals()
Expand Down
11 changes: 9 additions & 2 deletions apps/sim/background/cleanup-table-row-ttl.integration.ts
Original file line number Diff line number Diff line change
Expand Up @@ -110,8 +110,12 @@ describe.skipIf(!migrated)('Expiration with real PostgreSQL transactions', () =>
beforeAll(async () => {
await control`INSERT INTO "user" (id, name, email, email_verified, created_at, updated_at)
VALUES (${userId}, 'Expiration integration fixture', ${`${userId}@example.test`}, true, now(), now())`
await control`INSERT INTO workspace (id, name, owner_id, billed_account_user_id)
await control.begin(async (tx) => {
await tx`INSERT INTO workspace (id, name, owner_id, billed_account_user_id)
VALUES (${workspaceId}, 'Expiration integration fixtures', ${userId}, ${userId})`
await tx`INSERT INTO project (id, name, owner_id) VALUES (${workspaceId}, 'Fixture project', ${userId})`
await tx`INSERT INTO project_workspace (project_id, workspace_id) VALUES (${workspaceId}, ${workspaceId})`
})
})

beforeEach(async () => {
Expand All @@ -126,7 +130,10 @@ describe.skipIf(!migrated)('Expiration with real PostgreSQL transactions', () =>
})

afterAll(async () => {
await control`DELETE FROM workspace WHERE id = ${workspaceId}`
await control.begin(async (tx) => {
await tx`DELETE FROM workspace WHERE id = ${workspaceId}`
await tx`DELETE FROM project WHERE id = ${workspaceId}`
})
await control`DELETE FROM "user" WHERE id = ${userId}`
writeFileSync(
join(tmpdir(), 'expiration-qa-measurements.json'),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ import {
user,
workspace,
} from '@sim/db/schema'
import { deleteWorkspaceFixture, insertWorkspaceFixture } from '@sim/db/testing/workspace-fixtures'
import { sha256Hex } from '@sim/security/hash'
import { generateId } from '@sim/utils/id'
import { and, eq, inArray } from 'drizzle-orm'
Expand Down Expand Up @@ -106,7 +107,8 @@ describe('organization personal tokens', () => {
{ id: generateId(), organizationId: ids.org, userId: ids.owner, role: 'member' },
{ id: generateId(), organizationId: ids.org, userId: ids.other, role: 'admin' },
])
await db.insert(workspace).values(
await insertWorkspaceFixture(
db,
[ids.first, ids.second, ids.foreign].map((id) => ({
id,
name: 'Token fixture workspace',
Expand Down Expand Up @@ -187,7 +189,7 @@ describe('organization personal tokens', () => {
await db
.delete(credentialGroup)
.where(inArray(credentialGroup.id, [ids.group, ids.legacyGroup]))
await db.delete(workspace).where(inArray(workspace.id, [ids.first, ids.second, ids.foreign]))
await deleteWorkspaceFixture(db, inArray(workspace.id, [ids.first, ids.second, ids.foreign]))
await db.delete(organization).where(inArray(organization.id, [ids.org, ids.foreignOrg]))
await db.delete(user).where(inArray(user.id, [ids.owner, ids.other]))
})
Expand Down Expand Up @@ -287,7 +289,7 @@ describe('organization personal tokens', () => {
expect.objectContaining({ id: ids.token, workspaceId: null, organizationId: ids.org }),
])
}
await db.delete(workspace).where(eq(workspace.id, ids.first))
await deleteWorkspaceFixture(db, eq(workspace.id, ids.first))
await expect(resolve()).resolves.toMatchObject({ accessToken: tokenSecret })
})

Expand Down
5 changes: 3 additions & 2 deletions apps/sim/lib/environment/execution-environment.integration.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
*/
import { db } from '@sim/db'
import { environment, permissions, user, workspace, workspaceEnvironment } from '@sim/db/schema'
import { deleteWorkspaceFixture, insertWorkspaceFixture } from '@sim/db/testing/workspace-fixtures'
import { generateId } from '@sim/utils/id'
import { eq, inArray } from 'drizzle-orm'
import { afterAll, beforeAll, describe, expect, it } from 'vitest'
Expand Down Expand Up @@ -41,7 +42,7 @@ beforeAll(async () => {
...(id === suspended ? { banned: true } : {}),
}))
)
await db.insert(workspace).values({
await insertWorkspaceFixture(db, {
id: workspaceId,
name: 'Environment',
ownerId: owner,
Expand Down Expand Up @@ -73,7 +74,7 @@ beforeAll(async () => {
})

afterAll(async () => {
await db.delete(workspace).where(eq(workspace.id, workspaceId))
await deleteWorkspaceFixture(db, eq(workspace.id, workspaceId))
await db.delete(user).where(inArray(user.id, userIds))
})

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ import {
user,
workspace,
} from '@sim/db/schema'
import { deleteWorkspaceFixture } from '@sim/db/testing/workspace-fixtures'
import { generateId } from '@sim/utils/id'
import { and, eq } from 'drizzle-orm'
import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest'
Expand Down Expand Up @@ -278,7 +279,7 @@ describe('indexed source content through real application access', () => {
CONFLUENCE_CLIENT_ID: previousConfluenceClient.id,
CONFLUENCE_CLIENT_SECRET: previousConfluenceClient.secret,
})
await db.delete(workspace).where(eq(workspace.id, workspaceId))
await deleteWorkspaceFixture(db, eq(workspace.id, workspaceId))
await db.delete(user).where(eq(user.id, aliceId))
await db.delete(user).where(eq(user.id, bobId))
await rm(fixtures.storageRoot, { recursive: true, force: true })
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ import {
user,
workspace,
} from '@sim/db/schema'
import { deleteWorkspaceFixture } from '@sim/db/testing/workspace-fixtures'
import { sleep } from '@sim/utils/helpers'
import { generateId } from '@sim/utils/id'
import { serializeSignedCookie } from 'better-call'
Expand Down Expand Up @@ -225,7 +226,7 @@ describe.skipIf(!tokenPath || !fixturePath || !secondEmail)(
.where(eq(document.knowledgeBaseId, ids.knowledgeBaseId))
for (const row of rows)
if (row.storageKey) await deleteFile({ key: row.storageKey, context: 'knowledge-base' })
await db.delete(workspace).where(eq(workspace.id, ids.workspaceId))
await deleteWorkspaceFixture(db, eq(workspace.id, ids.workspaceId))
await db.delete(organization).where(eq(organization.id, ids.organizationId))
await db.delete(user).where(eq(user.id, ids.aliceId))
await db.delete(user).where(eq(user.id, ids.bobId))
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ import {
user,
workspace,
} from '@sim/db/schema'
import { deleteWorkspaceFixture } from '@sim/db/testing/workspace-fixtures'
import { generateId } from '@sim/utils/id'
import { and, eq } from 'drizzle-orm'
import { afterAll, beforeAll, beforeEach, describe, expect, it } from 'vitest'
Expand Down Expand Up @@ -99,7 +100,7 @@ describe('Confluence mirrored-identity self-enrollment', () => {
CONFLUENCE_CLIENT_SECRET: previousClient.secret,
})
for (const fixture of [ids, foreign]) {
await db.delete(workspace).where(eq(workspace.id, fixture.workspaceId))
await deleteWorkspaceFixture(db, eq(workspace.id, fixture.workspaceId))
await db.delete(user).where(eq(user.id, fixture.aliceId))
await db.delete(user).where(eq(user.id, fixture.bobId))
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ import {
user,
workspace,
} from '@sim/db/schema'
import { deleteWorkspaceFixture } from '@sim/db/testing/workspace-fixtures'
import { generateId } from '@sim/utils/id'
import { eq } from 'drizzle-orm'
import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest'
Expand Down Expand Up @@ -181,7 +182,7 @@ describe('Confluence identities with hidden directory email', () => {
afterAll(async () => {
vi.unstubAllGlobals()
Object.assign(env, previousConfluenceClient)
await db.delete(workspace).where(eq(workspace.id, ids.workspaceId))
await deleteWorkspaceFixture(db, eq(workspace.id, ids.workspaceId))
await db.delete(user).where(eq(user.id, ids.aliceId))
await db.delete(user).where(eq(user.id, ids.bobId))
await db.$client.end()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ import {
user,
workspace,
} from '@sim/db/schema'
import { deleteWorkspaceFixture } from '@sim/db/testing/workspace-fixtures'
import { generateId } from '@sim/utils/id'
import { eq, inArray } from 'drizzle-orm'
import { afterAll, beforeAll, beforeEach, describe, expect, it, vi } from 'vitest'
Expand Down Expand Up @@ -88,7 +89,7 @@ describe('durable connector capacity deferrals', () => {
})
afterAll(async () => {
await db.delete(knowledgeBase).where(eq(knowledgeBase.id, ids.knowledgeBaseId))
await db.delete(workspace).where(eq(workspace.id, ids.workspaceId))
await deleteWorkspaceFixture(db, eq(workspace.id, ids.workspaceId))
await db.delete(user).where(inArray(user.id, [ids.aliceId, ids.bobId]))
vi.restoreAllMocks()
vi.unstubAllGlobals()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ import {
} from '@sim/db/schema'
import { installProjectionSourceAcl } from '@sim/db/script-migrations/0021_embedding_search_connector'
import { installKnowledgeProjectionAsync } from '@sim/db/script-migrations/0024_knowledge_projection_async'
import { deleteWorkspaceFixture } from '@sim/db/testing/workspace-fixtures'
import { generateId } from '@sim/utils/id'
import { and, eq, inArray, sql } from 'drizzle-orm'
import postgres from 'postgres'
Expand Down Expand Up @@ -141,7 +142,7 @@ describe('connector lease ACL pages in PostgreSQL', () => {
await db.execute(sql`DROP TRIGGER IF EXISTS fail_after_acl_writes ON document`)
await db.execute(sql`DELETE FROM lease_page_acl_writes`)
await db.execute(sql`DELETE FROM lease_page_projection_writes`)
await db.delete(workspace).where(eq(workspace.id, ids.workspaceId))
await deleteWorkspaceFixture(db, eq(workspace.id, ids.workspaceId))
await db.delete(organization).where(eq(organization.id, ids.organizationId))
await db.delete(user).where(inArray(user.id, [ids.aliceId, ids.bobId]))
})
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ import {
user,
workspace,
} from '@sim/db/schema'
import { deleteWorkspaceFixture } from '@sim/db/testing/workspace-fixtures'
import { generateId } from '@sim/utils/id'
import { and, eq, inArray, isNotNull, sql } from 'drizzle-orm'
import { afterAll, afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
Expand Down Expand Up @@ -97,7 +98,7 @@ describe('source lifecycle KB guards', () => {
})
afterEach(async () => {
vi.restoreAllMocks()
await db.delete(workspace).where(eq(workspace.id, ids.workspaceId))
await deleteWorkspaceFixture(db, eq(workspace.id, ids.workspaceId))
await db.delete(organization).where(eq(organization.id, ids.organizationId))
await db.delete(user).where(inArray(user.id, [ids.aliceId, ids.bobId]))
})
Expand Down
Loading
Loading