From bd8678a8888a6cf71d3ba5badb7ab1654f593c55 Mon Sep 17 00:00:00 2001 From: whyour Date: Mon, 24 Aug 2026 21:13:15 +0800 Subject: [PATCH] feat(ql3): atomically apply cluster legacy env migration --- docs/QINGLONG_3_0_ARCHITECTURE_RFC.md | 26 + ...-config-reconciliation-and-task-binding.md | 4 +- ...luster-legacy-env-migration-plan-ledger.md | 4 + ...que-cluster-environment-bundle-delivery.md | 4 + ...luster-legacy-env-migration-application.md | 118 ++ docs/adr/README.md | 2 + .../test/application.test.cjs | 18 + .../test/bootstrap.test.cjs | 18 + packages/ql3-cluster-postgres/package.json | 5 + .../src/migration/migrationManifest.ts | 5 + .../src/migrations/index.ts | 2 + ...LegacyEnvMigrationApplicationRepository.ts | 1302 +++++++++++++++++ ...uster-legacy-env-migration-applications.ts | 182 +++ .../ql3-cluster-postgres/src/schema/schema.ts | 210 +++ .../src/schema/schemaContract.ts | 83 +- .../src/schema/schemaReadiness.ts | 39 + .../test/postgres.integration.test.cjs | 357 +++++ .../postgresqlMigrationDefinitions.test.cjs | 51 + .../test/postgresqlSchemaReadiness.test.cjs | 65 +- packages/ql3-runtime-core/package.json | 8 + .../clusterLegacyEnvMigrationApplication.ts | 660 +++++++++ ...sterLegacyEnvMigrationApplication.test.cjs | 222 +++ test/back/ql3PackageBoundaryAudit.test.cjs | 8 +- 23 files changed, 3372 insertions(+), 21 deletions(-) create mode 100644 docs/adr/ADR-0497-atomic-cluster-legacy-env-migration-application.md create mode 100644 packages/ql3-cluster-postgres/src/reconciliation/clusterLegacyEnvMigrationApplicationRepository.ts create mode 100644 packages/ql3-cluster-postgres/src/reconciliation/pg-0071-cluster-legacy-env-migration-applications.ts create mode 100644 packages/ql3-runtime-core/src/migration/clusterLegacyEnvMigrationApplication.ts create mode 100644 packages/ql3-runtime-core/test/clusterLegacyEnvMigrationApplication.test.cjs diff --git a/docs/QINGLONG_3_0_ARCHITECTURE_RFC.md b/docs/QINGLONG_3_0_ARCHITECTURE_RFC.md index bb9d9e3c..bc7fe349 100644 --- a/docs/QINGLONG_3_0_ARCHITECTURE_RFC.md +++ b/docs/QINGLONG_3_0_ARCHITECTURE_RFC.md @@ -148,6 +148,32 @@ arm64 physical HA 通过 146 gates、timeline `1→2`,报告 SHA-256 为 `0ee2199d0a52a02025bff017a07477d707353d12aefc3811f474e3775f2bb86b`。 +- D-402/ADR-0497(已验收):Cluster Legacy Env plan 现在可以在一个 Project-serialized + `SERIALIZABLE` transaction 中真正提交。profile-neutral contract 使用可重放、按 ID 排序的 + Task/Trigger mutation stream,并分别冻结 source revision-set 与 mutation-set digest;repository + 每批只保留 128 项,先迁移 Task,再迁移 Trigger,支持 100,000/500,000 上限而不在管理节点 + JS 堆保存全集。 + + Task current head、Plugin ownership、`command@v1` semantic、旧 digest 与数据库时间均在写前复验; + 新 Task revision 保留原字段,只追加固定 `environmentBundleRef`,enabled Task 同时生成 Cluster + execution revision。Trigger 复验自己的 current revision、旧 Task pin 和 schedule fence,随后固定 + 到同 application 的新 Task revision;合法的历史 pin 不要求等于迁移前 Task current head。 + Trigger revision、head CAS、schedule claim/due reset 和 fence 递增与 Task DML 原子提交。 + + `pg-0071-cluster-legacy-env-migration-applications` 将 PostgreSQL contract 推进到 v70,增加 application、 + Task item、Trigger item 三张 content-free append-only receipt 表。只有 Automation Manager 拥有 + `SELECT, INSERT`;runtime/admin 和其余角色无权限。exact replay 重新验证 durable heads、execution + revisions 和 schedules,但不会再次消费输入流。实现仍复用 runtime-core/cluster-postgres,没有新增 + package、生产依赖、daemon 或 Edge import。 + + Runtime Core `591/591`、Cluster PostgreSQL `361 total / 358 pass / 3 conditional skip / 0 fail`、 + v70 定向门 `74/74` 均为零失败。PostgreSQL 18.6 arm64 HA 再次通过 146 gates、timeline `1→2`, + 最终报告 SHA-256 为 `42ca97de43cfebd4611282b1fd5c0b09030eda89e88144967497902b01d18b3a`; + 真实用例覆盖 Trigger 固定 Task r1、Task current r2、原子生成 Task r3/Trigger r2 并重定向 pin, + 同时证明 bundle ref-only execution、schedule reset、无流消费 replay 与数据库角色隔离。D-402 + 关闭 mutation/receipt 边界;direct external custody、promotion 后 receipt replay 和固定低性能 Edge + 物理证据仍是 ADR-0491 转 Accepted 前的门禁。 + - D-396/ADR-0490(已验收):Run History 不再只有永久 `manual_external`,但也没有被错误实现为 Legacy 日志到 3.0 Run ledger 的回灌。 新的 Local adapter 以 ADR-0482 sealed capture bundle 作为 append-only 保全资产:Legacy history 必须逐事实选择 `retain_both`,Target history 必须选择 `retain_target`;receipt 只绑定 signed review、application、bundle fingerprint、领域 inventory 与有界 fact counts,不保存表名、Run ID、 diff --git a/docs/adr/ADR-0491-bounded-secret-config-reconciliation-and-task-binding.md b/docs/adr/ADR-0491-bounded-secret-config-reconciliation-and-task-binding.md index bb7d5d17..70713060 100644 --- a/docs/adr/ADR-0491-bounded-secret-config-reconciliation-and-task-binding.md +++ b/docs/adr/ADR-0491-bounded-secret-config-reconciliation-and-task-binding.md @@ -1,6 +1,6 @@ # ADR-0491:有界 Secret/Config Reconciliation 与任务环境绑定 -- 状态:Proposed(D-397 已实现 Legacy Env inspection、私有有界 row plan、durable plan publication、独立 signed decision、逐项 Automation adoption provenance、Local SQLite 原子 application publisher、Owner prepared/apply/rollback 编排、ADR-0492 completion v3;ADR-0494 完成 Cluster mounted-files provider live 子门,ADR-0495 完成 content-free Cluster plan ledger baseline,ADR-0496 完成 opaque environment bundle 的 Worker 内存展开与 HA-validated 数据面;真实 Edge 空间证据、Cluster Task/Trigger mutation/receipt、promotion 后 receipt replay 与直接外部 custody gate 尚未完成) +- 状态:Proposed(D-397 已实现 Legacy Env inspection、私有有界 row plan、durable plan publication、独立 signed decision、逐项 Automation adoption provenance、Local SQLite 原子 application publisher、Owner prepared/apply/rollback 编排、ADR-0492 completion v3;ADR-0494 完成 Cluster mounted-files provider live 子门,ADR-0495 完成 content-free Cluster plan ledger,ADR-0496 完成 opaque environment bundle 数据面,ADR-0497 完成 Cluster Task/Trigger 原子 mutation 与 receipt;真实 Edge 空间证据、promotion 后 receipt replay 与直接外部 custody gate 尚未完成) - 日期:2026-08-23 - 决策:D-397 - 关联:ADR-0073、ADR-0074、ADR-0092、ADR-0094、ADR-0480、ADR-0482、ADR-0483、ADR-0484、ADR-0485、ADR-0486、ADR-0487、ADR-0488、ADR-0490 @@ -147,4 +147,4 @@ D-397 当前八切片已经实现:absent、unsupported、Edge over-budget、2. ADR-0494 已完成 Cluster `mounted-files` provider live 子门:真实三节点 K3s 中两个 management replica、direct exact-key executor 和两个跨节点 provider observer 完成 PostgreSQL durable approval/binding、Kubernetes atomic projection rotation、无 Secret API 权限/ServiceAccount token、只读 `0440`、内容脱敏及删除后 fail-closed;v2 私有报告 24/24 gates 为 true,并保持 v1 verifier 兼容。该门不增加 Edge 闭包,也不等于直接 Vault/KMS/HSM custody。 -转为 Accepted 前仍必须完成:固定低性能 Edge 设备的真实空间/写放大/断电恢复证据,以及 Cluster Legacy Env migration 的逐项 Task/Trigger current-head revalidation、revision mutation/receipt、外部 custody adapter 和 HA promotion 后 receipt replay。ADR-0495 已完成专用 PostgreSQL SERIALIZABLE plan ledger baseline,但它只保存摘要、计数和 pinned SecretRef,不执行 Task/Trigger DML,也不接触 Secret material。ADR-0496 已验收只保存 pinned bundle ref、通过 fenced remote delivery 取回 typed carrier、在 Worker 内存展开的安全数据面,并通过 146-gate PostgreSQL HA;它仍不执行 migration DML 或生成 migration receipt。ADR-0492 已完成本机 completion schema 演进和 completed-head 后 rollback material 回收,ADR-0493 又让没有 Legacy 身份输入的 fresh v52 目标身份经 signed `retain_target` 正确形成 no-effect,并精确消除六张已知目标表的 `unknown` 误判。Legacy `Auths/Users` 或真正未知表仍保持 manual;本切片的 Local Owner 编排、ADR-0494 的 mounted-files gate、ADR-0495 的 plan ledger、ADR-0496 的数据面或 PostgreSQL HA 证据都不得冒充完整 Cluster migration 与外部密钥托管。 +转为 Accepted 前仍必须完成:固定低性能 Edge 设备的真实空间/写放大/断电恢复证据、直接外部 custody adapter,以及 HA promotion 后对既有 application receipt 的 exact replay。ADR-0495 已完成专用 PostgreSQL plan ledger,ADR-0496 已完成只保存 pinned bundle ref、通过 fenced remote delivery 取回 typed carrier 并在 Worker 内存展开的数据面;ADR-0497 又在一个 Project-serialized SERIALIZABLE transaction 中完成逐项 Task/Trigger current-head revalidation、revision/execution mutation、schedule reset 和 content-free append-only receipt,并支持合法历史 Task pin。它仍不写入 Secret material、不等于 direct Vault/KMS/HSM custody,也尚未在 promotion 后重放同一 application receipt。ADR-0492 已完成本机 completion schema 演进和 completed-head 后 rollback material 回收,ADR-0493 又让没有 Legacy 身份输入的 fresh v52 目标身份经 signed `retain_target` 正确形成 no-effect,并精确消除六张已知目标表的 `unknown` 误判。Legacy `Auths/Users` 或真正未知表仍保持 manual;Local Owner 编排、mounted-files gate、plan/application ledger 或通用 PostgreSQL HA 证据都不得冒充完整外部密钥托管。 diff --git a/docs/adr/ADR-0495-content-free-cluster-legacy-env-migration-plan-ledger.md b/docs/adr/ADR-0495-content-free-cluster-legacy-env-migration-plan-ledger.md index 719d64aa..932ae479 100644 --- a/docs/adr/ADR-0495-content-free-cluster-legacy-env-migration-plan-ledger.md +++ b/docs/adr/ADR-0495-content-free-cluster-legacy-env-migration-plan-ledger.md @@ -5,6 +5,10 @@ - 决策:D-400 - 关联:ADR-0104、ADR-0233、ADR-0259、ADR-0491、ADR-0494 +> 2026-08-24:ADR-0497/D-402 已完成本 ADR 所列的下一切片:在同一 Automation Manager +> SERIALIZABLE transaction 中复验并迁移 Task/Trigger,重置 schedule fence,并写入逐项 +> content-free receipt。本 ADR 的 plan ledger 决策保持不变。 + ## 背景 ADR-0491 已完成 Local SQLite 上的 Legacy Env 检查、人工裁决、Secret application、 diff --git a/docs/adr/ADR-0496-opaque-cluster-environment-bundle-delivery.md b/docs/adr/ADR-0496-opaque-cluster-environment-bundle-delivery.md index f889bb65..59a8704a 100644 --- a/docs/adr/ADR-0496-opaque-cluster-environment-bundle-delivery.md +++ b/docs/adr/ADR-0496-opaque-cluster-environment-bundle-delivery.md @@ -5,6 +5,10 @@ - 决策:D-401 前置切片 - 关联:ADR-0091、ADR-0092、ADR-0104、ADR-0113、ADR-0114、ADR-0491、ADR-0494、ADR-0495 +> 2026-08-24:ADR-0497/D-402 已完成本 ADR 所要求的 Task/Trigger mutation 与 migration +> receipt;本 ADR 继续定义 opaque bundle 的安全数据面,direct external custody 与 promotion +> 后 receipt replay 仍是后续门禁。 + ## 背景 ADR-0495 的 Cluster plan 只保存一个同 Project、固定 version 的 SecretRef,并明确禁止把 diff --git a/docs/adr/ADR-0497-atomic-cluster-legacy-env-migration-application.md b/docs/adr/ADR-0497-atomic-cluster-legacy-env-migration-application.md new file mode 100644 index 00000000..2a616a42 --- /dev/null +++ b/docs/adr/ADR-0497-atomic-cluster-legacy-env-migration-application.md @@ -0,0 +1,118 @@ +# ADR-0497:Cluster Legacy Env 的原子 Task/Trigger 迁移与只追加回执 + +- 状态:Accepted +- 日期:2026-08-24 +- 决策:D-402 +- 关联:ADR-0104、ADR-0233、ADR-0259、ADR-0491、ADR-0495、ADR-0496 + +## 背景 + +ADR-0495 冻结了无敏感内容的 Cluster migration plan,ADR-0496 又让 Task 与执行修订只需保存 +一个固定版本的 `environmentBundleRef`。此前仍缺少真正提交计划的 authority:如果分别调用现有 +Task 与 Trigger repository,会开启两个独立事务,出现 Task 已改而 Trigger、schedule 或 receipt +尚未提交的裂脑窗口;如果一次把 100,000 个 Task 和 500,000 个 Trigger 全部装入 JS 数组,又会 +让 Automation Manager 在低内存节点上不可用。 + +迁移还必须处理合法的历史 pin:Trigger 可能仍固定 Task r1,而 Task 当前 head 已推进到 r2。 +新 Trigger 应固定本次迁移产生的 Task r3,不能错误假设旧 Trigger pin 等于 Task current head。 + +## 决策 + +### 1. 使用可重放流与双摘要冻结输入 + +新增 profile-neutral application contract 与 +`qinglong/cluster-legacy-env-migration-application-receipt@v1`。调用方提供可重放的 Task/Trigger +mutation stream;每项包含 ordinal、实体身份、旧 revision/content digest 和独立 UUID mutation。 +Task 与 Trigger 分别计算 source revision-set digest 和 mutation-set digest:前者必须匹配已发布 +plan,后者必须匹配 application intent。 + +流必须从 ordinal 0 连续、按 Task/Trigger ID 严格递增,不允许重复、未知字段、非规范 UUID 或 +超过 100,000/500,000 项。Task/Trigger ID 延续现有定义契约的 128-byte、无控制字符边界,不用 +更窄的 ASCII 正则误伤合法历史数据。Trigger 可以为空,但 Task 至少一个。 + +### 2. 一个 Project-serialized SERIALIZABLE 事务完成全部写入 + +Automation Manager repository 先开启 `SERIALIZABLE`,取得 domain-separated Project advisory +transaction lock,再验证 active Project、plan ID/digest 和 exact mutation replay。每批最多只保留 +128 项,先完整处理 Task,再处理 Trigger;任何流摘要、head、ownership、spec 或 CAS 不一致都会 +回滚 receipt 和全部 revision DML。 + +Task 阶段对每个 current head 执行 `FOR UPDATE`,拒绝 Plugin-owned、非 `command@v1`、已有 +`environmentBundleRef`、revision/digest 漂移和数据库时间倒退;新 revision 保留 name、description、 +command config、labels、enabled 与 created time,只追加 plan 中的固定 bundle ref。enabled Task 同时 +生成新的 `remote_worker` Cluster execution revision,数据库只保存 ref,不保存 Env 名称或值。 + +Trigger 阶段复验 current Trigger revision、旧 Task pin 和 schedule revision,保留 spec 与 enabled, +但绑定同 application 中该 Task 的新 revision/content digest。新 Trigger revision、head CAS 与 +schedule reset 同事务提交;schedule 清空 due/claim 字段并递增 state/claim fence。旧 Trigger pin +可以是历史 revision,不要求等于 Task 迁移前 current head。 + +### 3. v70 提供三张只追加、无敏感内容的回执表 + +`pg-0071-cluster-legacy-env-migration-applications` 将 control contract 推进到 v70,新增: + +- `cluster_legacy_env_migration_application_receipts`:application/plan/mutation、四个 set digest、 + 固定 bundle ref、计数、数据库提交时间和精确 canonical receipt JSON; +- `cluster_legacy_env_migration_application_tasks`:每个 Task 的 before/after revision digest、mutation、 + 可选 execution digest 和 item digest; +- `cluster_legacy_env_migration_application_triggers`:每个 Trigger 的 before/after revision、旧 Task pin、 + 新 Task pin 和 item digest。 + +表中没有 Env name/value、bundle carrier、plaintext/ciphertext、key ID、Task/Trigger spec、命令或 +provider path。application/plan/mutation/receipt 唯一,ordinal 与 revision 关系有 named constraints; +子表通过 `(application_id, project_id)` 复合外键隔离 Project,Trigger item 又通过复合外键固定到 +同 application 的 Task item。三表均从 PUBLIC 撤销,只向 `ql3_automation_manager` 授予 +`SELECT, INSERT`,不授予 UPDATE/DELETE/TRUNCATE。 + +### 4. Exact replay 不重新消费输入流 + +相同 application mutation 先读取 durable receipt,精确比较 intent,然后聚合验证 Task heads、 +Task/execution revision、Trigger heads/revision 和 schedule 仍与逐项 receipt 一致。验证成功直接返回 +`existing`,不调用 Task/Trigger stream factory;intent drift、receipt 缺项、ordinal gap 或 current +head 已继续推进均失败关闭。序列化、死锁和 lock timeout 使用新 stream factory 最多重试三次。 + +### 5. 不为该能力新增微包或 Edge 常驻成本 + +纯契约位于既有 `runtime-core` 的显式 subpath,PostgreSQL authority 位于既有 +`cluster-postgres` 显式 subpath;不进入 runtime/admin/root entrypoint,不新增 workspace package、 +生产依赖、daemon、controller、timer、watcher 或全集缓存。低内存 Cluster 节点只承担固定 128 项 +batch;Edge/Standalone 默认 import graph 不加载 PostgreSQL authority。 + +## 被拒绝的替代方案 + +### 分别调用 Task 与 Trigger repository + +拒绝。两个事务无法保证 Task、Trigger、schedule 和 receipt 原子,response loss 也无法证明哪一半 +已经提交。 + +### 先收集全部候选再写数据库 + +拒绝。100,000/500,000 上限会把路由级设备或小型管理节点变成内存压力点;可重放有序流和固定 batch +已经能在事务回滚后重新计算摘要。 + +### 要求 Trigger 旧 Task pin 等于 Task current head + +拒绝。不可变 Trigger 合法固定历史 Task revision;只需复验该旧 pin,并把新 Trigger 显式重定向到 +本次 application 产生的新 Task revision。 + +### 把 Task/Trigger spec 或 Env 名称写进 receipt + +拒绝。revision 本身已经保存规范 spec;receipt 只需要 identity、digest、count 和 fence,复制内容会 +扩大 PostgreSQL backup、HA replica 和审计泄漏面。 + +## 当前验证与后续门禁 + +runtime-core 完整测试 `591/591`;cluster-postgres package 测试 `361 total / 358 pass / 3 +conditional skip / 0 fail`;v70 migration/schema/readiness 定向门 `74/74`。真实 PostgreSQL 18.6 +arm64 HA 多次通过 146 gates,timeline `1→2`;最终报告 SHA-256 为 +`42ca97de43cfebd4611282b1fd5c0b09030eda89e88144967497902b01d18b3a`。 + +真实数据库用例证明:带空格的合法 Task/Trigger ID 可迁移;Trigger 固定 Task r1、Task current r2 +时会原子生成 Task r3 与 Trigger r2 并让新 Trigger 固定 r3;execution plan 只出现 bundle ref; +schedule fence 递增且 claim/due 清空;exact replay 不消费 stream;Automation Manager UPDATE 以及 +runtime/admin SELECT 均以 `42501` 被拒绝。 + +本 ADR 关闭 D-402 的 Cluster Task/Trigger application/receipt 边界,但 ADR-0491 仍保持 Proposed。 +后续仍必须完成 direct Vault/KMS/HSM custody、migration Job 的短期身份与装配、在 PostgreSQL +promotion **之后**对既有 application receipt 执行 exact replay,以及固定低性能 Edge 设备的真实 +空间、写放大、断电与恢复证据。 diff --git a/docs/adr/README.md b/docs/adr/README.md index 8af13ade..6350e2fd 100644 --- a/docs/adr/README.md +++ b/docs/adr/README.md @@ -499,6 +499,8 @@ | [ADR-0493](./ADR-0493-target-native-reconciliation-domain-preservation.md) | 目标原生 Reconciliation 域分类与身份保留 | Accepted | | [ADR-0494](./ADR-0494-postgresql-secret-binding-and-mounted-provider-live-rotation.md) | PostgreSQL Secret Binding 与 Mounted Provider 在线轮换门 | Accepted | | [ADR-0495](./ADR-0495-content-free-cluster-legacy-env-migration-plan-ledger.md) | 无敏感内容的 Cluster Legacy Env 迁移计划账本 | Accepted | +| [ADR-0496](./ADR-0496-opaque-cluster-environment-bundle-delivery.md) | Cluster 不透明环境 Bundle 的有界交付与 Worker 内存展开 | Accepted | +| [ADR-0497](./ADR-0497-atomic-cluster-legacy-env-migration-application.md) | Cluster Legacy Env 的原子 Task/Trigger 迁移与只追加回执 | Accepted | ## 规则 diff --git a/packages/ql3-cluster-control/test/application.test.cjs b/packages/ql3-cluster-control/test/application.test.cjs index c7f11599..228222c9 100644 --- a/packages/ql3-cluster-control/test/application.test.cjs +++ b/packages/ql3-cluster-control/test/application.test.cjs @@ -218,6 +218,24 @@ function runtimePrivileges() { plugin_package_workflow_admission_steps: [true, true, false, false], plugin_package_workflow_task_attempt_admissions: [true, true, false, false], cluster_legacy_env_migration_plans: [false, false, false, false], + cluster_legacy_env_migration_application_receipts: [ + false, + false, + false, + false, + ], + cluster_legacy_env_migration_application_tasks: [ + false, + false, + false, + false, + ], + cluster_legacy_env_migration_application_triggers: [ + false, + false, + false, + false, + ], worker_execution_attestations: [true, false, false, false], run_events: [true, true, false, false], run_cancellation_dispatches: [true, true, true, false], diff --git a/packages/ql3-cluster-control/test/bootstrap.test.cjs b/packages/ql3-cluster-control/test/bootstrap.test.cjs index ef2cc429..c8fd92c7 100644 --- a/packages/ql3-cluster-control/test/bootstrap.test.cjs +++ b/packages/ql3-cluster-control/test/bootstrap.test.cjs @@ -132,6 +132,24 @@ function runtimePrivileges() { plugin_package_workflow_admission_steps: [true, true, false, false], plugin_package_workflow_task_attempt_admissions: [true, true, false, false], cluster_legacy_env_migration_plans: [false, false, false, false], + cluster_legacy_env_migration_application_receipts: [ + false, + false, + false, + false, + ], + cluster_legacy_env_migration_application_tasks: [ + false, + false, + false, + false, + ], + cluster_legacy_env_migration_application_triggers: [ + false, + false, + false, + false, + ], worker_execution_attestations: [true, false, false, false], run_events: [true, true, false, false], run_cancellation_dispatches: [true, true, true, false], diff --git a/packages/ql3-cluster-postgres/package.json b/packages/ql3-cluster-postgres/package.json index 5ba1296a..15f64e59 100644 --- a/packages/ql3-cluster-postgres/package.json +++ b/packages/ql3-cluster-postgres/package.json @@ -70,6 +70,11 @@ "require": "./dist/reconciliation/clusterLegacyEnvMigrationPlanRepository.js", "default": "./dist/reconciliation/clusterLegacyEnvMigrationPlanRepository.js" }, + "./cluster-legacy-env-migration-application": { + "types": "./dist/reconciliation/clusterLegacyEnvMigrationApplicationRepository.d.ts", + "require": "./dist/reconciliation/clusterLegacyEnvMigrationApplicationRepository.js", + "default": "./dist/reconciliation/clusterLegacyEnvMigrationApplicationRepository.js" + }, "./task-start": { "types": "./dist/task-start/taskStartRepository.d.ts", "require": "./dist/task-start/taskStartRepository.js", diff --git a/packages/ql3-cluster-postgres/src/migration/migrationManifest.ts b/packages/ql3-cluster-postgres/src/migration/migrationManifest.ts index 0fefdff7..eeded7ad 100644 --- a/packages/ql3-cluster-postgres/src/migration/migrationManifest.ts +++ b/packages/ql3-cluster-postgres/src/migration/migrationManifest.ts @@ -358,5 +358,10 @@ export const postgresqlMainMigrationManifest: MigrationStreamManifest = checksum: '7cd6d993f48e7bcebcd62c93571a738d5117c9bcde33b974c5ac8962e2a03fe4', }), + Object.freeze({ + id: 'pg-0071-cluster-legacy-env-migration-applications', + checksum: + '82538b5a244011a22afad7f9c6d266a8997da538bdd2bec2ee524997c3b85996', + }), ]), }); diff --git a/packages/ql3-cluster-postgres/src/migrations/index.ts b/packages/ql3-cluster-postgres/src/migrations/index.ts index 6f26f566..47a3b270 100644 --- a/packages/ql3-cluster-postgres/src/migrations/index.ts +++ b/packages/ql3-cluster-postgres/src/migrations/index.ts @@ -73,6 +73,7 @@ import { pg0067CancellationDispatchManagementMigration } from '../run-management import { pg0068CancellationDispatchProjectKeysetMigration } from '../run-management/pg-0068-cancellation-dispatch-project-keyset'; import { pg0069WorkerSessionManagementObservationMigration } from '../remote-execution/pg-0069-worker-session-management-observation'; import { pg0070ClusterLegacyEnvMigrationPlansMigration } from '../reconciliation/pg-0070-cluster-legacy-env-migration-plans'; +import { pg0071ClusterLegacyEnvMigrationApplicationsMigration } from '../reconciliation/pg-0071-cluster-legacy-env-migration-applications'; export const postgresqlMainMigrationStream: MigrationStreamDefinition = Object.freeze({ @@ -151,5 +152,6 @@ export const postgresqlMainMigrationStream: MigrationStreamDefinition; +type Queryable = Pick | Pick; + +const APPLICATION_BATCH_SIZE = 128; +const TASK_ITEM_DIGEST_DOMAIN = Buffer.from( + 'qinglong/cluster-legacy-env-migration-application-task-item@v1\0', + 'utf8', +); +const TRIGGER_ITEM_DIGEST_DOMAIN = Buffer.from( + 'qinglong/cluster-legacy-env-migration-application-trigger-item@v1\0', + 'utf8', +); + +const TASK_SELECT_FIELDS = ` + head.project_id AS "projectId", + head.task_id AS "taskId", + revision.revision, + revision.mutation_id AS "mutationId", + revision.name, + revision.description, + revision.kind, + revision.spec_json AS "specJson", + revision.labels_json AS "labelsJson", + revision.enabled, + revision.content_digest AS "contentDigest", + head.created_at_ms AS "createdAtMs", + revision.created_at_ms AS "updatedAtMs"`; + +const TRIGGER_SELECT_FIELDS = ` + head.project_id AS "projectId", + head.trigger_id AS "triggerId", + revision.revision, + revision.mutation_id AS "mutationId", + revision.task_id AS "taskId", + revision.task_revision AS "taskRevision", + revision.task_content_digest AS "taskContentDigest", + revision.spec_json AS "specJson", + revision.enabled, + revision.content_digest AS "contentDigest", + head.created_at_ms AS "createdAtMs", + revision.created_at_ms AS "updatedAtMs"`; + +export interface PostgresClusterLegacyEnvMigrationApplicationTransactionContext { + readonly intent: Readonly; + readonly replay: Readonly | null; + readonly receipt: Readonly; +} + +export type PostgresClusterLegacyEnvMigrationApplicationTransactionHook = ( + client: PostgresClient, + context: Readonly, +) => Promise; + +interface AppliedTask { + readonly source: Readonly; + readonly definition: Readonly; + readonly execution: Readonly | null; + readonly itemDigest: string; +} + +interface AppliedTrigger { + readonly source: Readonly; + readonly trigger: Readonly; + readonly itemDigest: string; +} + +function conflict(): ClusterLegacyEnvMigrationApplicationConflictError { + return new ClusterLegacyEnvMigrationApplicationConflictError(); +} + +function unavailable(): ClusterLegacyEnvMigrationApplicationUnavailableError { + return new ClusterLegacyEnvMigrationApplicationUnavailableError(); +} + +function invalid( + message: string, +): InvalidClusterLegacyEnvMigrationApplicationError { + return new InvalidClusterLegacyEnvMigrationApplicationError(message); +} + +function contentDigest(domain: Buffer, value: object): string { + return createHash('sha256') + .update(domain) + .update(JSON.stringify(value), 'utf8') + .digest('hex'); +} + +function taskItemDigest( + applicationId: string, + source: ClusterLegacyEnvMigrationTaskMutation, + definition: TaskDefinitionRecord, + execution: ClusterTaskExecutionRevision | null, +): string { + return contentDigest(TASK_ITEM_DIGEST_DOMAIN, { + applicationId, + ordinal: source.ordinal, + projectId: definition.projectId, + taskId: definition.taskId, + previousRevision: source.previousRevision, + previousContentDigest: source.previousContentDigest, + mutationId: source.mutationId, + revision: definition.revision, + contentDigest: definition.contentDigest, + executionContentDigest: execution?.contentDigest ?? null, + }); +} + +function triggerItemDigest( + applicationId: string, + source: ClusterLegacyEnvMigrationTriggerMutation, + trigger: TriggerRecord, +): string { + return contentDigest(TRIGGER_ITEM_DIGEST_DOMAIN, { + applicationId, + ordinal: source.ordinal, + projectId: trigger.projectId, + triggerId: trigger.triggerId, + taskId: trigger.taskId, + previousRevision: source.previousRevision, + previousContentDigest: source.previousContentDigest, + previousTaskRevision: source.previousTaskRevision, + previousTaskContentDigest: source.previousTaskContentDigest, + mutationId: source.mutationId, + revision: trigger.revision, + contentDigest: trigger.contentDigest, + taskRevision: trigger.taskRevision, + taskContentDigest: trigger.taskContentDigest, + }); +} + +function taskRecord(row: Row): Readonly { + try { + const description = row.description; + if (description !== null && typeof description !== 'string') { + throw unavailable(); + } + return normalizeTaskDefinitionRecord({ + projectId: postgresRequiredString(row.projectId, unavailable), + taskId: postgresRequiredString(row.taskId, unavailable), + revision: postgresRequiredInteger(row.revision, unavailable), + mutationId: postgresRequiredString(row.mutationId, unavailable), + name: postgresRequiredString(row.name, unavailable), + ...(description === null ? {} : { description }), + kind: postgresRequiredString( + row.kind, + unavailable, + ) as TaskDefinitionRecord['kind'], + spec: postgresRequiredJsonObject( + row.specJson, + unavailable, + ) as unknown as TaskDefinitionRecord['spec'], + labels: postgresRequiredJsonObject( + row.labelsJson, + unavailable, + ) as TaskDefinitionRecord['labels'], + enabled: postgresRequiredBoolean(row.enabled, unavailable), + contentDigest: postgresRequiredString(row.contentDigest, unavailable), + createdAtMs: postgresRequiredInteger(row.createdAtMs, unavailable), + updatedAtMs: postgresRequiredInteger(row.updatedAtMs, unavailable), + }); + } catch (error) { + if (error instanceof ClusterLegacyEnvMigrationApplicationUnavailableError) { + throw error; + } + throw unavailable(); + } +} + +function triggerRecord(row: Row): Readonly { + try { + return normalizeTriggerRecord({ + projectId: postgresRequiredString(row.projectId, unavailable), + triggerId: postgresRequiredString(row.triggerId, unavailable), + revision: postgresRequiredInteger(row.revision, unavailable), + mutationId: postgresRequiredString(row.mutationId, unavailable), + taskId: postgresRequiredString(row.taskId, unavailable), + taskRevision: postgresRequiredInteger(row.taskRevision, unavailable), + taskContentDigest: postgresRequiredString( + row.taskContentDigest, + unavailable, + ), + spec: postgresRequiredJsonObject( + row.specJson, + unavailable, + ) as unknown as TriggerRecord['spec'], + enabled: postgresRequiredBoolean(row.enabled, unavailable), + contentDigest: postgresRequiredString(row.contentDigest, unavailable), + createdAtMs: postgresRequiredInteger(row.createdAtMs, unavailable), + updatedAtMs: postgresRequiredInteger(row.updatedAtMs, unavailable), + }); + } catch (error) { + if (error instanceof ClusterLegacyEnvMigrationApplicationUnavailableError) { + throw error; + } + throw unavailable(); + } +} + +function planFromRow(row: Row): Readonly { + try { + return normalizeClusterLegacyEnvMigrationPlan( + postgresRequiredJsonObject( + row.planJson, + unavailable, + ) as unknown as ClusterLegacyEnvMigrationPlan, + ); + } catch (error) { + if (error instanceof ClusterLegacyEnvMigrationApplicationUnavailableError) { + throw error; + } + throw unavailable(); + } +} + +function receiptFromRow( + row: Row, +): Readonly { + try { + return normalizeClusterLegacyEnvMigrationApplicationReceipt( + postgresRequiredJsonObject( + row.receiptJson, + unavailable, + ) as unknown as ClusterLegacyEnvMigrationApplicationReceipt, + ); + } catch (error) { + if (error instanceof ClusterLegacyEnvMigrationApplicationUnavailableError) { + throw error; + } + throw unavailable(); + } +} + +async function findReceipt( + queryable: Queryable, + column: 'application_id' | 'mutation_id', + value: string, +): Promise | null> { + const result = await queryable.query( + `SELECT receipt_json AS "receiptJson" + FROM "ql3"."cluster_legacy_env_migration_application_receipts" + WHERE ${column} = $1 + LIMIT 2`, + [value], + ); + if (result.rows.length === 0) return null; + if (result.rows.length !== 1) throw unavailable(); + const receipt = receiptFromRow(result.rows[0]!); + if ( + (column === 'application_id' && receipt.applicationId !== value) || + (column === 'mutation_id' && receipt.mutationId !== value) + ) { + throw unavailable(); + } + return receipt; +} + +async function loadPlan( + client: PostgresClient, + intent: ClusterLegacyEnvMigrationApplicationIntent, +): Promise> { + const result = await client.query( + `SELECT plan_json AS "planJson" + FROM "ql3"."cluster_legacy_env_migration_plans" + WHERE plan_id = $1 + FOR SHARE`, + [intent.planId], + ); + if (result.rows.length !== 1) throw conflict(); + const plan = planFromRow(result.rows[0]!); + if ( + plan.planId !== intent.planId || + plan.projectId !== intent.projectId || + plan.planDigest !== intent.planDigest + ) { + throw conflict(); + } + return plan; +} + +function normalizeStreams( + value: Readonly, +): Readonly { + if ( + !value || + typeof value !== 'object' || + Array.isArray(value) || + (Object.getPrototypeOf(value) !== Object.prototype && + Object.getPrototypeOf(value) !== null) + ) { + throw invalid('mutation streams are invalid'); + } + const keys = Reflect.ownKeys(value); + const descriptors = Object.getOwnPropertyDescriptors(value); + const taskMutations = descriptors.taskMutations; + const triggerMutations = descriptors.triggerMutations; + if ( + keys.some((key) => typeof key !== 'string') || + keys.length !== 2 || + !keys.includes('taskMutations') || + !keys.includes('triggerMutations') || + taskMutations?.enumerable !== true || + taskMutations.get !== undefined || + taskMutations.set !== undefined || + typeof taskMutations.value !== 'function' || + triggerMutations?.enumerable !== true || + triggerMutations.get !== undefined || + triggerMutations.set !== undefined || + typeof triggerMutations.value !== 'function' + ) { + throw invalid('mutation stream shape is invalid'); + } + return Object.freeze({ + taskMutations: taskMutations.value, + triggerMutations: triggerMutations.value, + }); +} + +async function* mutationStream( + factory: () => Iterable | AsyncIterable, + label: string, +): AsyncIterable { + let source: Iterable | AsyncIterable; + try { + source = factory(); + } catch { + throw invalid(`${label} factory failed`); + } + if (!source || typeof source !== 'object') { + throw invalid(`${label} must be iterable`); + } + const asyncFactory = (source as AsyncIterable)[Symbol.asyncIterator]; + const syncFactory = (source as Iterable)[Symbol.iterator]; + let iterator: AsyncIterator | Iterator; + try { + iterator = + typeof asyncFactory === 'function' + ? asyncFactory.call(source) + : typeof syncFactory === 'function' + ? syncFactory.call(source) + : (() => { + throw invalid(`${label} must be iterable`); + })(); + } catch (error) { + if (error instanceof InvalidClusterLegacyEnvMigrationApplicationError) { + throw error; + } + throw invalid(`${label} iterator is invalid`); + } + while (true) { + let step: IteratorResult; + try { + step = await iterator.next(); + } catch { + throw invalid(`${label} iteration failed`); + } + if (!step || typeof step !== 'object' || typeof step.done !== 'boolean') { + throw invalid(`${label} iterator result is invalid`); + } + if (step.done) return; + yield step.value; + } +} + +function clusterExecutionPlanJson( + revision: ClusterTaskExecutionRevision, +): Readonly> { + return Object.freeze({ + command: revision.command, + environment: revision.environment, + ...(revision.environmentBundleRef === undefined + ? {} + : { environmentBundleRef: revision.environmentBundleRef }), + ...(revision.workingDirectory === undefined + ? {} + : { workingDirectory: revision.workingDirectory }), + ...(revision.timeoutMs === undefined + ? {} + : { timeoutMs: revision.timeoutMs }), + ...(revision.placement === undefined + ? {} + : { placement: revision.placement }), + }); +} + +function requireRowCount( + value: number | null | undefined, + expected: number, +): void { + if (value !== expected) throw conflict(); +} + +function indexRowsByStringId( + rows: readonly Row[], + field: 'taskId' | 'triggerId', +): ReadonlyMap { + const indexed = new Map(); + for (const row of rows) { + const id = postgresRequiredString(row[field], unavailable); + if (indexed.has(id)) throw unavailable(); + indexed.set(id, row); + } + return indexed; +} + +async function insertTaskBatch( + client: PostgresClient, + applicationId: string, + projectId: string, + environmentBundleRef: string, + committedAtMs: number, + batch: readonly ClusterLegacyEnvMigrationTaskMutation[], +): Promise { + const ids = batch.map((item) => item.taskId); + const loaded = await client.query( + `SELECT ${TASK_SELECT_FIELDS}, + EXISTS ( + SELECT 1 + FROM "ql3"."plugin_package_task_ownerships" AS ownership + WHERE ownership.project_id = head.project_id + AND ownership.task_id = head.task_id + ) AS "pluginOwned" + FROM "ql3"."task_definitions" AS head + JOIN "ql3"."task_definition_revisions" AS revision + ON revision.project_id = head.project_id + AND revision.task_id = head.task_id + AND revision.revision = head.current_revision + WHERE head.project_id = $1 AND head.task_id = ANY($2::varchar[]) + FOR UPDATE OF head`, + [projectId, ids], + ); + if (loaded.rows.length !== batch.length) throw conflict(); + const rowsByTaskId = indexRowsByStringId(loaded.rows, 'taskId'); + + const semanticRegistry = createBuiltInTaskSpecSemanticRegistry(); + const applied: AppliedTask[] = []; + for (const source of batch) { + const row = rowsByTaskId.get(source.taskId); + if (row === undefined) throw conflict(); + const current = taskRecord(row); + if ( + current.projectId !== projectId || + current.taskId !== source.taskId || + current.revision !== source.previousRevision || + current.contentDigest !== source.previousContentDigest || + row.pluginOwned !== false || + current.kind !== 'command' || + current.spec.schema !== BUILT_IN_COMMAND_TASK_SPEC_SCHEMA || + Object.hasOwn(current.spec.config, 'environmentBundleRef') || + committedAtMs < current.updatedAtMs + ) { + throw conflict(); + } + let definition: Readonly; + let execution: Readonly | null; + try { + const spec = semanticRegistry.normalize({ + projectId, + taskId: current.taskId, + kind: 'command', + spec: { + schema: BUILT_IN_COMMAND_TASK_SPEC_SCHEMA, + config: { + ...current.spec.config, + environmentBundleRef, + }, + } as TaskDefinitionSpec, + }); + definition = createTaskDefinitionRecord( + { + projectId, + taskId: current.taskId, + expectedRevision: current.revision, + mutationId: source.mutationId, + name: current.name, + ...(current.description === undefined + ? {} + : { description: current.description }), + kind: current.kind, + spec, + labels: current.labels, + enabled: current.enabled, + occurredAtMs: committedAtMs, + }, + current.createdAtMs, + ); + execution = definition.enabled + ? compileClusterCommandTaskDefinition(definition, semanticRegistry) + : null; + } catch { + throw conflict(); + } + applied.push({ + source, + definition, + execution, + itemDigest: taskItemDigest(applicationId, source, definition, execution), + }); + } + + const revisions = applied.map(({ definition }) => ({ + project_id: definition.projectId, + task_id: definition.taskId, + revision: definition.revision, + mutation_id: definition.mutationId, + name: definition.name, + description: definition.description ?? null, + kind: definition.kind, + spec_json: definition.spec, + labels_json: definition.labels, + enabled: definition.enabled, + content_digest: definition.contentDigest, + created_at_ms: definition.updatedAtMs, + })); + const revisionInsert = await client.query( + `INSERT INTO "ql3"."task_definition_revisions" ( + project_id, task_id, revision, mutation_id, name, description, kind, + spec_json, labels_json, enabled, content_digest, created_at_ms + ) + SELECT project_id, task_id, revision, mutation_id, name, description, + kind, spec_json, labels_json, enabled, content_digest, created_at_ms + FROM jsonb_to_recordset($1::jsonb) AS data( + project_id varchar(128), task_id varchar(128), revision integer, + mutation_id uuid, name varchar(255), description varchar(4096), + kind varchar(16), spec_json jsonb, labels_json jsonb, enabled boolean, + content_digest char(64), created_at_ms bigint + )`, + [JSON.stringify(revisions)], + ); + requireRowCount(revisionInsert.rowCount, applied.length); + + const executions = applied + .map(({ execution }) => execution) + .filter( + (value): value is Readonly => + value !== null, + ) + .map((execution) => ({ + project_id: execution.projectId, + task_id: execution.taskId, + source_revision: execution.sourceRevision, + task_revision: execution.taskRevision, + source_content_digest: execution.sourceContentDigest, + executor_type: execution.executorType, + plan_schema: execution.planSchema, + plan_json: clusterExecutionPlanJson(execution), + content_digest: execution.contentDigest, + created_at_ms: execution.createdAtMs, + })); + if (executions.length > 0) { + const executionInsert = await client.query( + `INSERT INTO "ql3"."task_execution_revisions" ( + project_id, task_id, source_revision, task_revision, + source_content_digest, executor_type, plan_schema, plan_json, + content_digest, created_at_ms + ) + SELECT project_id, task_id, source_revision, task_revision, + source_content_digest, executor_type, plan_schema, plan_json, + content_digest, created_at_ms + FROM jsonb_to_recordset($1::jsonb) AS data( + project_id varchar(128), task_id varchar(128), source_revision integer, + task_revision varchar(96), source_content_digest char(64), + executor_type varchar(32), plan_schema varchar(64), plan_json jsonb, + content_digest char(64), created_at_ms bigint + )`, + [JSON.stringify(executions)], + ); + requireRowCount(executionInsert.rowCount, executions.length); + } + + const heads = applied.map(({ source, definition }) => ({ + task_id: definition.taskId, + previous_revision: source.previousRevision, + revision: definition.revision, + updated_at_ms: definition.updatedAtMs, + })); + const headUpdate = await client.query( + `UPDATE "ql3"."task_definitions" AS head + SET current_revision = data.revision, + updated_at_ms = data.updated_at_ms + FROM jsonb_to_recordset($1::jsonb) AS data( + task_id varchar(128), previous_revision integer, + revision integer, updated_at_ms bigint + ) + WHERE head.project_id = $2 + AND head.task_id = data.task_id + AND head.current_revision = data.previous_revision`, + [JSON.stringify(heads), projectId], + ); + requireRowCount(headUpdate.rowCount, applied.length); + + const ledger = applied.map( + ({ source, definition, execution, itemDigest }) => ({ + application_id: applicationId, + ordinal: source.ordinal, + project_id: projectId, + task_id: source.taskId, + previous_revision: source.previousRevision, + previous_content_digest: source.previousContentDigest, + mutation_id: source.mutationId, + revision: definition.revision, + content_digest: definition.contentDigest, + execution_content_digest: execution?.contentDigest ?? null, + item_digest: itemDigest, + }), + ); + const ledgerInsert = await client.query( + `INSERT INTO "ql3"."cluster_legacy_env_migration_application_tasks" ( + application_id, ordinal, project_id, task_id, previous_revision, + previous_content_digest, mutation_id, revision, content_digest, + execution_content_digest, item_digest + ) + SELECT application_id, ordinal, project_id, task_id, previous_revision, + previous_content_digest, mutation_id, revision, content_digest, + execution_content_digest, item_digest + FROM jsonb_to_recordset($1::jsonb) AS data( + application_id varchar(128), ordinal integer, project_id varchar(128), + task_id varchar(128), previous_revision integer, + previous_content_digest char(64), mutation_id uuid, revision integer, + content_digest char(64), execution_content_digest char(64), + item_digest char(64) + )`, + [JSON.stringify(ledger)], + ); + requireRowCount(ledgerInsert.rowCount, applied.length); +} + +async function insertTriggerBatch( + client: PostgresClient, + applicationId: string, + projectId: string, + committedAtMs: number, + batch: readonly ClusterLegacyEnvMigrationTriggerMutation[], +): Promise { + const ids = batch.map((item) => item.triggerId); + const loaded = await client.query( + `SELECT ${TRIGGER_SELECT_FIELDS}, + schedule.trigger_revision AS "scheduleRevision", + migrated.revision AS "migratedTaskRevision", + migrated.content_digest AS "migratedTaskContentDigest" + FROM "ql3"."triggers" AS head + JOIN "ql3"."trigger_revisions" AS revision + ON revision.project_id = head.project_id + AND revision.trigger_id = head.trigger_id + AND revision.revision = head.current_revision + JOIN "ql3"."trigger_schedules" AS schedule + ON schedule.project_id = head.project_id + AND schedule.trigger_id = head.trigger_id + JOIN "ql3"."cluster_legacy_env_migration_application_tasks" AS migrated + ON migrated.application_id = $3 + AND migrated.project_id = head.project_id + AND migrated.task_id = revision.task_id + WHERE head.project_id = $1 AND head.trigger_id = ANY($2::varchar[]) + FOR UPDATE OF head, schedule`, + [projectId, ids, applicationId], + ); + if (loaded.rows.length !== batch.length) throw conflict(); + const rowsByTriggerId = indexRowsByStringId(loaded.rows, 'triggerId'); + + const semanticRegistry = createBuiltInTriggerSpecSemanticRegistry(); + const applied: AppliedTrigger[] = []; + for (const source of batch) { + const row = rowsByTriggerId.get(source.triggerId); + if (row === undefined) throw conflict(); + const current = triggerRecord(row); + const migratedTaskRevision = postgresRequiredInteger( + row.migratedTaskRevision, + unavailable, + ); + const migratedTaskContentDigest = postgresRequiredString( + row.migratedTaskContentDigest, + unavailable, + ); + if ( + current.projectId !== projectId || + current.triggerId !== source.triggerId || + current.taskId !== source.taskId || + current.revision !== source.previousRevision || + current.contentDigest !== source.previousContentDigest || + current.taskRevision !== source.previousTaskRevision || + current.taskContentDigest !== source.previousTaskContentDigest || + postgresRequiredInteger(row.scheduleRevision, unavailable) !== + current.revision || + committedAtMs < current.updatedAtMs + ) { + throw conflict(); + } + let trigger: Readonly; + try { + const spec = semanticRegistry.normalize({ + projectId, + triggerId: current.triggerId, + taskId: current.taskId, + taskRevision: migratedTaskRevision, + spec: current.spec, + }); + trigger = createTriggerRecord( + { + projectId, + triggerId: current.triggerId, + expectedRevision: current.revision, + mutationId: source.mutationId, + taskId: current.taskId, + taskRevision: migratedTaskRevision, + taskContentDigest: migratedTaskContentDigest, + spec, + enabled: current.enabled, + occurredAtMs: committedAtMs, + }, + current.createdAtMs, + ); + } catch { + throw conflict(); + } + applied.push({ + source, + trigger, + itemDigest: triggerItemDigest(applicationId, source, trigger), + }); + } + + const revisions = applied.map(({ trigger }) => ({ + project_id: trigger.projectId, + trigger_id: trigger.triggerId, + revision: trigger.revision, + mutation_id: trigger.mutationId, + task_id: trigger.taskId, + task_revision: trigger.taskRevision, + task_content_digest: trigger.taskContentDigest, + spec_json: trigger.spec, + enabled: trigger.enabled, + content_digest: trigger.contentDigest, + created_at_ms: trigger.updatedAtMs, + })); + const revisionInsert = await client.query( + `INSERT INTO "ql3"."trigger_revisions" ( + project_id, trigger_id, revision, mutation_id, task_id, task_revision, + task_content_digest, spec_json, enabled, content_digest, created_at_ms + ) + SELECT project_id, trigger_id, revision, mutation_id, task_id, + task_revision, task_content_digest, spec_json, enabled, + content_digest, created_at_ms + FROM jsonb_to_recordset($1::jsonb) AS data( + project_id varchar(128), trigger_id varchar(128), revision integer, + mutation_id uuid, task_id varchar(128), task_revision integer, + task_content_digest char(64), spec_json jsonb, enabled boolean, + content_digest char(64), created_at_ms bigint + )`, + [JSON.stringify(revisions)], + ); + requireRowCount(revisionInsert.rowCount, applied.length); + + const schedules = applied.map(({ source, trigger }) => ({ + trigger_id: trigger.triggerId, + previous_revision: source.previousRevision, + revision: trigger.revision, + updated_at_ms: trigger.updatedAtMs, + })); + const scheduleUpdate = await client.query( + `UPDATE "ql3"."trigger_schedules" AS schedule + SET trigger_revision = data.revision, + next_fire_at_ms = NULL, + last_scheduled_at_ms = NULL, + state_version = schedule.state_version + 1, + claim_owner = NULL, + claim_token = NULL, + claim_version = schedule.claim_version + 1, + claim_expires_at_ms = NULL, + updated_at_ms = data.updated_at_ms + FROM jsonb_to_recordset($1::jsonb) AS data( + trigger_id varchar(128), previous_revision integer, + revision integer, updated_at_ms bigint + ) + WHERE schedule.project_id = $2 + AND schedule.trigger_id = data.trigger_id + AND schedule.trigger_revision = data.previous_revision`, + [JSON.stringify(schedules), projectId], + ); + requireRowCount(scheduleUpdate.rowCount, applied.length); + + const heads = applied.map(({ source, trigger }) => ({ + trigger_id: trigger.triggerId, + task_id: trigger.taskId, + previous_revision: source.previousRevision, + revision: trigger.revision, + updated_at_ms: trigger.updatedAtMs, + })); + const headUpdate = await client.query( + `UPDATE "ql3"."triggers" AS head + SET current_revision = data.revision, + updated_at_ms = data.updated_at_ms + FROM jsonb_to_recordset($1::jsonb) AS data( + trigger_id varchar(128), task_id varchar(128), + previous_revision integer, revision integer, updated_at_ms bigint + ) + WHERE head.project_id = $2 + AND head.trigger_id = data.trigger_id + AND head.task_id = data.task_id + AND head.current_revision = data.previous_revision`, + [JSON.stringify(heads), projectId], + ); + requireRowCount(headUpdate.rowCount, applied.length); + + const ledger = applied.map(({ source, trigger, itemDigest }) => ({ + application_id: applicationId, + ordinal: source.ordinal, + project_id: projectId, + trigger_id: source.triggerId, + task_id: source.taskId, + previous_revision: source.previousRevision, + previous_content_digest: source.previousContentDigest, + previous_task_revision: source.previousTaskRevision, + previous_task_content_digest: source.previousTaskContentDigest, + mutation_id: source.mutationId, + revision: trigger.revision, + content_digest: trigger.contentDigest, + task_revision: trigger.taskRevision, + task_content_digest: trigger.taskContentDigest, + item_digest: itemDigest, + })); + const ledgerInsert = await client.query( + `INSERT INTO "ql3"."cluster_legacy_env_migration_application_triggers" ( + application_id, ordinal, project_id, trigger_id, task_id, + previous_revision, previous_content_digest, previous_task_revision, + previous_task_content_digest, mutation_id, revision, content_digest, + task_revision, task_content_digest, item_digest + ) + SELECT application_id, ordinal, project_id, trigger_id, task_id, + previous_revision, previous_content_digest, previous_task_revision, + previous_task_content_digest, mutation_id, revision, content_digest, + task_revision, task_content_digest, item_digest + FROM jsonb_to_recordset($1::jsonb) AS data( + application_id varchar(128), ordinal integer, project_id varchar(128), + trigger_id varchar(128), task_id varchar(128), + previous_revision integer, previous_content_digest char(64), + previous_task_revision integer, previous_task_content_digest char(64), + mutation_id uuid, revision integer, content_digest char(64), + task_revision integer, task_content_digest char(64), item_digest char(64) + )`, + [JSON.stringify(ledger)], + ); + requireRowCount(ledgerInsert.rowCount, applied.length); +} + +async function assertReplayCurrent( + client: PostgresClient, + receipt: ClusterLegacyEnvMigrationApplicationReceipt, +): Promise { + const tasks = await client.query( + `SELECT count(*) AS "totalCount", + count(*) FILTER (WHERE + head.current_revision = ledger.revision AND + revision.mutation_id = ledger.mutation_id AND + revision.content_digest = ledger.content_digest AND + ((ledger.execution_content_digest IS NULL AND execution.project_id IS NULL) OR + (ledger.execution_content_digest IS NOT NULL AND + execution.content_digest = ledger.execution_content_digest)) + ) AS "validCount", + min(ledger.ordinal) AS "minimumOrdinal", + max(ledger.ordinal) AS "maximumOrdinal" + FROM "ql3"."cluster_legacy_env_migration_application_tasks" AS ledger + LEFT JOIN "ql3"."task_definitions" AS head + ON head.project_id = ledger.project_id AND head.task_id = ledger.task_id + LEFT JOIN "ql3"."task_definition_revisions" AS revision + ON revision.project_id = ledger.project_id + AND revision.task_id = ledger.task_id + AND revision.revision = ledger.revision + LEFT JOIN "ql3"."task_execution_revisions" AS execution + ON execution.project_id = ledger.project_id + AND execution.task_id = ledger.task_id + AND execution.source_revision = ledger.revision + AND execution.executor_type = 'remote_worker' + WHERE ledger.application_id = $1`, + [receipt.applicationId], + ); + if (tasks.rows.length !== 1) throw unavailable(); + const taskRow = tasks.rows[0]!; + if ( + postgresRequiredInteger(taskRow.totalCount, unavailable) !== + receipt.taskCount || + postgresRequiredInteger(taskRow.validCount, unavailable) !== + receipt.taskCount || + postgresRequiredInteger(taskRow.minimumOrdinal, unavailable) !== 0 || + postgresRequiredInteger(taskRow.maximumOrdinal, unavailable) !== + receipt.taskCount - 1 + ) { + throw conflict(); + } + + const triggers = await client.query( + `SELECT count(*) AS "totalCount", + count(*) FILTER (WHERE + head.current_revision = ledger.revision AND + revision.mutation_id = ledger.mutation_id AND + revision.content_digest = ledger.content_digest AND + revision.task_revision = ledger.task_revision AND + revision.task_content_digest = ledger.task_content_digest AND + schedule.trigger_revision = ledger.revision + ) AS "validCount", + min(ledger.ordinal) AS "minimumOrdinal", + max(ledger.ordinal) AS "maximumOrdinal" + FROM "ql3"."cluster_legacy_env_migration_application_triggers" AS ledger + LEFT JOIN "ql3"."triggers" AS head + ON head.project_id = ledger.project_id + AND head.trigger_id = ledger.trigger_id + LEFT JOIN "ql3"."trigger_revisions" AS revision + ON revision.project_id = ledger.project_id + AND revision.trigger_id = ledger.trigger_id + AND revision.revision = ledger.revision + LEFT JOIN "ql3"."trigger_schedules" AS schedule + ON schedule.project_id = ledger.project_id + AND schedule.trigger_id = ledger.trigger_id + WHERE ledger.application_id = $1`, + [receipt.applicationId], + ); + if (triggers.rows.length !== 1) throw unavailable(); + const triggerRow = triggers.rows[0]!; + const triggerTotal = postgresRequiredInteger( + triggerRow.totalCount, + unavailable, + ); + const triggerValid = postgresRequiredInteger( + triggerRow.validCount, + unavailable, + ); + if ( + triggerTotal !== receipt.triggerCount || + triggerValid !== receipt.triggerCount || + (receipt.triggerCount === 0 + ? triggerRow.minimumOrdinal !== null || triggerRow.maximumOrdinal !== null + : postgresRequiredInteger(triggerRow.minimumOrdinal, unavailable) !== 0 || + postgresRequiredInteger(triggerRow.maximumOrdinal, unavailable) !== + receipt.triggerCount - 1) + ) { + throw conflict(); + } +} + +function mappedError(error: unknown): Error { + if ( + error instanceof InvalidClusterLegacyEnvMigrationApplicationError || + error instanceof ClusterLegacyEnvMigrationApplicationConflictError || + error instanceof ClusterLegacyEnvMigrationApplicationUnavailableError + ) { + return error; + } + const state = postgresSqlState(error); + if (state === '23503' || state === '23505' || state === '23514') { + return conflict(); + } + return unavailable(); +} + +/** + * Automation-manager-only, bounded-memory application authority. Task and + * Trigger heads, schedules, execution revisions and immutable receipts commit + * in one Project-serialized PostgreSQL transaction. + */ +export class PostgresClusterLegacyEnvMigrationApplicationRepository + implements ClusterLegacyEnvMigrationApplicationRepository +{ + constructor(private readonly pool: PostgresPool) { + if ( + !pool || + typeof pool.query !== 'function' || + typeof pool.connect !== 'function' + ) { + throw new TypeError( + 'PostgreSQL Cluster Legacy Env migration application pool is invalid', + ); + } + } + + async findByApplicationId( + applicationIdValue: string, + ): Promise | null> { + const applicationId = assertClusterLegacyEnvMigrationApplicationIdentifier( + applicationIdValue, + 'applicationId', + ); + try { + return await findReceipt(this.pool, 'application_id', applicationId); + } catch (error) { + throw mappedError(error); + } + } + + async apply( + intentValue: Readonly, + streamsValue: Readonly, + transactionHook?: PostgresClusterLegacyEnvMigrationApplicationTransactionHook, + ): Promise< + Readonly<{ + status: 'applied' | 'existing'; + receipt: Readonly; + }> + > { + if ( + transactionHook !== undefined && + typeof transactionHook !== 'function' + ) { + throw invalid('transaction hook is invalid'); + } + const intent = + normalizeClusterLegacyEnvMigrationApplicationIntent(intentValue); + const streams = normalizeStreams(streamsValue); + + for ( + let attempt = 0; + attempt < POSTGRES_DEFINITION_TRANSACTION_ATTEMPTS; + attempt += 1 + ) { + let client: PostgresClient; + try { + client = await this.pool.connect(); + } catch { + throw unavailable(); + } + let began = false; + let transactionHookError: unknown; + try { + await configurePostgresDefinitionTransaction(client); + began = true; + await client.query( + `SELECT pg_advisory_xact_lock(hashtextextended($1, 0))`, + [ + `qinglong/cluster-legacy-env-migration-application@v1:${intent.projectId}`, + ], + ); + + const replay = await findReceipt( + client, + 'mutation_id', + intent.mutationId, + ); + if (replay) { + if ( + !clusterLegacyEnvMigrationApplicationReceiptMatchesIntent( + replay, + intent, + ) + ) { + throw conflict(); + } + await assertReplayCurrent(client, replay); + if (transactionHook) { + try { + const hookResult = await transactionHook( + client, + Object.freeze({ intent, replay, receipt: replay }), + ); + if (hookResult !== undefined) { + throw invalid('transaction hook must not return a value'); + } + } catch (error) { + transactionHookError = error; + throw error; + } + } + await client.query('COMMIT'); + began = false; + return Object.freeze({ status: 'existing', receipt: replay }); + } + + const occupied = await findReceipt( + client, + 'application_id', + intent.applicationId, + ); + if (occupied) throw conflict(); + const project = await client.query<{ status: unknown }>( + `SELECT status FROM "ql3"."projects" WHERE id = $1 FOR SHARE`, + [intent.projectId], + ); + if (project.rows.length !== 1 || project.rows[0]?.status !== 'active') { + throw conflict(); + } + const plan = await loadPlan(client, intent); + const clock = await client.query<{ committedAtMs: unknown }>( + `SELECT floor(extract(epoch FROM transaction_timestamp()) * 1000)::bigint AS "committedAtMs"`, + ); + if (clock.rows.length !== 1) throw unavailable(); + const committedAtMs = postgresRequiredInteger( + clock.rows[0]?.committedAtMs, + unavailable, + ); + const receipt = createClusterLegacyEnvMigrationApplicationReceipt({ + applicationId: intent.applicationId, + mutationId: intent.mutationId, + projectId: intent.projectId, + planId: intent.planId, + planDigest: intent.planDigest, + environmentBundleRef: plan.target.secretRef, + taskRevisionSetDigest: plan.target.taskRevisionSetDigest, + triggerRevisionSetDigest: plan.target.triggerRevisionSetDigest, + taskMutationSetDigest: intent.taskMutationSetDigest, + triggerMutationSetDigest: intent.triggerMutationSetDigest, + taskCount: plan.target.taskCount, + triggerCount: plan.target.triggerCount, + committedAtMs, + }); + await client.query( + `INSERT INTO "ql3"."cluster_legacy_env_migration_application_receipts" ( + application_id, mutation_id, project_id, plan_id, plan_digest, + environment_bundle_ref, task_revision_set_digest, + trigger_revision_set_digest, task_mutation_set_digest, + trigger_mutation_set_digest, task_count, trigger_count, + committed_at_ms, receipt_digest, receipt_json + ) VALUES ( + $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, + $11, $12, $13, $14, $15::jsonb + )`, + [ + receipt.applicationId, + receipt.mutationId, + receipt.projectId, + receipt.planId, + receipt.planDigest, + receipt.environmentBundleRef, + receipt.taskRevisionSetDigest, + receipt.triggerRevisionSetDigest, + receipt.taskMutationSetDigest, + receipt.triggerMutationSetDigest, + receipt.taskCount, + receipt.triggerCount, + receipt.committedAtMs, + receipt.receiptDigest, + JSON.stringify(receipt), + ], + ); + + const taskDigester = + createClusterLegacyEnvMigrationTaskMutationSetDigester(); + let taskBatch: ClusterLegacyEnvMigrationTaskMutation[] = []; + for await (const raw of mutationStream( + streams.taskMutations, + 'Task mutation stream', + )) { + taskBatch.push(taskDigester.update(raw)); + if (taskBatch.length === APPLICATION_BATCH_SIZE) { + await insertTaskBatch( + client, + receipt.applicationId, + receipt.projectId, + receipt.environmentBundleRef, + receipt.committedAtMs, + taskBatch, + ); + taskBatch = []; + } + } + if (taskBatch.length > 0) { + await insertTaskBatch( + client, + receipt.applicationId, + receipt.projectId, + receipt.environmentBundleRef, + receipt.committedAtMs, + taskBatch, + ); + } + const taskSet = taskDigester.finish(); + if ( + taskSet.count !== receipt.taskCount || + taskSet.revisionSetDigest !== receipt.taskRevisionSetDigest || + taskSet.mutationSetDigest !== receipt.taskMutationSetDigest + ) { + throw conflict(); + } + + const triggerDigester = + createClusterLegacyEnvMigrationTriggerMutationSetDigester(); + let triggerBatch: ClusterLegacyEnvMigrationTriggerMutation[] = []; + for await (const raw of mutationStream( + streams.triggerMutations, + 'Trigger mutation stream', + )) { + triggerBatch.push(triggerDigester.update(raw)); + if (triggerBatch.length === APPLICATION_BATCH_SIZE) { + await insertTriggerBatch( + client, + receipt.applicationId, + receipt.projectId, + receipt.committedAtMs, + triggerBatch, + ); + triggerBatch = []; + } + } + if (triggerBatch.length > 0) { + await insertTriggerBatch( + client, + receipt.applicationId, + receipt.projectId, + receipt.committedAtMs, + triggerBatch, + ); + } + const triggerSet = triggerDigester.finish(); + if ( + triggerSet.count !== receipt.triggerCount || + triggerSet.revisionSetDigest !== receipt.triggerRevisionSetDigest || + triggerSet.mutationSetDigest !== receipt.triggerMutationSetDigest + ) { + throw conflict(); + } + + if (transactionHook) { + try { + const hookResult = await transactionHook( + client, + Object.freeze({ intent, replay: null, receipt }), + ); + if (hookResult !== undefined) { + throw invalid('transaction hook must not return a value'); + } + } catch (error) { + transactionHookError = error; + throw error; + } + } + await client.query('COMMIT'); + began = false; + return Object.freeze({ status: 'applied', receipt }); + } catch (error) { + if (began) await rollbackPostgresDefinitionTransaction(client); + const state = postgresSqlState(error); + if ( + state && + POSTGRES_DEFINITION_RETRYABLE_SQL_STATES.has(state) && + attempt + 1 < POSTGRES_DEFINITION_TRANSACTION_ATTEMPTS + ) { + continue; + } + if (error === transactionHookError && error instanceof Error) { + throw error; + } + throw mappedError(error); + } finally { + client.release(); + } + } + throw unavailable(); + } +} diff --git a/packages/ql3-cluster-postgres/src/reconciliation/pg-0071-cluster-legacy-env-migration-applications.ts b/packages/ql3-cluster-postgres/src/reconciliation/pg-0071-cluster-legacy-env-migration-applications.ts new file mode 100644 index 00000000..70dcb343 --- /dev/null +++ b/packages/ql3-cluster-postgres/src/reconciliation/pg-0071-cluster-legacy-env-migration-applications.ts @@ -0,0 +1,182 @@ +import { definePostgresSqlMigration } from '../migrations/sqlMigration'; +import { CAPABILITIES_V69 } from './pg-0070-cluster-legacy-env-migration-plans'; + +export const CAPABILITIES_V70 = CAPABILITIES_V69.replace( + '"cluster_legacy_env_migration_plan":1,', + '"cluster_legacy_env_migration_application":1,"cluster_legacy_env_migration_plan":1,', +); + +export const pg0071ClusterLegacyEnvMigrationApplicationsMigration = + definePostgresSqlMigration({ + id: 'pg-0071-cluster-legacy-env-migration-applications', + statements: [ + ` +CREATE TABLE "ql3"."cluster_legacy_env_migration_application_receipts" ( + application_id varchar(128) PRIMARY KEY, + mutation_id uuid NOT NULL, + project_id varchar(128) NOT NULL, + plan_id varchar(128) NOT NULL, + plan_digest char(64) NOT NULL, + environment_bundle_ref varchar(512) NOT NULL, + task_revision_set_digest char(64) NOT NULL, + trigger_revision_set_digest char(64) NOT NULL, + task_mutation_set_digest char(64) NOT NULL, + trigger_mutation_set_digest char(64) NOT NULL, + task_count integer NOT NULL, + trigger_count integer NOT NULL, + committed_at_ms bigint NOT NULL, + receipt_digest char(64) NOT NULL, + receipt_json jsonb NOT NULL, + CONSTRAINT ql3_cluster_legacy_env_application_project_fk + FOREIGN KEY (project_id) REFERENCES "ql3"."projects" (id) + ON DELETE RESTRICT ON UPDATE RESTRICT, + CONSTRAINT ql3_cluster_legacy_env_application_plan_fk + FOREIGN KEY (plan_id) REFERENCES "ql3"."cluster_legacy_env_migration_plans" (plan_id) + ON DELETE RESTRICT ON UPDATE RESTRICT, + CONSTRAINT ql3_cluster_legacy_env_application_identity_check CHECK ( + application_id ~ '^[A-Za-z0-9][A-Za-z0-9._:-]{0,127}$' AND + project_id ~ '^[A-Za-z0-9][A-Za-z0-9._:-]{0,127}$' AND + plan_id ~ '^[A-Za-z0-9][A-Za-z0-9._:-]{0,127}$' + ), + CONSTRAINT ql3_cluster_legacy_env_application_digest_check CHECK ( + plan_digest ~ '^[0-9a-f]{64}$' AND + task_revision_set_digest ~ '^[0-9a-f]{64}$' AND + trigger_revision_set_digest ~ '^[0-9a-f]{64}$' AND + task_mutation_set_digest ~ '^[0-9a-f]{64}$' AND + trigger_mutation_set_digest ~ '^[0-9a-f]{64}$' AND + receipt_digest ~ '^[0-9a-f]{64}$' + ), + CONSTRAINT ql3_cluster_legacy_env_application_target_check CHECK ( + environment_bundle_ref ~ '^qlsecret:v1:[A-Za-z0-9_-]+$' AND + octet_length(environment_bundle_ref) BETWEEN 14 AND 512 AND + task_count BETWEEN 1 AND 100000 AND + trigger_count BETWEEN 0 AND 500000 AND + committed_at_ms >= 0 + ), + CONSTRAINT ql3_cluster_legacy_env_application_json_check CHECK ( + jsonb_typeof(receipt_json) = 'object' AND + octet_length(receipt_json::text) BETWEEN 2 AND 8192 AND + receipt_json = jsonb_build_object( + 'schema', 'qinglong/cluster-legacy-env-migration-application-receipt@v1', + 'applicationId', application_id, + 'mutationId', mutation_id::text, + 'projectId', project_id, + 'planId', plan_id, + 'planDigest', plan_digest, + 'environmentBundleRef', environment_bundle_ref, + 'taskRevisionSetDigest', task_revision_set_digest, + 'triggerRevisionSetDigest', trigger_revision_set_digest, + 'taskMutationSetDigest', task_mutation_set_digest, + 'triggerMutationSetDigest', trigger_mutation_set_digest, + 'taskCount', task_count, + 'triggerCount', trigger_count, + 'committedAtMs', committed_at_ms, + 'receiptDigest', receipt_digest + ) + ) +) + `.trim(), + `CREATE UNIQUE INDEX ql3_cluster_legacy_env_application_mutation_uidx ON "ql3"."cluster_legacy_env_migration_application_receipts" (mutation_id)`, + `CREATE UNIQUE INDEX ql3_cluster_legacy_env_application_project_uidx ON "ql3"."cluster_legacy_env_migration_application_receipts" (application_id, project_id)`, + `CREATE UNIQUE INDEX ql3_cluster_legacy_env_application_plan_uidx ON "ql3"."cluster_legacy_env_migration_application_receipts" (plan_id)`, + `CREATE UNIQUE INDEX ql3_cluster_legacy_env_application_digest_uidx ON "ql3"."cluster_legacy_env_migration_application_receipts" (receipt_digest)`, + `CREATE INDEX ql3_cluster_legacy_env_application_project_idx ON "ql3"."cluster_legacy_env_migration_application_receipts" (project_id, committed_at_ms, application_id)`, + ` +CREATE TABLE "ql3"."cluster_legacy_env_migration_application_tasks" ( + application_id varchar(128) NOT NULL, + ordinal integer NOT NULL, + project_id varchar(128) NOT NULL, + task_id varchar(128) NOT NULL, + previous_revision integer NOT NULL, + previous_content_digest char(64) NOT NULL, + mutation_id uuid NOT NULL, + revision integer NOT NULL, + content_digest char(64) NOT NULL, + execution_content_digest char(64), + item_digest char(64) NOT NULL, + CONSTRAINT cluster_legacy_env_migration_application_tasks_pkey + PRIMARY KEY (application_id, ordinal), + CONSTRAINT ql3_cluster_legacy_env_application_task_receipt_fk + FOREIGN KEY (application_id, project_id) + REFERENCES "ql3"."cluster_legacy_env_migration_application_receipts" + (application_id, project_id) + ON DELETE RESTRICT ON UPDATE RESTRICT, + CONSTRAINT ql3_cluster_legacy_env_application_task_identity_check CHECK ( + project_id ~ '^[A-Za-z0-9][A-Za-z0-9._:-]{0,127}$' AND + char_length(task_id) BETWEEN 1 AND 128 AND task_id !~ '[[:cntrl:]]' + ), + CONSTRAINT ql3_cluster_legacy_env_application_task_revision_check CHECK ( + ordinal BETWEEN 0 AND 99999 AND + previous_revision BETWEEN 1 AND 2147483646 AND + revision = previous_revision + 1 + ), + CONSTRAINT ql3_cluster_legacy_env_application_task_digest_check CHECK ( + previous_content_digest ~ '^[0-9a-f]{64}$' AND + content_digest ~ '^[0-9a-f]{64}$' AND + (execution_content_digest IS NULL OR execution_content_digest ~ '^[0-9a-f]{64}$') AND + item_digest ~ '^[0-9a-f]{64}$' + ) +) + `.trim(), + `CREATE UNIQUE INDEX ql3_cluster_legacy_env_application_task_uidx ON "ql3"."cluster_legacy_env_migration_application_tasks" (application_id, task_id)`, + `CREATE UNIQUE INDEX ql3_cluster_legacy_env_application_task_revision_uidx ON "ql3"."cluster_legacy_env_migration_application_tasks" (application_id, project_id, task_id, revision, content_digest)`, + ` +CREATE TABLE "ql3"."cluster_legacy_env_migration_application_triggers" ( + application_id varchar(128) NOT NULL, + ordinal integer NOT NULL, + project_id varchar(128) NOT NULL, + trigger_id varchar(128) NOT NULL, + task_id varchar(128) NOT NULL, + previous_revision integer NOT NULL, + previous_content_digest char(64) NOT NULL, + previous_task_revision integer NOT NULL, + previous_task_content_digest char(64) NOT NULL, + mutation_id uuid NOT NULL, + revision integer NOT NULL, + content_digest char(64) NOT NULL, + task_revision integer NOT NULL, + task_content_digest char(64) NOT NULL, + item_digest char(64) NOT NULL, + CONSTRAINT cluster_legacy_env_migration_application_triggers_pkey + PRIMARY KEY (application_id, ordinal), + CONSTRAINT ql3_cluster_legacy_env_application_trigger_receipt_fk + FOREIGN KEY (application_id, project_id) + REFERENCES "ql3"."cluster_legacy_env_migration_application_receipts" + (application_id, project_id) + ON DELETE RESTRICT ON UPDATE RESTRICT, + CONSTRAINT ql3_cluster_legacy_env_application_trigger_task_fk + FOREIGN KEY (application_id, project_id, task_id, task_revision, task_content_digest) + REFERENCES "ql3"."cluster_legacy_env_migration_application_tasks" + (application_id, project_id, task_id, revision, content_digest) + ON DELETE RESTRICT ON UPDATE RESTRICT, + CONSTRAINT ql3_cluster_legacy_env_application_trigger_identity_check CHECK ( + project_id ~ '^[A-Za-z0-9][A-Za-z0-9._:-]{0,127}$' AND + char_length(trigger_id) BETWEEN 1 AND 128 AND trigger_id !~ '[[:cntrl:]]' AND + char_length(task_id) BETWEEN 1 AND 128 AND task_id !~ '[[:cntrl:]]' + ), + CONSTRAINT ql3_cluster_legacy_env_application_trigger_revision_check CHECK ( + ordinal BETWEEN 0 AND 499999 AND + previous_revision BETWEEN 1 AND 2147483646 AND + revision = previous_revision + 1 AND + previous_task_revision BETWEEN 1 AND 2147483647 AND + task_revision BETWEEN 2 AND 2147483647 + ), + CONSTRAINT ql3_cluster_legacy_env_application_trigger_digest_check CHECK ( + previous_content_digest ~ '^[0-9a-f]{64}$' AND + previous_task_content_digest ~ '^[0-9a-f]{64}$' AND + content_digest ~ '^[0-9a-f]{64}$' AND + task_content_digest ~ '^[0-9a-f]{64}$' AND + item_digest ~ '^[0-9a-f]{64}$' + ) +) + `.trim(), + `CREATE UNIQUE INDEX ql3_cluster_legacy_env_application_trigger_uidx ON "ql3"."cluster_legacy_env_migration_application_triggers" (application_id, trigger_id)`, + `REVOKE ALL ON "ql3"."cluster_legacy_env_migration_application_receipts" FROM PUBLIC`, + `REVOKE ALL ON "ql3"."cluster_legacy_env_migration_application_tasks" FROM PUBLIC`, + `REVOKE ALL ON "ql3"."cluster_legacy_env_migration_application_triggers" FROM PUBLIC`, + `GRANT SELECT, INSERT ON "ql3"."cluster_legacy_env_migration_application_receipts" TO ql3_automation_manager`, + `GRANT SELECT, INSERT ON "ql3"."cluster_legacy_env_migration_application_tasks" TO ql3_automation_manager`, + `GRANT SELECT, INSERT ON "ql3"."cluster_legacy_env_migration_application_triggers" TO ql3_automation_manager`, + `DO $ql3$ BEGIN UPDATE "ql3"."schema_capabilities" SET contract_version = 70, migration_id = 'pg-0071-cluster-legacy-env-migration-applications', capabilities = '${CAPABILITIES_V70}'::jsonb, updated_at_ms = floor(extract(epoch FROM transaction_timestamp()) * 1000)::bigint WHERE contract_name = 'control-core' AND contract_version = 69 AND migration_id = 'pg-0070-cluster-legacy-env-migration-plans' AND capabilities = '${CAPABILITIES_V69}'::jsonb; IF NOT FOUND THEN RAISE EXCEPTION 'control-core capability is not at version 69' USING ERRCODE = 'check_violation'; END IF; END $ql3$`, + ], + }); diff --git a/packages/ql3-cluster-postgres/src/schema/schema.ts b/packages/ql3-cluster-postgres/src/schema/schema.ts index 900747a6..2f3556a1 100644 --- a/packages/ql3-cluster-postgres/src/schema/schema.ts +++ b/packages/ql3-cluster-postgres/src/schema/schema.ts @@ -180,6 +180,213 @@ export const clusterLegacyEnvMigrationPlans = ql3Schema.table( ], ); +export const clusterLegacyEnvMigrationApplicationReceipts = ql3Schema.table( + 'cluster_legacy_env_migration_application_receipts', + { + applicationId: varchar('application_id', { length: 128 }).primaryKey(), + mutationId: uuid('mutation_id').notNull(), + projectId: varchar('project_id', { length: 128 }).notNull(), + planId: varchar('plan_id', { length: 128 }).notNull(), + planDigest: char('plan_digest', { length: 64 }).notNull(), + environmentBundleRef: varchar('environment_bundle_ref', { + length: 512, + }).notNull(), + taskRevisionSetDigest: char('task_revision_set_digest', { + length: 64, + }).notNull(), + triggerRevisionSetDigest: char('trigger_revision_set_digest', { + length: 64, + }).notNull(), + taskMutationSetDigest: char('task_mutation_set_digest', { + length: 64, + }).notNull(), + triggerMutationSetDigest: char('trigger_mutation_set_digest', { + length: 64, + }).notNull(), + taskCount: integer('task_count').notNull(), + triggerCount: integer('trigger_count').notNull(), + committedAtMs: bigint('committed_at_ms', { mode: 'number' }).notNull(), + receiptDigest: char('receipt_digest', { length: 64 }).notNull(), + receiptJson: jsonb('receipt_json') + .$type>() + .notNull(), + }, + (table) => [ + foreignKey({ + name: 'ql3_cluster_legacy_env_application_project_fk', + columns: [table.projectId], + foreignColumns: [projects.id], + }).onDelete('restrict'), + foreignKey({ + name: 'ql3_cluster_legacy_env_application_plan_fk', + columns: [table.planId], + foreignColumns: [clusterLegacyEnvMigrationPlans.planId], + }).onDelete('restrict'), + check( + 'ql3_cluster_legacy_env_application_identity_check', + sql`${table.applicationId} ~ '^[A-Za-z0-9][A-Za-z0-9._:-]{0,127}$' and ${table.projectId} ~ '^[A-Za-z0-9][A-Za-z0-9._:-]{0,127}$' and ${table.planId} ~ '^[A-Za-z0-9][A-Za-z0-9._:-]{0,127}$'`, + ), + check( + 'ql3_cluster_legacy_env_application_digest_check', + sql`${table.planDigest} ~ '^[0-9a-f]{64}$' and ${table.taskRevisionSetDigest} ~ '^[0-9a-f]{64}$' and ${table.triggerRevisionSetDigest} ~ '^[0-9a-f]{64}$' and ${table.taskMutationSetDigest} ~ '^[0-9a-f]{64}$' and ${table.triggerMutationSetDigest} ~ '^[0-9a-f]{64}$' and ${table.receiptDigest} ~ '^[0-9a-f]{64}$'`, + ), + check( + 'ql3_cluster_legacy_env_application_target_check', + sql`${table.environmentBundleRef} ~ '^qlsecret:v1:[A-Za-z0-9_-]+$' and octet_length(${table.environmentBundleRef}) between 14 and 512 and ${table.taskCount} between 1 and 100000 and ${table.triggerCount} between 0 and 500000 and ${table.committedAtMs} >= 0`, + ), + check( + 'ql3_cluster_legacy_env_application_json_check', + sql`jsonb_typeof(${table.receiptJson}) = 'object' and octet_length(${table.receiptJson}::text) between 2 and 8192 and ${table.receiptJson} = jsonb_build_object('schema', 'qinglong/cluster-legacy-env-migration-application-receipt@v1', 'applicationId', ${table.applicationId}, 'mutationId', ${table.mutationId}::text, 'projectId', ${table.projectId}, 'planId', ${table.planId}, 'planDigest', ${table.planDigest}, 'environmentBundleRef', ${table.environmentBundleRef}, 'taskRevisionSetDigest', ${table.taskRevisionSetDigest}, 'triggerRevisionSetDigest', ${table.triggerRevisionSetDigest}, 'taskMutationSetDigest', ${table.taskMutationSetDigest}, 'triggerMutationSetDigest', ${table.triggerMutationSetDigest}, 'taskCount', ${table.taskCount}, 'triggerCount', ${table.triggerCount}, 'committedAtMs', ${table.committedAtMs}, 'receiptDigest', ${table.receiptDigest})`, + ), + uniqueIndex('ql3_cluster_legacy_env_application_mutation_uidx').on( + table.mutationId, + ), + uniqueIndex('ql3_cluster_legacy_env_application_project_uidx').on( + table.applicationId, + table.projectId, + ), + uniqueIndex('ql3_cluster_legacy_env_application_plan_uidx').on( + table.planId, + ), + uniqueIndex('ql3_cluster_legacy_env_application_digest_uidx').on( + table.receiptDigest, + ), + index('ql3_cluster_legacy_env_application_project_idx').on( + table.projectId, + table.committedAtMs, + table.applicationId, + ), + ], +); + +export const clusterLegacyEnvMigrationApplicationTasks = ql3Schema.table( + 'cluster_legacy_env_migration_application_tasks', + { + applicationId: varchar('application_id', { length: 128 }).notNull(), + ordinal: integer('ordinal').notNull(), + projectId: varchar('project_id', { length: 128 }).notNull(), + taskId: varchar('task_id', { length: 128 }).notNull(), + previousRevision: integer('previous_revision').notNull(), + previousContentDigest: char('previous_content_digest', { + length: 64, + }).notNull(), + mutationId: uuid('mutation_id').notNull(), + revision: integer('revision').notNull(), + contentDigest: char('content_digest', { length: 64 }).notNull(), + executionContentDigest: char('execution_content_digest', { length: 64 }), + itemDigest: char('item_digest', { length: 64 }).notNull(), + }, + (table) => [ + primaryKey({ + name: 'cluster_legacy_env_migration_application_tasks_pkey', + columns: [table.applicationId, table.ordinal], + }), + foreignKey({ + name: 'ql3_cluster_legacy_env_application_task_receipt_fk', + columns: [table.applicationId, table.projectId], + foreignColumns: [ + clusterLegacyEnvMigrationApplicationReceipts.applicationId, + clusterLegacyEnvMigrationApplicationReceipts.projectId, + ], + }).onDelete('restrict'), + check( + 'ql3_cluster_legacy_env_application_task_identity_check', + sql`${table.projectId} ~ '^[A-Za-z0-9][A-Za-z0-9._:-]{0,127}$' and char_length(${table.taskId}) between 1 and 128 and ${table.taskId} !~ '[[:cntrl:]]'`, + ), + check( + 'ql3_cluster_legacy_env_application_task_revision_check', + sql`${table.ordinal} between 0 and 99999 and ${table.previousRevision} between 1 and 2147483646 and ${table.revision} = ${table.previousRevision} + 1`, + ), + check( + 'ql3_cluster_legacy_env_application_task_digest_check', + sql`${table.previousContentDigest} ~ '^[0-9a-f]{64}$' and ${table.contentDigest} ~ '^[0-9a-f]{64}$' and (${table.executionContentDigest} is null or ${table.executionContentDigest} ~ '^[0-9a-f]{64}$') and ${table.itemDigest} ~ '^[0-9a-f]{64}$'`, + ), + uniqueIndex('ql3_cluster_legacy_env_application_task_uidx').on( + table.applicationId, + table.taskId, + ), + uniqueIndex('ql3_cluster_legacy_env_application_task_revision_uidx').on( + table.applicationId, + table.projectId, + table.taskId, + table.revision, + table.contentDigest, + ), + ], +); + +export const clusterLegacyEnvMigrationApplicationTriggers = ql3Schema.table( + 'cluster_legacy_env_migration_application_triggers', + { + applicationId: varchar('application_id', { length: 128 }).notNull(), + ordinal: integer('ordinal').notNull(), + projectId: varchar('project_id', { length: 128 }).notNull(), + triggerId: varchar('trigger_id', { length: 128 }).notNull(), + taskId: varchar('task_id', { length: 128 }).notNull(), + previousRevision: integer('previous_revision').notNull(), + previousContentDigest: char('previous_content_digest', { + length: 64, + }).notNull(), + previousTaskRevision: integer('previous_task_revision').notNull(), + previousTaskContentDigest: char('previous_task_content_digest', { + length: 64, + }).notNull(), + mutationId: uuid('mutation_id').notNull(), + revision: integer('revision').notNull(), + contentDigest: char('content_digest', { length: 64 }).notNull(), + taskRevision: integer('task_revision').notNull(), + taskContentDigest: char('task_content_digest', { length: 64 }).notNull(), + itemDigest: char('item_digest', { length: 64 }).notNull(), + }, + (table) => [ + primaryKey({ + name: 'cluster_legacy_env_migration_application_triggers_pkey', + columns: [table.applicationId, table.ordinal], + }), + foreignKey({ + name: 'ql3_cluster_legacy_env_application_trigger_receipt_fk', + columns: [table.applicationId, table.projectId], + foreignColumns: [ + clusterLegacyEnvMigrationApplicationReceipts.applicationId, + clusterLegacyEnvMigrationApplicationReceipts.projectId, + ], + }).onDelete('restrict'), + foreignKey({ + name: 'ql3_cluster_legacy_env_application_trigger_task_fk', + columns: [ + table.applicationId, + table.projectId, + table.taskId, + table.taskRevision, + table.taskContentDigest, + ], + foreignColumns: [ + clusterLegacyEnvMigrationApplicationTasks.applicationId, + clusterLegacyEnvMigrationApplicationTasks.projectId, + clusterLegacyEnvMigrationApplicationTasks.taskId, + clusterLegacyEnvMigrationApplicationTasks.revision, + clusterLegacyEnvMigrationApplicationTasks.contentDigest, + ], + }).onDelete('restrict'), + check( + 'ql3_cluster_legacy_env_application_trigger_identity_check', + sql`${table.projectId} ~ '^[A-Za-z0-9][A-Za-z0-9._:-]{0,127}$' and char_length(${table.triggerId}) between 1 and 128 and ${table.triggerId} !~ '[[:cntrl:]]' and char_length(${table.taskId}) between 1 and 128 and ${table.taskId} !~ '[[:cntrl:]]'`, + ), + check( + 'ql3_cluster_legacy_env_application_trigger_revision_check', + sql`${table.ordinal} between 0 and 499999 and ${table.previousRevision} between 1 and 2147483646 and ${table.revision} = ${table.previousRevision} + 1 and ${table.previousTaskRevision} between 1 and 2147483647 and ${table.taskRevision} between 2 and 2147483647`, + ), + check( + 'ql3_cluster_legacy_env_application_trigger_digest_check', + sql`${table.previousContentDigest} ~ '^[0-9a-f]{64}$' and ${table.previousTaskContentDigest} ~ '^[0-9a-f]{64}$' and ${table.contentDigest} ~ '^[0-9a-f]{64}$' and ${table.taskContentDigest} ~ '^[0-9a-f]{64}$' and ${table.itemDigest} ~ '^[0-9a-f]{64}$'`, + ), + uniqueIndex('ql3_cluster_legacy_env_application_trigger_uidx').on( + table.applicationId, + table.triggerId, + ), + ], +); + export const pluginPackageInstalls = ql3Schema.table( 'plugin_package_installs', { @@ -6289,6 +6496,9 @@ export const ql3PostgresTables = [ schemaCapabilities, projects, clusterLegacyEnvMigrationPlans, + clusterLegacyEnvMigrationApplicationReceipts, + clusterLegacyEnvMigrationApplicationTasks, + clusterLegacyEnvMigrationApplicationTriggers, pluginPackageInstalls, pluginPackageInstallHeads, pluginPackageInstallMutations, diff --git a/packages/ql3-cluster-postgres/src/schema/schemaContract.ts b/packages/ql3-cluster-postgres/src/schema/schemaContract.ts index 9ef3165c..e1ea967f 100644 --- a/packages/ql3-cluster-postgres/src/schema/schemaContract.ts +++ b/packages/ql3-cluster-postgres/src/schema/schemaContract.ts @@ -21,8 +21,8 @@ export interface PostgresSchemaContractTrigger { export interface PostgresSchemaContract { readonly schema: 'ql3'; readonly contractName: 'control-core'; - readonly contractVersion: 69; - readonly migrationId: 'pg-0070-cluster-legacy-env-migration-plans'; + readonly contractVersion: 70; + readonly migrationId: 'pg-0071-cluster-legacy-env-migration-applications'; readonly minimumServerMajor: 16; readonly maximumServerMajor: 18; readonly capabilities: Readonly<{ @@ -50,6 +50,7 @@ export interface PostgresSchemaContract { cluster_scheduler_admission: 1; database_role_grants: 1; cluster_execution_revision: 1; + cluster_legacy_env_migration_application: 1; cluster_legacy_env_migration_plan: 1; identity_admin: 1; plugin_package_admission: 1; @@ -123,8 +124,8 @@ export const postgresqlControlSchemaContract: PostgresSchemaContract = Object.freeze({ schema: 'ql3', contractName: 'control-core', - contractVersion: 69, - migrationId: 'pg-0070-cluster-legacy-env-migration-plans', + contractVersion: 70, + migrationId: 'pg-0071-cluster-legacy-env-migration-applications', minimumServerMajor: 16, maximumServerMajor: 18, capabilities: Object.freeze({ @@ -138,6 +139,7 @@ export const postgresqlControlSchemaContract: PostgresSchemaContract = automation_management_boundary: 1, automation_management_identity_keyset_ledger: 1, cluster_execution_revision: 1, + cluster_legacy_env_migration_application: 1, cluster_legacy_env_migration_plan: 1, cluster_recovery: 1, cluster_recovery_claim: 1, @@ -251,6 +253,53 @@ export const postgresqlControlSchemaContract: PostgresSchemaContract = 'planned_at_ms', 'plan_json', ]), + table('cluster_legacy_env_migration_application_receipts', [ + 'application_id', + 'mutation_id', + 'project_id', + 'plan_id', + 'plan_digest', + 'environment_bundle_ref', + 'task_revision_set_digest', + 'trigger_revision_set_digest', + 'task_mutation_set_digest', + 'trigger_mutation_set_digest', + 'task_count', + 'trigger_count', + 'committed_at_ms', + 'receipt_digest', + 'receipt_json', + ]), + table('cluster_legacy_env_migration_application_tasks', [ + 'application_id', + 'ordinal', + 'project_id', + 'task_id', + 'previous_revision', + 'previous_content_digest', + 'mutation_id', + 'revision', + 'content_digest', + 'execution_content_digest', + 'item_digest', + ]), + table('cluster_legacy_env_migration_application_triggers', [ + 'application_id', + 'ordinal', + 'project_id', + 'trigger_id', + 'task_id', + 'previous_revision', + 'previous_content_digest', + 'previous_task_revision', + 'previous_task_content_digest', + 'mutation_id', + 'revision', + 'content_digest', + 'task_revision', + 'task_content_digest', + 'item_digest', + ]), table('plugin_package_installs', [ 'installation_id', 'project_id', @@ -1566,6 +1615,17 @@ export const postgresqlControlSchemaContract: PostgresSchemaContract = 'ql3_cluster_legacy_env_plan_mutation_uidx', 'ql3_cluster_legacy_env_plan_digest_uidx', 'ql3_cluster_legacy_env_plan_project_idx', + 'cluster_legacy_env_migration_application_receipts_pkey', + 'ql3_cluster_legacy_env_application_mutation_uidx', + 'ql3_cluster_legacy_env_application_project_uidx', + 'ql3_cluster_legacy_env_application_plan_uidx', + 'ql3_cluster_legacy_env_application_digest_uidx', + 'ql3_cluster_legacy_env_application_project_idx', + 'cluster_legacy_env_migration_application_tasks_pkey', + 'ql3_cluster_legacy_env_application_task_uidx', + 'ql3_cluster_legacy_env_application_task_revision_uidx', + 'cluster_legacy_env_migration_application_triggers_pkey', + 'ql3_cluster_legacy_env_application_trigger_uidx', 'plugin_package_installs_pkey', 'ql3_plugin_package_installs_quarantine_target_key', 'ql3_plugin_package_installs_recovery_idx', @@ -1896,6 +1956,16 @@ export const postgresqlControlSchemaContract: PostgresSchemaContract = 'ql3_cluster_legacy_env_plan_target_check', 'ql3_cluster_legacy_env_plan_time_check', 'ql3_cluster_legacy_env_plan_json_check', + 'ql3_cluster_legacy_env_application_identity_check', + 'ql3_cluster_legacy_env_application_digest_check', + 'ql3_cluster_legacy_env_application_target_check', + 'ql3_cluster_legacy_env_application_json_check', + 'ql3_cluster_legacy_env_application_task_identity_check', + 'ql3_cluster_legacy_env_application_task_revision_check', + 'ql3_cluster_legacy_env_application_task_digest_check', + 'ql3_cluster_legacy_env_application_trigger_identity_check', + 'ql3_cluster_legacy_env_application_trigger_revision_check', + 'ql3_cluster_legacy_env_application_trigger_digest_check', 'ql3_plugin_package_installs_identity_check', 'ql3_plugin_package_installs_operation_check', 'ql3_plugin_package_installs_state_check', @@ -2376,6 +2446,11 @@ export const postgresqlControlSchemaContract: PostgresSchemaContract = foreignKeys: Object.freeze([ 'ql3_schema_capabilities_migration_fk', 'ql3_cluster_legacy_env_plan_project_fk', + 'ql3_cluster_legacy_env_application_project_fk', + 'ql3_cluster_legacy_env_application_plan_fk', + 'ql3_cluster_legacy_env_application_task_receipt_fk', + 'ql3_cluster_legacy_env_application_trigger_receipt_fk', + 'ql3_cluster_legacy_env_application_trigger_task_fk', 'ql3_plugin_package_installs_project_fk', 'ql3_plugin_package_install_heads_project_fk', 'ql3_plugin_package_install_heads_install_fk', diff --git a/packages/ql3-cluster-postgres/src/schema/schemaReadiness.ts b/packages/ql3-cluster-postgres/src/schema/schemaReadiness.ts index 42221528..d6dc12aa 100644 --- a/packages/ql3-cluster-postgres/src/schema/schemaReadiness.ts +++ b/packages/ql3-cluster-postgres/src/schema/schemaReadiness.ts @@ -156,6 +156,24 @@ const REQUIRED_RUNTIME_PRIVILEGES = Object.freeze({ update: false, delete: false, }), + cluster_legacy_env_migration_application_receipts: Object.freeze({ + select: false, + insert: false, + update: false, + delete: false, + }), + cluster_legacy_env_migration_application_tasks: Object.freeze({ + select: false, + insert: false, + update: false, + delete: false, + }), + cluster_legacy_env_migration_application_triggers: Object.freeze({ + select: false, + insert: false, + update: false, + delete: false, + }), plugin_package_installs: Object.freeze({ select: false, insert: false, @@ -711,6 +729,24 @@ const REQUIRED_ADMIN_PRIVILEGES = Object.freeze({ update: false, delete: false, }), + cluster_legacy_env_migration_application_receipts: Object.freeze({ + select: false, + insert: false, + update: false, + delete: false, + }), + cluster_legacy_env_migration_application_tasks: Object.freeze({ + select: false, + insert: false, + update: false, + delete: false, + }), + cluster_legacy_env_migration_application_triggers: Object.freeze({ + select: false, + insert: false, + update: false, + delete: false, + }), plugin_package_installs: Object.freeze({ select: false, insert: false, @@ -1430,6 +1466,9 @@ const REQUIRED_AUTOMATION_MANAGER_PRIVILEGES: RequiredPrivileges = } : name === 'security_audit_events' || name === 'cluster_legacy_env_migration_plans' || + name === 'cluster_legacy_env_migration_application_receipts' || + name === 'cluster_legacy_env_migration_application_tasks' || + name === 'cluster_legacy_env_migration_application_triggers' || name === 'task_definition_revisions' || name === 'task_execution_revisions' || name === 'trigger_revisions' diff --git a/packages/ql3-cluster-postgres/test/postgres.integration.test.cjs b/packages/ql3-cluster-postgres/test/postgres.integration.test.cjs index 6a53fe3d..42c1093b 100644 --- a/packages/ql3-cluster-postgres/test/postgres.integration.test.cjs +++ b/packages/ql3-cluster-postgres/test/postgres.integration.test.cjs @@ -74,6 +74,14 @@ const { const { PostgresClusterLegacyEnvMigrationPlanRepository, } = require('../dist/reconciliation/clusterLegacyEnvMigrationPlanRepository'); +const { + PostgresClusterLegacyEnvMigrationApplicationRepository, +} = require('../dist/reconciliation/clusterLegacyEnvMigrationApplicationRepository'); +const { + InvalidClusterLegacyEnvMigrationApplicationError, + createClusterLegacyEnvMigrationTaskMutationSetDigester, + createClusterLegacyEnvMigrationTriggerMutationSetDigester, +} = require('@qinglong/runtime-core/cluster-legacy-env-migration-application'); function nextMinute(schedule, afterMs) { if (schedule.expression !== '* * * * *' || schedule.timezone !== 'UTC') { @@ -666,6 +674,49 @@ async function observeContractPublisherTrust(pool) { }); } +test('Cluster Legacy Env application rejects accessor stream factories before PostgreSQL', async () => { + let getterCalls = 0; + let connectCalls = 0; + const repository = new PostgresClusterLegacyEnvMigrationApplicationRepository( + { + async query() { + throw new Error('query must not be called'); + }, + async connect() { + connectCalls += 1; + throw new Error('connect must not be called'); + }, + }, + ); + const streams = {}; + for (const name of ['taskMutations', 'triggerMutations']) { + Object.defineProperty(streams, name, { + enumerable: true, + get() { + getterCalls += 1; + return () => []; + }, + }); + } + await assert.rejects( + repository.apply( + { + applicationId: 'accessor-stream-application', + mutationId: '10000000-0000-4000-8000-000000000001', + projectId: 'accessor-stream-project', + planId: 'accessor-stream-plan', + planDigest: '1'.repeat(64), + taskMutationSetDigest: '2'.repeat(64), + triggerMutationSetDigest: '3'.repeat(64), + }, + streams, + ), + InvalidClusterLegacyEnvMigrationApplicationError, + ); + assert.equal(getterCalls, 0); + assert.equal(connectCalls, 0); +}); + if (!migrationConnectionString) { test('PostgreSQL integration requires QL3_TEST_POSTGRES_URL', { skip: true, @@ -6685,4 +6736,310 @@ if (!migrationConnectionString) { ]); } }); + + test('atomically applies one Legacy Env plan to Task, Trigger and schedule receipts', async () => { + const projectId = `legacy-env-application-${process.pid}`; + const taskId = `legacy task ${process.pid}`; + const triggerId = `legacy trigger ${process.pid}`; + const planId = `legacy-env-application-plan-${process.pid}`; + const planMutationId = `legacy-env-application-plan-mutation-${process.pid}`; + const applicationId = `legacy-env-application-receipt-${process.pid}`; + const taskMutationId = '719f7900-0000-4000-8000-000000000001'; + const triggerMutationId = '719f7900-0000-4000-8000-000000000002'; + const applicationMutationId = '719f7900-0000-4000-8000-000000000003'; + const migrationDatabase = await open('migration'); + const automationDatabase = await open('automation-manager'); + const runtimeDatabase = await open('runtime'); + const adminDatabase = await open('admin'); + try { + await runPostgresMigrations({ pool: migrationDatabase.pool }); + await migrationDatabase.pool.query( + `TRUNCATE TABLE + "ql3"."cluster_legacy_env_migration_application_triggers", + "ql3"."cluster_legacy_env_migration_application_tasks", + "ql3"."cluster_legacy_env_migration_application_receipts", + "ql3"."cluster_legacy_env_migration_plans", + "ql3"."trigger_schedules", "ql3"."triggers", + "ql3"."task_execution_revisions", "ql3"."task_definitions" + CASCADE`, + ); + await migrationDatabase.pool.query( + `INSERT INTO "ql3"."projects" + (id, name, slug, status, version, created_at_ms, updated_at_ms) + VALUES ($1, $1, $2, 'active', 1, 1, 1) + ON CONFLICT (id) DO NOTHING`, + [projectId, `legacy-env-application-${process.pid}`], + ); + const clock = await migrationDatabase.pool.query( + `SELECT floor(extract(epoch FROM clock_timestamp()) * 1000)::bigint + AS "observedAtMs"`, + ); + const occurredAtMs = Number(clock.rows[0].observedAtMs) - 1000; + const initialTask = ( + await new PostgresTaskDefinitionRepository( + migrationDatabase.pool, + ).appendTaskDefinitionRevision({ + projectId, + taskId, + expectedRevision: null, + mutationId: '719f7900-0000-4000-8000-000000000010', + name: 'Legacy command', + description: 'preserved by atomic migration', + kind: 'command', + spec: { + schema: 'qinglong/command@v1', + config: { + command: { + kind: 'argv', + file: '/bin/echo', + args: ['legacy'], + }, + timeoutMs: 30_000, + }, + }, + labels: { 'qinglong.io/source': 'legacy' }, + enabled: true, + occurredAtMs, + }) + ).definition; + const trigger = ( + await new PostgresTriggerRepository( + migrationDatabase.pool, + ).appendTriggerRevision({ + projectId, + triggerId, + expectedRevision: null, + mutationId: '719f7900-0000-4000-8000-000000000011', + taskId, + taskRevision: initialTask.revision, + taskContentDigest: initialTask.contentDigest, + spec: { + schema: 'qinglong/cron@v1', + config: { + expression: '* * * * *', + timezone: 'UTC', + misfirePolicy: 'skip', + }, + }, + enabled: true, + occurredAtMs: occurredAtMs + 1, + }) + ).trigger; + const task = ( + await new PostgresTaskDefinitionRepository( + migrationDatabase.pool, + ).appendTaskDefinitionRevision({ + projectId, + taskId, + expectedRevision: initialTask.revision, + mutationId: '719f7900-0000-4000-8000-000000000012', + name: initialTask.name, + description: initialTask.description, + kind: initialTask.kind, + spec: initialTask.spec, + labels: initialTask.labels, + enabled: initialTask.enabled, + occurredAtMs: occurredAtMs + 2, + }) + ).definition; + + const taskMutations = [ + { + ordinal: 0, + taskId, + previousRevision: task.revision, + previousContentDigest: task.contentDigest, + mutationId: taskMutationId, + }, + ]; + const triggerMutations = [ + { + ordinal: 0, + triggerId, + taskId, + previousRevision: trigger.revision, + previousContentDigest: trigger.contentDigest, + previousTaskRevision: initialTask.revision, + previousTaskContentDigest: initialTask.contentDigest, + mutationId: triggerMutationId, + }, + ]; + const taskDigester = + createClusterLegacyEnvMigrationTaskMutationSetDigester(); + taskMutations.forEach((value) => taskDigester.update(value)); + const taskSet = taskDigester.finish(); + const triggerDigester = + createClusterLegacyEnvMigrationTriggerMutationSetDigester(); + triggerMutations.forEach((value) => triggerDigester.update(value)); + const triggerSet = triggerDigester.finish(); + const secretRef = createSecretRef({ + projectId, + name: 'legacy-env-bundle', + version: 1, + }); + const plan = ( + await new PostgresClusterLegacyEnvMigrationPlanRepository( + automationDatabase.pool, + ).publish({ + planId, + mutationId: planMutationId, + projectId, + source: { + reconciliationBundleDigest: '1'.repeat(64), + decisionDigest: '2'.repeat(64), + candidateSetDigest: '3'.repeat(64), + sourceRowCount: 1, + activeRowCount: 1, + disabledRowCount: 0, + effectiveBindingCount: 1, + }, + target: { + secretRef, + taskRevisionSetDigest: taskSet.revisionSetDigest, + triggerRevisionSetDigest: triggerSet.revisionSetDigest, + taskCount: taskSet.count, + triggerCount: triggerSet.count, + totalEffectiveBytes: 128, + }, + }) + ).plan; + const intent = { + applicationId, + mutationId: applicationMutationId, + projectId, + planId, + planDigest: plan.planDigest, + taskMutationSetDigest: taskSet.mutationSetDigest, + triggerMutationSetDigest: triggerSet.mutationSetDigest, + }; + let taskStreamCalls = 0; + let triggerStreamCalls = 0; + const streams = { + taskMutations() { + taskStreamCalls += 1; + return taskMutations; + }, + triggerMutations() { + triggerStreamCalls += 1; + return triggerMutations; + }, + }; + const applications = + new PostgresClusterLegacyEnvMigrationApplicationRepository( + automationDatabase.pool, + ); + const applied = await applications.apply(intent, streams); + assert.equal(applied.status, 'applied'); + assert.equal(taskStreamCalls, 1); + assert.equal(triggerStreamCalls, 1); + assert.deepEqual( + await applications.findByApplicationId(applicationId), + applied.receipt, + ); + + const state = await automationDatabase.pool.query( + `SELECT task.current_revision AS "taskRevision", + task_revision.spec_json AS "taskSpec", + task_revision.name AS "taskName", + task_revision.description AS "taskDescription", + task_revision.labels_json AS "taskLabels", + execution.plan_json AS "executionPlan", + trigger.current_revision AS "triggerRevision", + trigger_revision.task_revision AS "triggerTaskRevision", + trigger_revision.task_content_digest AS "triggerTaskContentDigest", + schedule.trigger_revision AS "scheduleRevision", + schedule.next_fire_at_ms AS "nextFireAtMs", + schedule.last_scheduled_at_ms AS "lastScheduledAtMs", + schedule.state_version AS "scheduleStateVersion", + schedule.claim_version AS "scheduleClaimVersion" + FROM "ql3"."task_definitions" AS task + JOIN "ql3"."task_definition_revisions" AS task_revision + ON task_revision.project_id = task.project_id + AND task_revision.task_id = task.task_id + AND task_revision.revision = task.current_revision + JOIN "ql3"."task_execution_revisions" AS execution + ON execution.project_id = task.project_id + AND execution.task_id = task.task_id + AND execution.source_revision = task.current_revision + JOIN "ql3"."triggers" AS trigger + ON trigger.project_id = task.project_id + AND trigger.task_id = task.task_id + JOIN "ql3"."trigger_revisions" AS trigger_revision + ON trigger_revision.project_id = trigger.project_id + AND trigger_revision.trigger_id = trigger.trigger_id + AND trigger_revision.revision = trigger.current_revision + JOIN "ql3"."trigger_schedules" AS schedule + ON schedule.project_id = trigger.project_id + AND schedule.trigger_id = trigger.trigger_id + WHERE task.project_id = $1 AND task.task_id = $2 + AND trigger.trigger_id = $3`, + [projectId, taskId, triggerId], + ); + assert.equal(state.rowCount, 1); + const row = state.rows[0]; + assert.equal(row.taskRevision, 3); + assert.equal(row.taskSpec.config.environmentBundleRef, secretRef); + assert.equal(row.taskSpec.config.timeoutMs, 30_000); + assert.equal(row.taskName, task.name); + assert.equal(row.taskDescription, task.description); + assert.deepEqual(row.taskLabels, task.labels); + assert.equal(row.executionPlan.environmentBundleRef, secretRef); + assert.equal(row.triggerRevision, 2); + assert.equal(row.triggerTaskRevision, 3); + assert.match(row.triggerTaskContentDigest, /^[0-9a-f]{64}$/); + assert.equal(row.scheduleRevision, 2); + assert.equal(row.nextFireAtMs, null); + assert.equal(row.lastScheduledAtMs, null); + assert.equal(row.scheduleStateVersion, 1); + assert.equal(row.scheduleClaimVersion, 1); + + const replay = await applications.apply(intent, streams); + assert.equal(replay.status, 'existing'); + assert.deepEqual(replay.receipt, applied.receipt); + assert.equal(taskStreamCalls, 1); + assert.equal(triggerStreamCalls, 1); + + const ledger = await automationDatabase.pool.query( + `SELECT + (SELECT count(*)::integer + FROM "ql3"."cluster_legacy_env_migration_application_receipts" + WHERE application_id = $1) AS receipts, + (SELECT count(*)::integer + FROM "ql3"."cluster_legacy_env_migration_application_tasks" + WHERE application_id = $1) AS tasks, + (SELECT count(*)::integer + FROM "ql3"."cluster_legacy_env_migration_application_triggers" + WHERE application_id = $1) AS triggers`, + [applicationId], + ); + assert.deepEqual(ledger.rows, [{ receipts: 1, tasks: 1, triggers: 1 }]); + await assert.rejects( + automationDatabase.pool.query( + `UPDATE "ql3"."cluster_legacy_env_migration_application_receipts" + SET committed_at_ms = committed_at_ms + WHERE application_id = $1`, + [applicationId], + ), + (error) => error?.code === '42501', + ); + for (const database of [runtimeDatabase, adminDatabase]) { + await assert.rejects( + database.pool.query( + `SELECT application_id + FROM "ql3"."cluster_legacy_env_migration_application_receipts" + WHERE application_id = $1`, + [applicationId], + ), + (error) => error?.code === '42501', + ); + } + } finally { + await Promise.all([ + adminDatabase.close(), + runtimeDatabase.close(), + automationDatabase.close(), + migrationDatabase.close(), + ]); + } + }); } diff --git a/packages/ql3-cluster-postgres/test/postgresqlMigrationDefinitions.test.cjs b/packages/ql3-cluster-postgres/test/postgresqlMigrationDefinitions.test.cjs index cba23cc6..3ef852f0 100644 --- a/packages/ql3-cluster-postgres/test/postgresqlMigrationDefinitions.test.cjs +++ b/packages/ql3-cluster-postgres/test/postgresqlMigrationDefinitions.test.cjs @@ -121,6 +121,7 @@ test('defines the immutable PostgreSQL capability and Run core stream', async () 'pg-0068-cancellation-dispatch-project-keyset', 'pg-0069-worker-session-management-observation', 'pg-0070-cluster-legacy-env-migration-plans', + 'pg-0071-cluster-legacy-env-migration-applications', ], ); for (const migration of postgresqlMainMigrationStream.migrations) { @@ -609,6 +610,11 @@ test('freezes every published PostgreSQL migration checksum', () => { checksum: '7cd6d993f48e7bcebcd62c93571a738d5117c9bcde33b974c5ac8962e2a03fe4', }, + { + id: 'pg-0071-cluster-legacy-env-migration-applications', + checksum: + '82538b5a244011a22afad7f9c6d266a8997da538bdd2bec2ee524997c3b85996', + }, ]; assert.deepEqual( postgresqlMainMigrationStream.migrations.map(({ id, checksum }) => ({ @@ -2489,3 +2495,48 @@ test('advances capability v69 with a content-free Legacy Env plan ledger', async /migration_id = 'pg-0069-worker-session-management-observation'/, ); }); + +test('advances capability v70 with atomic Legacy Env application receipts', async () => { + const migration = migrationById( + 'pg-0071-cluster-legacy-env-migration-applications', + ); + const statements = []; + await migration.up({ + async query(statement) { + statements.push(statement); + return { rows: [] }; + }, + }); + const sql = statements.join('\n'); + for (const table of [ + 'cluster_legacy_env_migration_application_receipts', + 'cluster_legacy_env_migration_application_tasks', + 'cluster_legacy_env_migration_application_triggers', + ]) { + assert.match(sql, new RegExp(`CREATE TABLE "ql3"\\."${table}"`)); + assert.match( + sql, + new RegExp( + `GRANT SELECT, INSERT ON "ql3"\\."${table}" TO ql3_automation_manager`, + ), + ); + } + assert.match(sql, /task_count BETWEEN 1 AND 100000/); + assert.match(sql, /trigger_count BETWEEN 0 AND 500000/); + assert.match(sql, /revision = previous_revision \+ 1/); + assert.match( + sql, + /qinglong\/cluster-legacy-env-migration-application-receipt@v1/, + ); + assert.doesNotMatch( + sql, + /GRANT (?:UPDATE|DELETE|TRUNCATE)[^;]+cluster_legacy_env_migration_application/, + ); + assert.match(sql, /contract_version = 70/); + assert.match(sql, /"cluster_legacy_env_migration_application":1/); + assert.match(sql, /contract_version = 69/); + assert.match( + sql, + /migration_id = 'pg-0070-cluster-legacy-env-migration-plans'/, + ); +}); diff --git a/packages/ql3-cluster-postgres/test/postgresqlSchemaReadiness.test.cjs b/packages/ql3-cluster-postgres/test/postgresqlSchemaReadiness.test.cjs index 3314491e..2ff49b87 100644 --- a/packages/ql3-cluster-postgres/test/postgresqlSchemaReadiness.test.cjs +++ b/packages/ql3-cluster-postgres/test/postgresqlSchemaReadiness.test.cjs @@ -34,6 +34,24 @@ function validPrivileges() { schema_capabilities: [true, false, false, false], projects: [true, true, true, false], cluster_legacy_env_migration_plans: [false, false, false, false], + cluster_legacy_env_migration_application_receipts: [ + false, + false, + false, + false, + ], + cluster_legacy_env_migration_application_tasks: [ + false, + false, + false, + false, + ], + cluster_legacy_env_migration_application_triggers: [ + false, + false, + false, + false, + ], task_definitions: [true, false, false, false], task_definition_revisions: [true, false, false, false], task_execution_revisions: [true, false, false, false], @@ -169,6 +187,24 @@ function validAdminPrivileges() { schema_capabilities: [true, false, false, false], projects: [true, false, false, false], cluster_legacy_env_migration_plans: [false, false, false, false], + cluster_legacy_env_migration_application_receipts: [ + false, + false, + false, + false, + ], + cluster_legacy_env_migration_application_tasks: [ + false, + false, + false, + false, + ], + cluster_legacy_env_migration_application_triggers: [ + false, + false, + false, + false, + ], task_definitions: [false, false, false, false], task_definition_revisions: [false, false, false, false], task_execution_revisions: [false, false, false, false], @@ -457,6 +493,9 @@ function automationManagerPrivileges() { 'plugin_package_task_ownerships', 'plugin_package_identity_keyset_ledger', 'cluster_legacy_env_migration_plans', + 'cluster_legacy_env_migration_application_receipts', + 'cluster_legacy_env_migration_application_tasks', + 'cluster_legacy_env_migration_application_triggers', 'security_audit_events', 'task_definitions', 'task_definition_revisions', @@ -475,6 +514,9 @@ function automationManagerPrivileges() { 'trigger_schedules', 'plugin_package_identity_keyset_ledger', 'cluster_legacy_env_migration_plans', + 'cluster_legacy_env_migration_application_receipts', + 'cluster_legacy_env_migration_application_tasks', + 'cluster_legacy_env_migration_application_triggers', ]); return postgresqlControlSchemaContract.tables.map(({ name: tableName }) => ({ tableName, @@ -838,7 +880,7 @@ test('accepts the exact PostgreSQL control schema and least-privilege runtime ro serverMajor: 16, currentUser: 'ql3_runtime', contractName: 'control-core', - contractVersion: 69, + contractVersion: 70, migrationIds: [ 'pg-0001-schema-capability', 'pg-0002-run-core', @@ -910,6 +952,7 @@ test('accepts the exact PostgreSQL control schema and least-privilege runtime ro 'pg-0068-cancellation-dispatch-project-keyset', 'pg-0069-worker-session-management-observation', 'pg-0070-cluster-legacy-env-migration-plans', + 'pg-0071-cluster-legacy-env-migration-applications', ], }); }); @@ -940,10 +983,10 @@ test('accepts the exact schema and isolated least-privilege admin role', async ( }), ); assert.equal(report.currentUser, 'ql3_admin'); - assert.equal(report.contractVersion, 69); + assert.equal(report.contractVersion, 70); assert.equal( report.migrationIds.at(-1), - 'pg-0070-cluster-legacy-env-migration-plans', + 'pg-0071-cluster-legacy-env-migration-applications', ); }); @@ -956,10 +999,10 @@ test('accepts the isolated least-privilege automation manager role', async () => }), ); assert.equal(report.currentUser, 'ql3_automation_manager'); - assert.equal(report.contractVersion, 69); + assert.equal(report.contractVersion, 70); assert.equal( report.migrationIds.at(-1), - 'pg-0070-cluster-legacy-env-migration-plans', + 'pg-0071-cluster-legacy-env-migration-applications', ); const widened = automationManagerPrivileges(); @@ -988,10 +1031,10 @@ test('accepts the isolated least-privilege human Approval manager role', async ( }), ); assert.equal(report.currentUser, 'ql3_approval_manager'); - assert.equal(report.contractVersion, 69); + assert.equal(report.contractVersion, 70); assert.equal( report.migrationIds.at(-1), - 'pg-0070-cluster-legacy-env-migration-plans', + 'pg-0071-cluster-legacy-env-migration-applications', ); const widened = approvalManagerPrivileges(); @@ -1022,10 +1065,10 @@ test('accepts the isolated least-privilege Run manager role', async () => { }), ); assert.equal(report.currentUser, 'ql3_run_manager'); - assert.equal(report.contractVersion, 69); + assert.equal(report.contractVersion, 70); assert.equal( report.migrationIds.at(-1), - 'pg-0070-cluster-legacy-env-migration-plans', + 'pg-0071-cluster-legacy-env-migration-applications', ); const widened = runManagerPrivileges(); @@ -1186,10 +1229,10 @@ test('accepts the exact schema and isolated Worker ingress role', async () => { }), ); assert.equal(report.currentUser, 'ql3_worker_ingress'); - assert.equal(report.contractVersion, 69); + assert.equal(report.contractVersion, 70); assert.equal( report.migrationIds.at(-1), - 'pg-0070-cluster-legacy-env-migration-plans', + 'pg-0071-cluster-legacy-env-migration-applications', ); }); diff --git a/packages/ql3-runtime-core/package.json b/packages/ql3-runtime-core/package.json index 957920f3..a93e9fb3 100644 --- a/packages/ql3-runtime-core/package.json +++ b/packages/ql3-runtime-core/package.json @@ -266,6 +266,9 @@ "cluster-legacy-env-migration-plan": [ "dist/migration/clusterLegacyEnvMigrationPlan.d.ts" ], + "cluster-legacy-env-migration-application": [ + "dist/migration/clusterLegacyEnvMigrationApplication.d.ts" + ], "run": [ "dist/run/run.d.ts" ], @@ -743,6 +746,11 @@ "require": "./dist/migration/clusterLegacyEnvMigrationPlan.js", "default": "./dist/migration/clusterLegacyEnvMigrationPlan.js" }, + "./cluster-legacy-env-migration-application": { + "types": "./dist/migration/clusterLegacyEnvMigrationApplication.d.ts", + "require": "./dist/migration/clusterLegacyEnvMigrationApplication.js", + "default": "./dist/migration/clusterLegacyEnvMigrationApplication.js" + }, "./cluster-execution-revision": { "types": "./dist/task-definition/clusterExecutionRevision.d.ts", "require": "./dist/task-definition/clusterExecutionRevision.js", diff --git a/packages/ql3-runtime-core/src/migration/clusterLegacyEnvMigrationApplication.ts b/packages/ql3-runtime-core/src/migration/clusterLegacyEnvMigrationApplication.ts new file mode 100644 index 00000000..0b90c12f --- /dev/null +++ b/packages/ql3-runtime-core/src/migration/clusterLegacyEnvMigrationApplication.ts @@ -0,0 +1,660 @@ +import { createHash, type Hash } from 'node:crypto'; + +import { + MAX_CLUSTER_LEGACY_ENV_TASKS, + MAX_CLUSTER_LEGACY_ENV_TRIGGERS, +} from './clusterLegacyEnvMigrationPlan'; +import { parseSecretRef } from '../secret/secretReference'; + +export const CLUSTER_LEGACY_ENV_MIGRATION_APPLICATION_RECEIPT_SCHEMA = + 'qinglong/cluster-legacy-env-migration-application-receipt@v1' as const; +export const MAX_CLUSTER_LEGACY_ENV_MIGRATION_RECEIPT_JSON_BYTES = 8 * 1024; + +export interface ClusterLegacyEnvMigrationTaskMutation { + readonly ordinal: number; + readonly taskId: string; + readonly previousRevision: number; + readonly previousContentDigest: string; + readonly mutationId: string; +} + +export interface ClusterLegacyEnvMigrationTriggerMutation { + readonly ordinal: number; + readonly triggerId: string; + readonly taskId: string; + readonly previousRevision: number; + readonly previousContentDigest: string; + readonly previousTaskRevision: number; + readonly previousTaskContentDigest: string; + readonly mutationId: string; +} + +export interface ClusterLegacyEnvMigrationApplicationIntent { + readonly applicationId: string; + readonly mutationId: string; + readonly projectId: string; + readonly planId: string; + readonly planDigest: string; + readonly taskMutationSetDigest: string; + readonly triggerMutationSetDigest: string; +} + +export interface ClusterLegacyEnvMigrationApplicationReceipt + extends ClusterLegacyEnvMigrationApplicationIntent { + readonly schema: typeof CLUSTER_LEGACY_ENV_MIGRATION_APPLICATION_RECEIPT_SCHEMA; + readonly environmentBundleRef: string; + readonly taskRevisionSetDigest: string; + readonly triggerRevisionSetDigest: string; + readonly taskCount: number; + readonly triggerCount: number; + readonly committedAtMs: number; + readonly receiptDigest: string; +} + +export interface ClusterLegacyEnvMigrationMutationStreams { + readonly taskMutations: () => + | Iterable + | AsyncIterable; + readonly triggerMutations: () => + | Iterable + | AsyncIterable; +} + +export interface ClusterLegacyEnvMigrationApplicationRepository { + apply( + intent: Readonly, + streams: Readonly, + ): Promise< + Readonly<{ + status: 'applied' | 'existing'; + receipt: Readonly; + }> + >; + findByApplicationId( + applicationId: string, + ): Promise | null>; +} + +export interface ClusterLegacyEnvMigrationTaskMutationSetDigestResult { + readonly count: number; + readonly revisionSetDigest: string; + readonly mutationSetDigest: string; +} + +export interface ClusterLegacyEnvMigrationTriggerMutationSetDigestResult { + readonly count: number; + readonly revisionSetDigest: string; + readonly mutationSetDigest: string; +} + +export class InvalidClusterLegacyEnvMigrationApplicationError extends TypeError { + readonly code = 'CLUSTER_LEGACY_ENV_MIGRATION_APPLICATION_INVALID'; + + constructor(message: string) { + super(`Cluster Legacy Env migration application is invalid: ${message}`); + this.name = 'InvalidClusterLegacyEnvMigrationApplicationError'; + } +} + +export class ClusterLegacyEnvMigrationApplicationConflictError extends Error { + readonly code = 'CLUSTER_LEGACY_ENV_MIGRATION_APPLICATION_CONFLICT'; + + constructor() { + super( + 'Cluster Legacy Env migration application conflicts with durable state', + ); + this.name = 'ClusterLegacyEnvMigrationApplicationConflictError'; + } +} + +export class ClusterLegacyEnvMigrationApplicationUnavailableError extends Error { + readonly code = 'CLUSTER_LEGACY_ENV_MIGRATION_APPLICATION_UNAVAILABLE'; + + constructor(options?: ErrorOptions) { + super('Cluster Legacy Env migration application is unavailable', options); + this.name = 'ClusterLegacyEnvMigrationApplicationUnavailableError'; + } +} + +const ID_PATTERN = /^[A-Za-z0-9][A-Za-z0-9._:-]{0,127}$/; +const MUTATION_ID_PATTERN = + /^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$/; +const DIGEST_PATTERN = /^[0-9a-f]{64}$/; +const TASK_REVISION_SET_DOMAIN = Buffer.from( + 'qinglong/cluster-legacy-env-task-revision-set@v1\0', + 'utf8', +); +const TASK_MUTATION_SET_DOMAIN = Buffer.from( + 'qinglong/cluster-legacy-env-task-mutation-set@v1\0', + 'utf8', +); +const TRIGGER_REVISION_SET_DOMAIN = Buffer.from( + 'qinglong/cluster-legacy-env-trigger-revision-set@v1\0', + 'utf8', +); +const TRIGGER_MUTATION_SET_DOMAIN = Buffer.from( + 'qinglong/cluster-legacy-env-trigger-mutation-set@v1\0', + 'utf8', +); +const RECEIPT_DIGEST_DOMAIN = Buffer.from( + 'qinglong/cluster-legacy-env-migration-application-receipt-digest@v1\0', + 'utf8', +); + +function invalid(message: string): never { + throw new InvalidClusterLegacyEnvMigrationApplicationError(message); +} + +function record(value: unknown, label: string): Record { + if ( + !value || + typeof value !== 'object' || + Array.isArray(value) || + (Object.getPrototypeOf(value) !== Object.prototype && + Object.getPrototypeOf(value) !== null) + ) { + return invalid(`${label} must be an object`); + } + const descriptors = Object.getOwnPropertyDescriptors(value); + if ( + Object.values(descriptors).some( + (descriptor) => + descriptor.get !== undefined || + descriptor.set !== undefined || + descriptor.enumerable !== true, + ) + ) { + return invalid(`${label} must contain enumerable data properties`); + } + return value as Record; +} + +function exactKeys( + value: object, + expected: readonly string[], + label: string, +): void { + const keys = Reflect.ownKeys(value); + const canonical = [...expected].sort(); + if ( + keys.some((key) => typeof key !== 'string') || + keys.length !== canonical.length || + keys + .map(String) + .sort() + .some((key, index) => key !== canonical[index]) + ) { + invalid(`${label} shape is invalid`); + } +} + +export function assertClusterLegacyEnvMigrationApplicationIdentifier( + value: unknown, + label: 'applicationId' | 'planId' | 'projectId', +): string { + if (typeof value !== 'string' || !ID_PATTERN.test(value)) { + return invalid(`${label} is invalid`); + } + return value; +} + +function entityIdentifier( + value: unknown, + label: 'taskId' | 'triggerId', +): string { + if ( + typeof value !== 'string' || + value.length < 1 || + Buffer.byteLength(value, 'utf8') > 128 || + /[\u0000-\u001f\u007f]/.test(value) + ) { + return invalid(`${label} is invalid`); + } + return value; +} + +function mutationId(value: unknown, label: string): string { + if (typeof value !== 'string' || !MUTATION_ID_PATTERN.test(value)) { + return invalid(`${label} is invalid`); + } + return value; +} + +function digest(value: unknown, label: string): string { + if (typeof value !== 'string' || !DIGEST_PATTERN.test(value)) { + return invalid(`${label} is invalid`); + } + return value; +} + +function revision(value: unknown, label: string): number { + if ( + !Number.isSafeInteger(value) || + (value as number) < 1 || + (value as number) >= 2_147_483_647 + ) { + return invalid(`${label} is invalid`); + } + return value as number; +} + +function ordinal(value: unknown, label: string, maximum: number): number { + if ( + !Number.isSafeInteger(value) || + (value as number) < 0 || + (value as number) >= maximum + ) { + return invalid(`${label} is invalid`); + } + return value as number; +} + +function timestamp(value: unknown): number { + if (!Number.isSafeInteger(value) || (value as number) < 0) { + return invalid('committedAtMs is invalid'); + } + return value as number; +} + +export function normalizeClusterLegacyEnvMigrationTaskMutation( + value: ClusterLegacyEnvMigrationTaskMutation, +): Readonly { + const candidate = record(value, 'Task mutation'); + exactKeys( + candidate, + [ + 'mutationId', + 'ordinal', + 'previousContentDigest', + 'previousRevision', + 'taskId', + ], + 'Task mutation', + ); + return Object.freeze({ + ordinal: ordinal( + value.ordinal, + 'Task mutation ordinal', + MAX_CLUSTER_LEGACY_ENV_TASKS, + ), + taskId: entityIdentifier(value.taskId, 'taskId'), + previousRevision: revision(value.previousRevision, 'previousRevision'), + previousContentDigest: digest( + value.previousContentDigest, + 'previousContentDigest', + ), + mutationId: mutationId(value.mutationId, 'Task mutationId'), + }); +} + +export function normalizeClusterLegacyEnvMigrationTriggerMutation( + value: ClusterLegacyEnvMigrationTriggerMutation, +): Readonly { + const candidate = record(value, 'Trigger mutation'); + exactKeys( + candidate, + [ + 'mutationId', + 'ordinal', + 'previousContentDigest', + 'previousRevision', + 'previousTaskContentDigest', + 'previousTaskRevision', + 'taskId', + 'triggerId', + ], + 'Trigger mutation', + ); + return Object.freeze({ + ordinal: ordinal( + value.ordinal, + 'Trigger mutation ordinal', + MAX_CLUSTER_LEGACY_ENV_TRIGGERS, + ), + triggerId: entityIdentifier(value.triggerId, 'triggerId'), + taskId: entityIdentifier(value.taskId, 'taskId'), + previousRevision: revision(value.previousRevision, 'previousRevision'), + previousContentDigest: digest( + value.previousContentDigest, + 'previousContentDigest', + ), + previousTaskRevision: revision( + value.previousTaskRevision, + 'previousTaskRevision', + ), + previousTaskContentDigest: digest( + value.previousTaskContentDigest, + 'previousTaskContentDigest', + ), + mutationId: mutationId(value.mutationId, 'Trigger mutationId'), + }); +} + +export function normalizeClusterLegacyEnvMigrationApplicationIntent( + value: ClusterLegacyEnvMigrationApplicationIntent, +): Readonly { + const candidate = record(value, 'application intent'); + exactKeys( + candidate, + [ + 'applicationId', + 'mutationId', + 'planDigest', + 'planId', + 'projectId', + 'taskMutationSetDigest', + 'triggerMutationSetDigest', + ], + 'application intent', + ); + return Object.freeze({ + applicationId: assertClusterLegacyEnvMigrationApplicationIdentifier( + value.applicationId, + 'applicationId', + ), + mutationId: mutationId(value.mutationId, 'mutationId'), + projectId: assertClusterLegacyEnvMigrationApplicationIdentifier( + value.projectId, + 'projectId', + ), + planId: assertClusterLegacyEnvMigrationApplicationIdentifier( + value.planId, + 'planId', + ), + planDigest: digest(value.planDigest, 'planDigest'), + taskMutationSetDigest: digest( + value.taskMutationSetDigest, + 'taskMutationSetDigest', + ), + triggerMutationSetDigest: digest( + value.triggerMutationSetDigest, + 'triggerMutationSetDigest', + ), + }); +} + +function updateHash(hash: Hash, value: object): void { + hash.update(JSON.stringify(value), 'utf8').update('\n', 'utf8'); +} + +function finalizeHash(hash: Hash, count: number): string { + updateHash(hash, { count }); + return hash.digest('hex'); +} + +export function createClusterLegacyEnvMigrationTaskMutationSetDigester(): Readonly<{ + update( + value: ClusterLegacyEnvMigrationTaskMutation, + ): Readonly; + finish(): Readonly; +}> { + const revisionHash = createHash('sha256').update(TASK_REVISION_SET_DOMAIN); + const mutationHash = createHash('sha256').update(TASK_MUTATION_SET_DOMAIN); + let count = 0; + let previousId: string | undefined; + let finished = false; + return Object.freeze({ + update(value: ClusterLegacyEnvMigrationTaskMutation) { + if (finished) return invalid('Task mutation digester is finished'); + const item = normalizeClusterLegacyEnvMigrationTaskMutation(value); + if ( + item.ordinal !== count || + (previousId !== undefined && item.taskId <= previousId) + ) { + return invalid( + 'Task mutations must be contiguous and ordered by taskId', + ); + } + updateHash(revisionHash, { + ordinal: item.ordinal, + taskId: item.taskId, + revision: item.previousRevision, + contentDigest: item.previousContentDigest, + }); + updateHash(mutationHash, item); + previousId = item.taskId; + count += 1; + return item; + }, + finish() { + if (finished) return invalid('Task mutation digester is finished'); + finished = true; + return Object.freeze({ + count, + revisionSetDigest: finalizeHash(revisionHash, count), + mutationSetDigest: finalizeHash(mutationHash, count), + }); + }, + }); +} + +export function createClusterLegacyEnvMigrationTriggerMutationSetDigester(): Readonly<{ + update( + value: ClusterLegacyEnvMigrationTriggerMutation, + ): Readonly; + finish(): Readonly; +}> { + const revisionHash = createHash('sha256').update(TRIGGER_REVISION_SET_DOMAIN); + const mutationHash = createHash('sha256').update(TRIGGER_MUTATION_SET_DOMAIN); + let count = 0; + let previousId: string | undefined; + let finished = false; + return Object.freeze({ + update(value: ClusterLegacyEnvMigrationTriggerMutation) { + if (finished) return invalid('Trigger mutation digester is finished'); + const item = normalizeClusterLegacyEnvMigrationTriggerMutation(value); + if ( + item.ordinal !== count || + (previousId !== undefined && item.triggerId <= previousId) + ) { + return invalid( + 'Trigger mutations must be contiguous and ordered by triggerId', + ); + } + updateHash(revisionHash, { + ordinal: item.ordinal, + triggerId: item.triggerId, + taskId: item.taskId, + revision: item.previousRevision, + contentDigest: item.previousContentDigest, + taskRevision: item.previousTaskRevision, + taskContentDigest: item.previousTaskContentDigest, + }); + updateHash(mutationHash, item); + previousId = item.triggerId; + count += 1; + return item; + }, + finish() { + if (finished) return invalid('Trigger mutation digester is finished'); + finished = true; + return Object.freeze({ + count, + revisionSetDigest: finalizeHash(revisionHash, count), + mutationSetDigest: finalizeHash(mutationHash, count), + }); + }, + }); +} + +function unsignedReceipt( + value: Omit, +): object { + return { + schema: value.schema, + applicationId: value.applicationId, + mutationId: value.mutationId, + projectId: value.projectId, + planId: value.planId, + planDigest: value.planDigest, + environmentBundleRef: value.environmentBundleRef, + taskRevisionSetDigest: value.taskRevisionSetDigest, + triggerRevisionSetDigest: value.triggerRevisionSetDigest, + taskMutationSetDigest: value.taskMutationSetDigest, + triggerMutationSetDigest: value.triggerMutationSetDigest, + taskCount: value.taskCount, + triggerCount: value.triggerCount, + committedAtMs: value.committedAtMs, + }; +} + +export function clusterLegacyEnvMigrationApplicationReceiptDigest( + value: Omit, +): string { + return createHash('sha256') + .update(RECEIPT_DIGEST_DOMAIN) + .update(JSON.stringify(unsignedReceipt(value)), 'utf8') + .digest('hex'); +} + +export function createClusterLegacyEnvMigrationApplicationReceipt( + value: Omit< + ClusterLegacyEnvMigrationApplicationReceipt, + 'schema' | 'receiptDigest' + >, +): Readonly { + const candidate = record(value, 'application receipt input'); + exactKeys( + candidate, + [ + 'applicationId', + 'committedAtMs', + 'environmentBundleRef', + 'mutationId', + 'planDigest', + 'planId', + 'projectId', + 'taskCount', + 'taskMutationSetDigest', + 'taskRevisionSetDigest', + 'triggerCount', + 'triggerMutationSetDigest', + 'triggerRevisionSetDigest', + ], + 'application receipt input', + ); + const intent = normalizeClusterLegacyEnvMigrationApplicationIntent({ + applicationId: value.applicationId, + mutationId: value.mutationId, + projectId: value.projectId, + planId: value.planId, + planDigest: value.planDigest, + taskMutationSetDigest: value.taskMutationSetDigest, + triggerMutationSetDigest: value.triggerMutationSetDigest, + }); + let environmentBundleReference; + try { + environmentBundleReference = parseSecretRef(value.environmentBundleRef); + } catch { + return invalid('environmentBundleRef is invalid'); + } + if ( + !Number.isSafeInteger(value.taskCount) || + value.taskCount < 1 || + value.taskCount > MAX_CLUSTER_LEGACY_ENV_TASKS || + !Number.isSafeInteger(value.triggerCount) || + value.triggerCount < 0 || + value.triggerCount > MAX_CLUSTER_LEGACY_ENV_TRIGGERS || + environmentBundleReference.projectId !== intent.projectId || + environmentBundleReference.version === undefined + ) { + return invalid('receipt target is invalid'); + } + const unsigned = Object.freeze({ + schema: CLUSTER_LEGACY_ENV_MIGRATION_APPLICATION_RECEIPT_SCHEMA, + ...intent, + environmentBundleRef: value.environmentBundleRef, + taskRevisionSetDigest: digest( + value.taskRevisionSetDigest, + 'taskRevisionSetDigest', + ), + triggerRevisionSetDigest: digest( + value.triggerRevisionSetDigest, + 'triggerRevisionSetDigest', + ), + taskCount: value.taskCount, + triggerCount: value.triggerCount, + committedAtMs: timestamp(value.committedAtMs), + }); + const receipt = Object.freeze({ + ...unsigned, + receiptDigest: clusterLegacyEnvMigrationApplicationReceiptDigest(unsigned), + }); + if ( + Buffer.byteLength(JSON.stringify(receipt), 'utf8') > + MAX_CLUSTER_LEGACY_ENV_MIGRATION_RECEIPT_JSON_BYTES + ) { + return invalid('encoded receipt exceeds the size limit'); + } + return receipt; +} + +export function normalizeClusterLegacyEnvMigrationApplicationReceipt( + value: ClusterLegacyEnvMigrationApplicationReceipt, +): Readonly { + const candidate = record(value, 'application receipt'); + exactKeys( + candidate, + [ + 'applicationId', + 'committedAtMs', + 'environmentBundleRef', + 'mutationId', + 'planDigest', + 'planId', + 'projectId', + 'receiptDigest', + 'schema', + 'taskCount', + 'taskMutationSetDigest', + 'taskRevisionSetDigest', + 'triggerCount', + 'triggerMutationSetDigest', + 'triggerRevisionSetDigest', + ], + 'application receipt', + ); + if ( + value.schema !== CLUSTER_LEGACY_ENV_MIGRATION_APPLICATION_RECEIPT_SCHEMA + ) { + return invalid('receipt schema is invalid'); + } + const expected = createClusterLegacyEnvMigrationApplicationReceipt({ + applicationId: value.applicationId, + mutationId: value.mutationId, + projectId: value.projectId, + planId: value.planId, + planDigest: value.planDigest, + environmentBundleRef: value.environmentBundleRef, + taskRevisionSetDigest: value.taskRevisionSetDigest, + triggerRevisionSetDigest: value.triggerRevisionSetDigest, + taskMutationSetDigest: value.taskMutationSetDigest, + triggerMutationSetDigest: value.triggerMutationSetDigest, + taskCount: value.taskCount, + triggerCount: value.triggerCount, + committedAtMs: value.committedAtMs, + }); + if (value.receiptDigest !== expected.receiptDigest) { + return invalid('receiptDigest does not match receipt'); + } + return expected; +} + +export function clusterLegacyEnvMigrationApplicationReceiptMatchesIntent( + receiptValue: ClusterLegacyEnvMigrationApplicationReceipt, + intentValue: ClusterLegacyEnvMigrationApplicationIntent, +): boolean { + const receipt = + normalizeClusterLegacyEnvMigrationApplicationReceipt(receiptValue); + const intent = + normalizeClusterLegacyEnvMigrationApplicationIntent(intentValue); + return ( + receipt.applicationId === intent.applicationId && + receipt.mutationId === intent.mutationId && + receipt.projectId === intent.projectId && + receipt.planId === intent.planId && + receipt.planDigest === intent.planDigest && + receipt.taskMutationSetDigest === intent.taskMutationSetDigest && + receipt.triggerMutationSetDigest === intent.triggerMutationSetDigest + ); +} diff --git a/packages/ql3-runtime-core/test/clusterLegacyEnvMigrationApplication.test.cjs b/packages/ql3-runtime-core/test/clusterLegacyEnvMigrationApplication.test.cjs new file mode 100644 index 00000000..a91d1bfe --- /dev/null +++ b/packages/ql3-runtime-core/test/clusterLegacyEnvMigrationApplication.test.cjs @@ -0,0 +1,222 @@ +const assert = require('node:assert/strict'); +const { test } = require('node:test'); + +const { + CLUSTER_LEGACY_ENV_MIGRATION_APPLICATION_RECEIPT_SCHEMA, + InvalidClusterLegacyEnvMigrationApplicationError, + clusterLegacyEnvMigrationApplicationReceiptMatchesIntent, + createClusterLegacyEnvMigrationApplicationReceipt, + createClusterLegacyEnvMigrationTaskMutationSetDigester, + createClusterLegacyEnvMigrationTriggerMutationSetDigester, + normalizeClusterLegacyEnvMigrationApplicationReceipt, +} = require('@qinglong/runtime-core/cluster-legacy-env-migration-application'); +const { createSecretRef } = require('@qinglong/runtime-core/secret-reference'); + +const IDS = Object.freeze({ + application: '11111111-1111-4111-8111-111111111111', + applicationMutation: '22222222-2222-4222-8222-222222222222', + taskOneMutation: '33333333-3333-4333-8333-333333333333', + taskTwoMutation: '44444444-4444-4444-8444-444444444444', + triggerMutation: '55555555-5555-4555-8555-555555555555', +}); + +function taskMutations() { + return [ + { + ordinal: 0, + taskId: 'task-a', + previousRevision: 3, + previousContentDigest: '1'.repeat(64), + mutationId: IDS.taskOneMutation, + }, + { + ordinal: 1, + taskId: 'task-b', + previousRevision: 8, + previousContentDigest: '2'.repeat(64), + mutationId: IDS.taskTwoMutation, + }, + ]; +} + +function triggerMutations() { + return [ + { + ordinal: 0, + triggerId: 'trigger-a', + taskId: 'task-a', + previousRevision: 4, + previousContentDigest: '3'.repeat(64), + previousTaskRevision: 3, + previousTaskContentDigest: '1'.repeat(64), + mutationId: IDS.triggerMutation, + }, + ]; +} + +function digestTaskMutations(values = taskMutations()) { + const digester = createClusterLegacyEnvMigrationTaskMutationSetDigester(); + for (const value of values) digester.update(value); + return digester.finish(); +} + +function digestTriggerMutations(values = triggerMutations()) { + const digester = createClusterLegacyEnvMigrationTriggerMutationSetDigester(); + for (const value of values) digester.update(value); + return digester.finish(); +} + +function receiptInput(overrides = {}) { + const task = digestTaskMutations(); + const trigger = digestTriggerMutations(); + return { + applicationId: 'legacy-env-application-a', + mutationId: IDS.applicationMutation, + projectId: 'project-a', + planId: 'legacy-env-plan-a', + planDigest: '6'.repeat(64), + environmentBundleRef: createSecretRef({ + projectId: 'project-a', + name: 'legacy-env-bundle', + version: 7, + }), + taskRevisionSetDigest: task.revisionSetDigest, + triggerRevisionSetDigest: trigger.revisionSetDigest, + taskMutationSetDigest: task.mutationSetDigest, + triggerMutationSetDigest: trigger.mutationSetDigest, + taskCount: task.count, + triggerCount: trigger.count, + committedAtMs: 12_345, + ...overrides, + }; +} + +test('streams canonical Task and Trigger revision and mutation sets', () => { + const firstTask = digestTaskMutations(); + const secondTask = digestTaskMutations(); + const firstTrigger = digestTriggerMutations(); + const secondTrigger = digestTriggerMutations(); + + assert.deepEqual(firstTask, secondTask); + assert.deepEqual(firstTrigger, secondTrigger); + assert.equal(firstTask.count, 2); + assert.equal(firstTrigger.count, 1); + assert.notEqual(firstTask.revisionSetDigest, firstTask.mutationSetDigest); + assert.notEqual( + firstTrigger.revisionSetDigest, + firstTrigger.mutationSetDigest, + ); +}); + +test('rejects gaps, duplicate ordering, malformed UUIDs and unknown fields', () => { + const gap = taskMutations(); + gap[1] = { ...gap[1], ordinal: 2 }; + assert.throws(() => digestTaskMutations(gap), { + name: 'InvalidClusterLegacyEnvMigrationApplicationError', + }); + + const unordered = taskMutations(); + unordered[1] = { ...unordered[1], taskId: 'task-a' }; + assert.throws( + () => digestTaskMutations(unordered), + InvalidClusterLegacyEnvMigrationApplicationError, + ); + + const malformedMutation = taskMutations(); + malformedMutation[0] = { ...malformedMutation[0], mutationId: 'not-a-uuid' }; + assert.throws( + () => digestTaskMutations(malformedMutation), + InvalidClusterLegacyEnvMigrationApplicationError, + ); + + const extra = taskMutations(); + extra[0] = { ...extra[0], command: 'must-not-enter-ledger' }; + assert.throws( + () => digestTaskMutations(extra), + InvalidClusterLegacyEnvMigrationApplicationError, + ); +}); + +test('creates a frozen content-free receipt and detects tampering', () => { + const receipt = createClusterLegacyEnvMigrationApplicationReceipt( + receiptInput(), + ); + assert.equal( + receipt.schema, + CLUSTER_LEGACY_ENV_MIGRATION_APPLICATION_RECEIPT_SCHEMA, + ); + assert.equal(Object.isFrozen(receipt), true); + assert.match(receipt.receiptDigest, /^[0-9a-f]{64}$/); + assert.deepEqual( + normalizeClusterLegacyEnvMigrationApplicationReceipt(receipt), + receipt, + ); + assert.equal( + clusterLegacyEnvMigrationApplicationReceiptMatchesIntent(receipt, { + applicationId: receipt.applicationId, + mutationId: receipt.mutationId, + projectId: receipt.projectId, + planId: receipt.planId, + planDigest: receipt.planDigest, + taskMutationSetDigest: receipt.taskMutationSetDigest, + triggerMutationSetDigest: receipt.triggerMutationSetDigest, + }), + true, + ); + assert.throws( + () => + normalizeClusterLegacyEnvMigrationApplicationReceipt({ + ...receipt, + taskCount: receipt.taskCount + 1, + }), + InvalidClusterLegacyEnvMigrationApplicationError, + ); +}); + +test('requires one Task and a canonical version-pinned bundle in the same Project', () => { + assert.throws( + () => + createClusterLegacyEnvMigrationApplicationReceipt( + receiptInput({ taskCount: 0 }), + ), + InvalidClusterLegacyEnvMigrationApplicationError, + ); + assert.throws( + () => + createClusterLegacyEnvMigrationApplicationReceipt( + receiptInput({ + environmentBundleRef: createSecretRef({ + projectId: 'project-a', + name: 'legacy-env-bundle', + }), + }), + ), + InvalidClusterLegacyEnvMigrationApplicationError, + ); + assert.throws( + () => + createClusterLegacyEnvMigrationApplicationReceipt( + receiptInput({ + environmentBundleRef: createSecretRef({ + projectId: 'project-b', + name: 'legacy-env-bundle', + version: 7, + }), + }), + ), + InvalidClusterLegacyEnvMigrationApplicationError, + ); +}); + +test('supports an empty Trigger mutation set without weakening Task authority', () => { + const trigger = digestTriggerMutations([]); + const receipt = createClusterLegacyEnvMigrationApplicationReceipt( + receiptInput({ + triggerRevisionSetDigest: trigger.revisionSetDigest, + triggerMutationSetDigest: trigger.mutationSetDigest, + triggerCount: 0, + }), + ); + assert.equal(receipt.taskCount, 2); + assert.equal(receipt.triggerCount, 0); +}); diff --git a/test/back/ql3PackageBoundaryAudit.test.cjs b/test/back/ql3PackageBoundaryAudit.test.cjs index 2f8bf5d4..04205a7e 100644 --- a/test/back/ql3PackageBoundaryAudit.test.cjs +++ b/test/back/ql3PackageBoundaryAudit.test.cjs @@ -299,10 +299,10 @@ test('current QL3 workspace has exactly eighteen reviewed package boundaries', ( rootSourceFileRoles: runtimeCore.rootSourceFileRoles, }, { - sourceFiles: 173, + sourceFiles: 174, rootSourceFiles: 1, rootSourceLines: 160, - nestedSourceFiles: 172, + nestedSourceFiles: 173, rootSourceFileRoles: { 'index.ts': 'public_export' }, }, ); @@ -421,10 +421,10 @@ test('current QL3 workspace has exactly eighteen reviewed package boundaries', ( rootSourceFileRoles: clusterPostgres.rootSourceFileRoles, }, { - sourceFiles: 177, + sourceFiles: 179, rootSourceFiles: 1, rootSourceLines: 126, - nestedSourceFiles: 176, + nestedSourceFiles: 178, rootSourceFileRoles: { 'index.ts': 'public_export' }, }, );