Skip to content

Commit 5e11d7d

Browse files
committed
refactor(projects): backfill through bounded SQL migration batches
1 parent 2e94a13 commit 5e11d7d

5 files changed

Lines changed: 513 additions & 51 deletions

File tree

‎.github/scripts/check-project-rollout.py‎

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -2,8 +2,8 @@
22
"""Read-only, fail-closed ECS retirement check for the Project contract migration.
33
44
The expected digest is the operator's release-scoped acknowledgment that Project
5-
writers are enabled, relevant old worker jobs are drained, and backfill verification
6-
passed. AWS checks below independently verify ECS retirement, not those assertions.
5+
compatible writers are deployed and relevant old worker jobs are drained.
6+
AWS checks below independently verify ECS retirement, not worker drainage.
77
"""
88
import argparse
99
import json
@@ -24,7 +24,7 @@ def aws(region, *args):
2424

2525
def verify(environment, region, digest):
2626
if not re.fullmatch(r'sha256:[0-9a-f]{64}', digest):
27-
raise RuntimeError('Set the environment-specific PROJECT_ENFORCEMENT_READY_IMAGE_DIGEST after reviewing rollout and backfill evidence')
27+
raise RuntimeError('Set the environment-specific PROJECT_ENFORCEMENT_READY_IMAGE_DIGEST after reviewing compatible rollout and worker-drain evidence')
2828
pipeline = f'sim-{environment}-{region}-app-deployment'
2929
executions = aws(region, 'codepipeline', 'list-pipeline-executions', '--pipeline-name', pipeline).get('pipelineExecutionSummaries', [])
3030
if not executions or executions[0].get('status') != 'Succeeded':
@@ -65,7 +65,7 @@ def verify(environment, region, digest):
6565
if not latest or latest[0].get('pipelineExecutionId') != execution_id or latest[0].get('status') != 'Succeeded':
6666
raise RuntimeError('Application deployment changed during preflight')
6767
print(json.dumps({'ecsRetired': True, 'expectedImageDigest': digest, 'pipelineExecutionId': execution_id,
68-
'operatorAcknowledgedWritersWorkersAndBackfill': True}))
68+
'operatorAcknowledgedCompatibleWritersAndWorkers': True}))
6969

7070

7171
if __name__ == '__main__':

‎.github/workflows/migrations.yml‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -68,8 +68,8 @@ jobs:
6868
role-to-assume: ${{ inputs.environment == 'production' && secrets.AWS_ROLE_TO_ASSUME || secrets.STAGING_AWS_ROLE_TO_ASSUME }}
6969
aws-region: ${{ inputs.environment == 'production' && secrets.AWS_REGION || secrets.STAGING_AWS_REGION }}
7070

71-
# Set this release-scoped acknowledgment only after writers are enabled everywhere,
72-
# relevant old Trigger.dev runs are drained, and the reviewed backfill verifies ready.
71+
# Set this release-scoped acknowledgment only after compatible writers are deployed
72+
# everywhere and relevant old Trigger.dev runs are drained. SQL performs the backfill.
7373
# The preflight independently checks ECS retirement; it does not inspect worker runs.
7474
- name: Require completed compatible rollout before Project enforcement
7575
if: steps.project-contract.outputs.required == 'true'

‎apps/sim/lib/projects/__integration__/foundation.integration.ts‎

Lines changed: 8 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -202,15 +202,17 @@ afterAll(async () => {
202202
if (restoreEnforcement) {
203203
const client = postgres(readTestDatabaseUrl(), { max: 1, onnotice: () => undefined })
204204
try {
205-
await client.unsafe(
205+
for (const statement of (
206206
await readFile(
207207
new URL(
208208
'../../../../../packages/db/migrations/0394_project_membership_enforcement.sql',
209209
import.meta.url
210210
),
211211
'utf8'
212212
)
213-
)
213+
).split('--> statement-breakpoint')) {
214+
await client.unsafe(statement)
215+
}
214216
} finally {
215217
await client.end()
216218
}
@@ -223,15 +225,17 @@ describe('Project foundation at the database and application boundary', () => {
223225
async () => {
224226
const client = postgres(readTestDatabaseUrl(), { max: 1, onnotice: () => undefined })
225227
try {
226-
await client.unsafe(
228+
for (const statement of (
227229
await readFile(
228230
new URL(
229231
'../../../../../packages/db/migrations/0394_project_membership_enforcement.sql',
230232
import.meta.url
231233
),
232234
'utf8'
233235
)
234-
)
236+
).split('--> statement-breakpoint')) {
237+
await client.unsafe(statement)
238+
}
235239
const f = await fixture(true, 1)
236240
const organizationId = f.organizationId
237241
if (!organizationId) throw new Error('Missing organization fixture')

‎packages/db/migrations/0394_project_membership_enforcement.sql‎

Lines changed: 154 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -1,15 +1,144 @@
1-
BEGIN;
1+
-- Each family commits separately; never retain locks across the full backfill.
2+
-- The migration runner serializes runners with its session advisory lock.
3+
COMMIT;
4+
--> statement-breakpoint
5+
SET statement_timeout = '15min';
6+
--> statement-breakpoint
7+
CREATE TEMP TABLE IF NOT EXISTS project_backfill_roots (id text PRIMARY KEY) ON COMMIT PRESERVE ROWS;
8+
TRUNCATE project_backfill_roots;
9+
INSERT INTO project_backfill_roots SELECT id FROM workspace WHERE forked_from_workspace_id IS NULL;
10+
--> statement-breakpoint
11+
CREATE OR REPLACE PROCEDURE pg_temp.backfill_project_families() LANGUAGE plpgsql AS $$
12+
DECLARE
13+
root_id text;
14+
previous_root text;
15+
environment_id text;
16+
family_ids text[];
17+
current_ids text[];
18+
project_ids text[];
19+
target_project text;
20+
root workspace%ROWTYPE;
21+
existing project%ROWTYPE;
22+
archive_time timestamp;
23+
attempts integer := 0;
24+
completed integer := 0;
25+
assigned integer := 0;
26+
inserted integer;
27+
retry boolean;
28+
BEGIN
29+
LOOP
30+
SELECT id INTO root_id FROM project_backfill_roots
31+
WHERE previous_root IS NULL OR id > previous_root ORDER BY id LIMIT 1;
32+
EXIT WHEN root_id IS NULL;
33+
-- A pathological family cannot hold environment locks indefinitely (PostgreSQL 17+).
34+
PERFORM set_config('transaction_timeout', '5s', true);
35+
PERFORM set_config('lock_timeout', '1s', true);
36+
retry := false;
37+
BEGIN
38+
WITH RECURSIVE family(id) AS (
39+
SELECT id FROM workspace WHERE id = root_id AND forked_from_workspace_id IS NULL
40+
UNION
41+
SELECT w.id FROM workspace w JOIN family f ON w.forked_from_workspace_id = f.id
42+
) SELECT array_agg(id ORDER BY id) INTO family_ids FROM (SELECT id FROM family LIMIT 1001) bounded;
43+
IF cardinality(family_ids) > 1000 THEN
44+
RAISE EXCEPTION 'Project backfill family exceeds 1000 environments' USING ERRCODE = '54000', DETAIL = root_id;
45+
END IF;
46+
IF family_ids IS NOT NULL THEN
47+
-- Application writers take the shared form before reading membership, including absence.
48+
FOREACH environment_id IN ARRAY family_ids LOOP
49+
IF NOT pg_try_advisory_xact_lock(hashtextextended('project-backfill:' || environment_id, 0)) THEN
50+
RAISE EXCEPTION 'Project backfill environment is busy' USING ERRCODE = '55P03';
51+
END IF;
52+
END LOOP;
53+
-- Do not wait while holding a partial lock set; unrelated environments remain writable.
54+
PERFORM id FROM workspace WHERE id = ANY(family_ids) ORDER BY id FOR NO KEY UPDATE NOWAIT;
55+
WITH RECURSIVE family(id) AS (
56+
SELECT id FROM workspace WHERE id = root_id AND forked_from_workspace_id IS NULL
57+
UNION
58+
SELECT w.id FROM workspace w JOIN family f ON w.forked_from_workspace_id = f.id
59+
) SELECT array_agg(id ORDER BY id) INTO current_ids FROM (SELECT id FROM family LIMIT 1001) bounded;
60+
IF family_ids IS DISTINCT FROM current_ids THEN
61+
RAISE EXCEPTION 'Project backfill lineage changed during discovery' USING ERRCODE = '55P03';
62+
END IF;
63+
SELECT * INTO STRICT root FROM workspace WHERE id = root_id;
64+
IF EXISTS (SELECT 1 FROM workspace WHERE id = ANY(family_ids) AND organization_id IS DISTINCT FROM root.organization_id) THEN
65+
RAISE EXCEPTION 'Project backfill family spans organizations; reconcile before retrying' USING ERRCODE = '55000', DETAIL = root_id;
66+
END IF;
67+
SELECT array_agg(DISTINCT project_id ORDER BY project_id) INTO project_ids
68+
FROM project_workspace WHERE workspace_id = ANY(family_ids);
69+
IF cardinality(project_ids) > 1 THEN
70+
RAISE EXCEPTION 'Project backfill family spans Projects; reconcile before retrying' USING ERRCODE = '55000', DETAIL = root_id;
71+
END IF;
72+
SELECT CASE WHEN count(*) = count(archived_at) THEN max(archived_at) END INTO archive_time
73+
FROM workspace WHERE id = ANY(family_ids);
74+
target_project := project_ids[1];
75+
IF target_project IS NOT NULL THEN
76+
IF NOT pg_try_advisory_xact_lock(hashtextextended('project:' || target_project, 0)) THEN
77+
RAISE EXCEPTION 'Project backfill Project is busy' USING ERRCODE = '55P03';
78+
END IF;
79+
SELECT * INTO STRICT existing FROM project WHERE id = target_project FOR UPDATE NOWAIT;
80+
IF existing.organization_id IS DISTINCT FROM root.organization_id
81+
OR (existing.archived_at IS NULL) <> (archive_time IS NULL)
82+
OR EXISTS (SELECT 1 FROM project_workspace WHERE project_id = target_project AND NOT workspace_id = ANY(family_ids)) THEN
83+
RAISE EXCEPTION 'Project backfill existing assignment has incompatible scope, archive state, or lineage' USING ERRCODE = '55000', DETAIL = root_id;
84+
END IF;
85+
ELSE
86+
target_project := gen_random_uuid()::text;
87+
INSERT INTO project (id, name, owner_id, organization_id, archived_at)
88+
VALUES (target_project, left(coalesce(nullif(btrim(root.name), ''), 'Untitled'), 90) || ' - Project', root.owner_id, root.organization_id, archive_time);
89+
END IF;
90+
INSERT INTO project_workspace (project_id, workspace_id)
91+
SELECT target_project, id FROM unnest(family_ids) ids(id)
92+
WHERE NOT EXISTS (SELECT 1 FROM project_workspace pw WHERE pw.workspace_id = ids.id);
93+
GET DIAGNOSTICS inserted = ROW_COUNT;
94+
assigned := assigned + inserted;
95+
END IF;
96+
EXCEPTION WHEN lock_not_available OR deadlock_detected OR serialization_failure THEN
97+
retry := true;
98+
END;
99+
-- Commit outside the exception subtransaction, releasing every row/advisory lock before retry.
100+
COMMIT;
101+
IF retry THEN
102+
attempts := attempts + 1;
103+
IF attempts >= 20 THEN
104+
RAISE EXCEPTION 'Project backfill family remains busy; safe to retry migration' USING ERRCODE = '55P03', DETAIL = root_id;
105+
END IF;
106+
PERFORM pg_sleep(0.05 + random() * 0.15);
107+
COMMIT;
108+
ELSE
109+
attempts := 0;
110+
completed := completed + 1;
111+
previous_root := root_id;
112+
IF completed % 100 = 0 THEN
113+
RAISE NOTICE 'Project backfill: % families processed, % memberships assigned', completed, assigned;
114+
END IF;
115+
END IF;
116+
END LOOP;
117+
RAISE NOTICE 'Project backfill complete: % families processed, % memberships assigned', completed, assigned;
118+
END;
119+
$$;
2120
--> statement-breakpoint
3-
SET LOCAL lock_timeout = '5s';
121+
CALL pg_temp.backfill_project_families();
4122
--> statement-breakpoint
5-
SET LOCAL statement_timeout = '60s';
123+
DROP PROCEDURE pg_temp.backfill_project_families();
124+
-- migration-safe: session-local scratch table created above; no application readers or persistent data.
125+
DROP TABLE pg_temp.project_backfill_roots;
6126
--> statement-breakpoint
7-
LOCK TABLE workspace, project, project_workspace, workflow IN SHARE ROW EXCLUSIVE MODE;
127+
SET statement_timeout = '60s';
8128
--> statement-breakpoint
9-
DO $$
129+
CREATE OR REPLACE FUNCTION pg_temp.validate_project_membership() RETURNS void LANGUAGE plpgsql AS $$
10130
BEGIN
131+
IF EXISTS (
132+
WITH RECURSIVE reachable(id) AS (
133+
SELECT id FROM workspace WHERE forked_from_workspace_id IS NULL
134+
UNION
135+
SELECT w.id FROM workspace w JOIN reachable r ON w.forked_from_workspace_id = r.id
136+
) SELECT 1 FROM workspace w LEFT JOIN reachable r ON r.id = w.id WHERE r.id IS NULL
137+
) THEN
138+
RAISE EXCEPTION 'Project backfill found a fork cycle or missing parent; reconcile before retrying' USING ERRCODE = '55000';
139+
END IF;
11140
IF EXISTS (SELECT 1 FROM workspace w LEFT JOIN project_workspace pw ON pw.workspace_id = w.id WHERE pw.workspace_id IS NULL) THEN
12-
RAISE EXCEPTION 'Project enforcement requires the reviewed Project backfill: missing environment memberships' USING ERRCODE = '55000';
141+
RAISE EXCEPTION 'Project enforcement found environments unreachable from a valid fork root or concurrently detached; reconcile and retry' USING ERRCODE = '55000';
13142
END IF;
14143
IF EXISTS (
15144
SELECT 1 FROM project p LEFT JOIN project_workspace pw ON pw.project_id = p.id
@@ -18,7 +147,7 @@ BEGIN
18147
OR (p.archived_at IS NULL AND count(w.id) FILTER (WHERE w.archived_at IS NULL) = 0)
19148
OR (p.archived_at IS NOT NULL AND count(w.id) FILTER (WHERE w.archived_at IS NULL) > 0)
20149
) THEN
21-
RAISE EXCEPTION 'Project enforcement requires nonempty Projects with consistent archive state; reconcile and verify the backfill' USING ERRCODE = '55000';
150+
RAISE EXCEPTION 'Project enforcement requires nonempty Projects with consistent archive state; reconcile and retry the migration' USING ERRCODE = '55000';
22151
END IF;
23152
IF EXISTS (
24153
SELECT 1 FROM project_workspace pw JOIN project p ON p.id = pw.project_id
@@ -40,6 +169,8 @@ BEGIN
40169
END;
41170
$$;
42171
--> statement-breakpoint
172+
SELECT pg_temp.validate_project_membership();
173+
--> statement-breakpoint
43174
CREATE OR REPLACE FUNCTION project_contract_assert_project(target_id text) RETURNS void
44175
LANGUAGE plpgsql AS $$
45176
DECLARE
@@ -183,6 +314,15 @@ BEGIN
183314
END;
184315
$$;
185316
--> statement-breakpoint
317+
BEGIN;
318+
--> statement-breakpoint
319+
SET LOCAL lock_timeout = '1s';
320+
--> statement-breakpoint
321+
SET LOCAL statement_timeout = '5s';
322+
--> statement-breakpoint
323+
-- Only trigger installation holds table locks. Both data scans run outside this transaction.
324+
LOCK TABLE workspace, project, project_workspace, workflow IN ACCESS EXCLUSIVE MODE NOWAIT;
325+
--> statement-breakpoint
186326
DROP TRIGGER IF EXISTS project_contract_lock ON project;
187327
CREATE TRIGGER project_contract_lock BEFORE INSERT OR UPDATE OR DELETE ON project
188328
FOR EACH ROW EXECUTE FUNCTION project_contract_before_write();
@@ -212,3 +352,10 @@ CREATE CONSTRAINT TRIGGER project_contract_check AFTER INSERT OR UPDATE OF works
212352
DEFERRABLE INITIALLY DEFERRED FOR EACH ROW EXECUTE FUNCTION project_contract_after_write();
213353
--> statement-breakpoint
214354
COMMIT;
355+
--> statement-breakpoint
356+
-- Installed triggers keep new writes valid while this read-only scan checks existing rows.
357+
SELECT pg_temp.validate_project_membership();
358+
--> statement-breakpoint
359+
DROP FUNCTION pg_temp.validate_project_membership();
360+
--> statement-breakpoint
361+
SET statement_timeout = 0;

0 commit comments

Comments
 (0)