From d8489d8a0b38b880dcc2991267620450b62b509b Mon Sep 17 00:00:00 2001 From: whyour Date: Fri, 14 Aug 2026 10:48:37 +0800 Subject: [PATCH] feat(ql3): recover terminal secret action jobs --- docs/QINGLONG_3_0_ARCHITECTURE_RFC.md | 1 + ...ransition-plugin-package-secret-binding.md | 4 + .../executor/pluginPackageExecutorProcess.ts | 12 + ...PackageKubernetesSecretActionController.ts | 284 ++++++++- .../pluginPackageExecutorProcess.test.cjs | 6 + ...eKubernetesSecretActionController.test.cjs | 552 +++++++++++++++++- 6 files changed, 833 insertions(+), 26 deletions(-) diff --git a/docs/QINGLONG_3_0_ARCHITECTURE_RFC.md b/docs/QINGLONG_3_0_ARCHITECTURE_RFC.md index dbbbc84c..fd1cacae 100644 --- a/docs/QINGLONG_3_0_ARCHITECTURE_RFC.md +++ b/docs/QINGLONG_3_0_ARCHITECTURE_RFC.md @@ -32,6 +32,7 @@ - D-306B2/ADR-0396(进行中):Secret rebind/rotation/revocation 不更新历史 binding,而是作为下一 Package generation 的 activation 前置事实。共享 transition plan v1 同时绑定上一 active target、可选的上一 binding、durable install history 的最后尝试 generation、新 target、可选下一 binding plan、逐 requirement 与 SecretRef 差异及独立 digest;上一 active Manifest 没有 Secret requirement 时 binding 可空,但 target/lock/generation lineage 不可省略。失败 install 也永久消耗 generation,重试必须使用 `lastAttemptGeneration + 1`,active lineage 继续由 `previousActiveLockDigest` 指回旧代。SQLite capability v49(0097/0098)已完成 immutable transition receipt ledger、Local Owner staged `plan→execute`、单事务 binding/audit/receipt 和 activation prerequisite。PostgreSQL capability v62(pg-0063)具有 receipt ledger,Cluster management 与 package-executor 也已完成 separation-of-duty 产品编排:executor 在一个 SERIALIZABLE transaction 中复验 current staged head、上一 active lineage、durable 最大 generation 与可选上一 binding,并原子提交可选目标 binding 和 immutable receipt;数据库 trigger 和最小角色 readiness 防止绕过,recovery/直接 activation 缺 receipt 均失败关闭。Cluster startup recovery 现把 binding/transition receipt 编译为 content-blind Kubernetes active pointer v3:只保存 source Secret 名、不可逆 SHA-256 key/path、`0440`、binding/receipt/projection digest 和逻辑 assignment,不保存 SecretRef 或明文;同一次 ConfigMap `resourceVersion` CAS 同时切换 Package generation 与投影声明。rotate 只投影下一代 exact key;revoke 生成显式空 projection,Pod volume renderer 对空项返回不挂载,避免 Kubernetes 空 `items` 被解释为投影全部 key。无 Secret generation 继续发布兼容 v2;publisher 不获得 Secret `get/list`,不新增 watcher/controller。响应丢失通过 durable pointer inspect 精确收敛,projection source 不可用或 digest 漂移时旧 active pointer 保持不变。真实三节点 K3s `v1.34.3+k3s1` 已完成两个受限 actor 同一 resourceVersion 的 v3 rotation 竞争:1 成功/1 冲突、最终恰好 1 pointer、1 个 exact projection item、Secret API read denied,projection digest `22add8965accf3f736935167963b9dbdeab8fba05f739f24d39dec941aae9680`、transition receipt digest `355b89ecd54422af33fa573780c8b70ed41da4df226ec301bdd5b6c71de609e1`;随后生产 renderer 驱动 2 副本 workload 分布到 2 个节点,源 Secret 同时保留目标 key 与 decoy key,而 Pod 只看到 exact path、文件模式 `0440`。第三个同权限 actor 以 generation 3 transition receipt `252e8cd1d0f8a2b63f8c0af6247861c08e5b7ae8886c0d65cd90537be3e9f9ac` 发布显式空 projection;Deployment 滚动产生全新 Pod UID,两个新 Pod 均无 Secret volume/mount 和投影根目录,源 Secret 保留,actor 的 Secret `get` 仍为 403;临时容器/网络已清理。实现没有新增 workspace package,Secret projection/renderer 内聚在既有 `cluster-admin/plugin-package/secret-binding`,18-package boundary 仍为 `singleSourcePackages=[]`、`shallowSourcePackages=[]`,cluster dependency 与 edge import 审计无 finding。阶段提交 `9d7431c2` 后,完整 18-package 串行 build/test 退出 0,backend 1,192 pass/2 skip/0 fail;PostgreSQL 18.4 physical HA 125 gate、timeline `1→2` 通过,报告 SHA-256 为 `2f0d1107cc6d447868bfaeb3650284acd803924c6cd3cdee978b3eb5882eb26c`;Edge 本机观测的模块加载 RSS 增量约 6.0 MiB、1 万行输出峰值增量约 3.9 MiB,但尚无固定物理低配门。executor 精确投影重构已开始落地:共享 dispatcher 增加不扫描队列的 `dispatchById`,在租约领取前拒绝未配置 handler 的动作;executor 的 exact mode 跳过全部 Approval consumer,只执行一个 durable dispatch。既有 `cluster-admin/plugin-package/executor` 内新增 immutable Kubernetes Job renderer,以 dispatch 与 approval plan 联合校验生成确定性名称,只接受 digest-pinned 镜像,并只挂载去重后的 exact SHA-256 `items`(`0440`、`optional:false`、无 API token);零 SecretRef 的 revoke 使用 1 KiB `emptyDir`,绝不以空 `items` 误挂全量 Secret。常规 batch CronJob 已移除 Package values volume、Secret root 和 dispatch-id authority;未挂 Secret root 时 dispatcher 根本不注册 binding/transition handler,因此相关 durable execution 保持 pending,不会被错误领取后 blocked。该切片不增加 package、依赖或常驻进程,并把 action Job 内存 request/limit 固定为 48/192 MiB、数据库池固定 1,适合低配节点。生产 rollout controller 尚不能直接取得不受约束的 `jobs.create`:该权限可通过自定义 PodSpec 间接放大到任意 Secret/镜像,必须先用 admission policy 把 ServiceAccount、digest 镜像、command、source Secret、exact item/path 与数据库 SecretRef 固定,再接入 create/get-only adapter 和恢复状态机。本阶段完整 18-package clean build/test 退出 0;backend 1195 项为 1193 pass、2 条条件 skip、0 fail;package boundary 确认为 18 个 package、`singleSourcePackages=[]`、`shallowSourcePackages=[]`,cluster dependency、edge import 和 cluster deployment 审计均无 finding。Edge arm64 本机观测模块加载 RSS 增量 8,945,664 bytes,1 万行输出峰值增量 5,226,496 bytes;仍仅为观测而非固定物理低配门。PostgreSQL `18.4` arm64 physical HA 125 项 gate、timeline `1→2` 通过,报告 SHA-256 为 `45fab400eb449774d50429103dd766a2755166530ac54ddc1056f777bc16c15f`,临时 Docker 资源已清理。升级失败自动回滚、真实 controller/RBAC/admission 现场门和固定低配物理证据仍待完成。 - 2026-08-14 收口更新(取代上一句关于 controller 尚未启用的描述):生产 controller 已在既有 `cluster-admin/plugin-package/executor` 内接入 create/get-only Kubernetes Job adapter。它先按确定性名称 GET;仅当 durable execution 仍可创建且 approval 未过期时使用 Strict CREATE;CREATE 的 409 或响应丢失均回到 exact GET 收敛,`executing` 但 Job 缺失、过期 plan、终态 Job 与 renderer contract 漂移全部返回 `recoveryRequired`,不会盲目重建。Controller ServiceAccount 的 RBAC 仅允许 Job `create|get`,action ServiceAccount 无 API token;`admissionregistration.k8s.io/v1` ValidatingAdmissionPolicy 以固定参数 ConfigMap 和请求者身份锁死 digest 镜像、command、两个 ServiceAccount、source Secret、exact SHA-256 item/path、PostgreSQL SecretRef、安全上下文、资源额度和 volume/mount 形状,参数缺失与策略错误均失败关闭。基础 NetworkPolicy 继续只允许 DNS,集群 overlay 必须显式提供 API Server 的精确 CIDR/TCP 443 出口补丁。 - 真实 K3s `v1.34.3+k3s1` 已完成 admission 编译与现场门:合规 Job 的 server dry-run 通过,篡改镜像被策略拒绝,删除参数 ConfigMap 后创建被拒绝;controller SA 的 `list|watch|delete jobs`、Pod 创建和 Secret 读取均被拒绝,action SA 的 Job/Pod 创建与 Secret 读取也均被拒绝。实现仍保持 18 个 package,未新增 workspace package、Edge daemon/timer/watcher 或低配设备常驻负担;短生命周期 controller 与按需 action Job 仅属于 Cluster profile。完整 18-package clean build/test 退出 0;backend 1196 项为 1194 pass、2 条条件 skip、0 fail;cluster-admin 339 pass/3 skip、cluster-postgres 328 pass/2 skip,package boundary、cluster dependency、edge import 和 cluster deployment 审计均无 finding。PostgreSQL `18.4` arm64 physical HA 125 项、timeline `1→2` 通过,报告 SHA-256 为 `a3d34e61ea2064e1cde574e533137186e09fdce9048455da64f582906037fa0d`,临时 Docker 资源已清理。ADR 仍为 Proposed:升级失败自动回滚、终态 Job 的 durable 恢复决议和固定物理低配设备证据尚未完成。 + - 2026-08-14 终态恢复更新:Secret Action controller 不再把所有终态 Job 仅计为瞬时 `recoveryRequired`。Job 到达 Complete/Failed 后已停止执行,controller 会用 started execution 的原 lease fence 复验不可变业务结果:首次 binding 必须与 approval plan、`startedAtMs` 推导出的 binding 完全一致;transition 必须与 plan、authority evidence、commit time 推导出的 receipt 完全一致。精确 durable result 存在时补写 `succeeded`,即使 Job 已被 TTL 清理也能收敛;Failed 且无 durable mutation 时写 `failed`;Complete 但无 receipt 时以 `indeterminate` 写 `blocked`。Job 在 start barrier 前终态或审批过期且尚未创建时,controller 复用既有 claim→release fence 写 `blocked`,不让坏 Job 永久占据 reconciler 页首。任何 stored result 漂移继续抛出 conflict,`executing + Job 缺失 + receipt 缺失` 继续要求人工处理,绝不自动重建可能已产生副作用的动作。该切片不修改共享 execution schema、PostgreSQL migration 或角色权限,不新增 package、连接与常驻进程;Cluster controller 复用现有 package-executor Pool,Edge 零变化。controller/process 定向 21/21,cluster-admin 全包 348 项为 345 pass/3 条件 skip/0 fail;完整 18-package 串行 build/test 退出 0;backend 1196 项为 1194 pass/2 条件 skip/0 fail;package boundary、cluster dependency、edge import、cluster deployment 均无 finding,部署/包边界聚焦测试 61/61。PostgreSQL `18.4` arm64 physical HA 125 项、timeline `1→2` 通过,报告 SHA-256 为 `bec512767fbbd7774baa9366698f60c25c8b017ed66f459b154d143fe86293bc`,临时 Docker 资源已清理。 - D-302/ADR-0390(已接受) Cluster operator context 增加无网络、无 mutation 的内建 `ql3-cluster-admin context validate` 预检。它先复用 owner-private context reader,再让每个 entry 经过与真实请求相同的 production HTTPS/Kubernetes configuration preparation,验证精确 route、hostname、CA、 diff --git a/docs/adr/ADR-0396-generation-transition-plugin-package-secret-binding.md b/docs/adr/ADR-0396-generation-transition-plugin-package-secret-binding.md index a24fd000..5a67d023 100644 --- a/docs/adr/ADR-0396-generation-transition-plugin-package-secret-binding.md +++ b/docs/adr/ADR-0396-generation-transition-plugin-package-secret-binding.md @@ -47,3 +47,7 @@ D-306B1 只允许给当前 active 且尚未绑定的 Package generation 做首 - 真实 K3s `v1.34.3+k3s1` 现场门已证明策略可由 API Server 编译:合规 Job dry-run 被接受,镜像漂移和参数 ConfigMap 缺失被拒绝;controller SA 的 `list|watch|delete jobs`、Pod 创建和 Secret 读取均被拒绝,action SA 的 Job/Pod 创建和 Secret 读取也均被拒绝。实现没有新增 workspace package、Edge daemon/timer/watcher 或低配设备常驻负担,18-package boundary 继续为 `singleSourcePackages=[]`、`shallowSourcePackages=[]`。 - 本切片完整性门:controller/renderer/process 定向 18/18;cluster-admin 339 pass/3 条件 skip、cluster-postgres 328 pass/2 条件 skip;完整 18-package clean build/test 退出 0;backend 1196 项为 1194 pass、2 条条件 skip、0 fail;package boundary、cluster dependency、edge import、cluster deployment 均零 finding。PostgreSQL `18.4` arm64 physical HA 125 项、timeline `1→2` 通过,报告 SHA-256 `a3d34e61ea2064e1cde574e533137186e09fdce9048455da64f582906037fa0d`,临时 Docker 资源已清理。 - ADR 继续保持 Proposed:升级失败自动回滚、终态 Job 的 durable 恢复决议和固定物理低配设备证据尚未完成。当前 controller 明确暴露恢复要求,不以不安全的自动重试冒充闭环。 +- 后续终态恢复切片取代上一条“终态 Job 的 durable 恢复决议尚未完成”的描述。Controller 只在 Job Complete/Failed 后执行恢复,此时 Kubernetes 已证明该 Job 不再运行;它用原 execution lease/start fence 重新推导 expected binding 或 transition receipt,并与 package-executor 只读取得的 immutable durable result 做 exact compare。匹配时补写 `succeeded`;Failed 且无 mutation 时写 `failed`;Complete 且无 receipt 时通过 `indeterminate` 写 `blocked`。Job 已被 TTL 清理但 receipt 存在也可收敛为 succeeded。 +- start barrier 前的终态 Job、以及 approval 过期且 Job 尚未创建,复用共享 execution repository 的 claim→release-before-start 转换持久化为 `blocked`;controller 崩溃后 lease 可回收,不增加新状态、表、migration 或专用恢复 daemon。Durable result 漂移仍是全局 conflict;`executing + Job 缺失 + receipt 缺失` 不能排除孤儿 Pod 或未知副作用,继续保持 `recoveryRequired` 且绝不重建。 +- 本切片定向 controller/process 21/21;cluster-admin 全包 348 项为 345 pass、3 条件 skip、0 fail;完整 18-package 串行 build/test 退出 0;backend 1196 项为 1194 pass、2 条件 skip、0 fail;package boundary、cluster dependency、edge import、cluster deployment 均无 finding,部署/包边界聚焦测试 61/61。PostgreSQL `18.4` arm64 physical HA 125 项、timeline `1→2` 通过,报告 SHA-256 `bec512767fbbd7774baa9366698f60c25c8b017ed66f459b154d143fe86293bc`,临时 Docker 资源已清理。共享 `approved_action_executions` contract、PostgreSQL 权限与 Worker Credential 调用链均未修改;实现继续位于既有 `cluster-admin/plugin-package/executor`,Edge/Standalone 不加载 controller,也没有新增 workspace package、连接、timer、watcher 或常驻内存。 +- ADR 继续保持 Proposed:升级失败自动回滚、`executing + Job/receipt 均缺失` 的显式人工处置产品路径,以及固定物理低配设备证据仍待完成。 diff --git a/packages/ql3-cluster-admin/src/plugin-package/executor/pluginPackageExecutorProcess.ts b/packages/ql3-cluster-admin/src/plugin-package/executor/pluginPackageExecutorProcess.ts index 3dd98c42..2db43815 100644 --- a/packages/ql3-cluster-admin/src/plugin-package/executor/pluginPackageExecutorProcess.ts +++ b/packages/ql3-cluster-admin/src/plugin-package/executor/pluginPackageExecutorProcess.ts @@ -17,7 +17,9 @@ import { loadPostgresConnectionEnvironment, PostgresApprovedActionExecutionRepository, PostgresPluginPackageSecretBindingApprovalPlanReader, + PostgresPluginPackageSecretBindingRepository, PostgresPluginPackageSecretBindingTransitionApprovalPlanReader, + PostgresPluginPackageSecretBindingTransitionRepository, type PostgresConnectionOptions, type PostgresPoolOptions, type PostgresSchemaReadinessReport, @@ -631,6 +633,9 @@ function emptySecretActionJobSummary(): Readonly { findByActionRef(actionRef: string): Promise | null>; } -export interface PluginPackageKubernetesSecretActionExecutionReader { +export interface PluginPackageKubernetesSecretActionExecutionPort { listReconciliableExecutions(query: Readonly<{ nowMs: number; limit: number; @@ -35,6 +52,19 @@ export interface PluginPackageKubernetesSecretActionExecutionReader { executions: readonly Readonly[]; truncated: boolean; }>>; + completeExecution( + command: CompleteApprovedActionExecutionCommand, + ): Promise>; + claimExecution( + command: ClaimApprovedActionExecutionCommand, + ): Promise; + releaseExecutionBeforeStart( + command: ReleaseApprovedActionExecutionBeforeStartCommand, + ): Promise>; +} + +export interface PluginPackageSecretActionDurableResultReader { + find(generationDigest: string): Promise | null>; } export interface PluginPackageKubernetesSecretActionJobResource { @@ -71,9 +101,11 @@ export interface PluginPackageKubernetesSecretActionJobApi { } export interface PluginPackageKubernetesSecretActionControllerOptions { - readonly executions: PluginPackageKubernetesSecretActionExecutionReader; + readonly executions: PluginPackageKubernetesSecretActionExecutionPort; readonly bindingPlans: PluginPackageSecretActionApprovalPlanReader; readonly transitionPlans: PluginPackageSecretActionApprovalPlanReader; + readonly bindings: PluginPackageSecretActionDurableResultReader; + readonly transitionReceipts: PluginPackageSecretActionDurableResultReader; readonly jobs: PluginPackageKubernetesSecretActionJobApi; readonly job: Readonly; readonly now?: () => number; @@ -84,6 +116,9 @@ export interface PluginPackageKubernetesSecretActionControllerSummary { readonly created: number; readonly existing: number; readonly active: number; + readonly recoveredSucceeded: number; + readonly recoveredFailed: number; + readonly recoveredBlocked: number; readonly recoveryRequired: number; readonly unavailable: number; readonly truncated: boolean; @@ -214,10 +249,17 @@ export class PluginPackageKubernetesSecretActionController { typeof options !== 'object' || !options.executions || typeof options.executions.listReconciliableExecutions !== 'function' || + typeof options.executions.completeExecution !== 'function' || + typeof options.executions.claimExecution !== 'function' || + typeof options.executions.releaseExecutionBeforeStart !== 'function' || !options.bindingPlans || typeof options.bindingPlans.findByActionRef !== 'function' || !options.transitionPlans || typeof options.transitionPlans.findByActionRef !== 'function' || + !options.bindings || + typeof options.bindings.find !== 'function' || + !options.transitionReceipts || + typeof options.transitionReceipts.find !== 'function' || !options.jobs || typeof options.jobs.createNamespacedJob !== 'function' || typeof options.jobs.readNamespacedJob !== 'function' || @@ -269,6 +311,9 @@ export class PluginPackageKubernetesSecretActionController { created: 0, existing: 0, active: 0, + recoveredSucceeded: 0, + recoveredFailed: 0, + recoveredBlocked: 0, recoveryRequired: 0, unavailable: 0, truncated: page.truncated, @@ -293,7 +338,15 @@ export class PluginPackageKubernetesSecretActionController { async #reconcileOne( snapshot: Readonly, nowMs: number, - ): Promise<'created' | 'existing' | 'active' | 'recoveryRequired'> { + ): Promise< + | 'created' + | 'existing' + | 'active' + | 'recoveredSucceeded' + | 'recoveredFailed' + | 'recoveredBlocked' + | 'recoveryRequired' + > { const plan = await this.#plan(snapshot); if (!plan) { throw new PluginPackageKubernetesSecretActionControllerUnavailableError(); @@ -313,10 +366,22 @@ export class PluginPackageKubernetesSecretActionController { } catch (error) { if (apiStatus(error) !== 404) throw error; if ( - snapshot.execution.status === 'executing' || - nowMs > plan.expiresAtMs + snapshot.execution.status === 'executing' ) { - return 'recoveryRequired'; + return this.#recoverMissingExecuting( + snapshot, + plan, + name, + nowMs, + ); + } + if (nowMs > plan.expiresAtMs) { + return this.#blockBeforeStart( + snapshot, + name, + nowMs, + 'package_secret_action_approval_expired', + ); } try { observed = await this.options.jobs.createNamespacedJob({ @@ -343,10 +408,213 @@ export class PluginPackageKubernetesSecretActionController { } assertObservedJob(desired, observed); const terminal = terminalStatus(observed); - if (terminal !== 'active') return 'recoveryRequired'; + if (terminal !== 'active') { + return this.#recoverTerminal(snapshot, plan, name, terminal, nowMs); + } return disposition; } + async #recoverTerminal( + snapshot: Readonly, + plan: Readonly, + jobName: string, + terminal: 'complete' | 'failed', + nowMs: number, + ): Promise< + 'recoveredSucceeded' | 'recoveredFailed' | 'recoveredBlocked' | 'recoveryRequired' + > { + const execution = snapshot.execution; + if (execution.status !== 'executing') { + return this.#blockBeforeStart( + snapshot, + jobName, + nowMs, + terminal === 'failed' + ? 'package_secret_action_job_failed_before_start' + : 'package_secret_action_job_completed_before_start', + ); + } + if ( + execution.startedAtMs === null || + execution.leaseOwner === null || + execution.leaseToken === null + ) { + throw new PluginPackageKubernetesSecretActionControllerConflictError(); + } + const durableResultDigest = await this.#durableResultDigest( + plan, + execution.startedAtMs, + ); + const outcome = durableResultDigest + ? 'succeeded' + : terminal === 'failed' + ? 'failed' + : 'indeterminate'; + return this.#completeExecuting( + snapshot, + jobName, + nowMs, + outcome, + durableResultDigest + ? snapshot.dispatch.action.actionType === + PLUGIN_PACKAGE_SECRET_BINDING_ACTION_TYPE + ? 'package_secret_binding_job_recovered' + : 'package_secret_transition_job_recovered' + : terminal === 'failed' + ? 'package_secret_action_job_failed' + : 'package_secret_action_receipt_missing', + durableResultDigest, + ); + } + + async #recoverMissingExecuting( + snapshot: Readonly, + plan: Readonly, + jobName: string, + nowMs: number, + ): Promise<'recoveredSucceeded' | 'recoveryRequired'> { + const execution = snapshot.execution; + if ( + execution.status !== 'executing' || + execution.startedAtMs === null || + execution.leaseOwner === null || + execution.leaseToken === null + ) { + throw new PluginPackageKubernetesSecretActionControllerConflictError(); + } + const durableResultDigest = await this.#durableResultDigest( + plan, + execution.startedAtMs, + ); + if (!durableResultDigest) return 'recoveryRequired'; + const recovered = await this.#completeExecuting( + snapshot, + jobName, + nowMs, + 'succeeded', + snapshot.dispatch.action.actionType === + PLUGIN_PACKAGE_SECRET_BINDING_ACTION_TYPE + ? 'package_secret_binding_job_recovered' + : 'package_secret_transition_job_recovered', + durableResultDigest, + ); + if (recovered !== 'recoveredSucceeded') { + throw new PluginPackageKubernetesSecretActionControllerConflictError(); + } + return recovered; + } + + async #completeExecuting( + snapshot: Readonly, + jobName: string, + nowMs: number, + outcome: 'succeeded' | 'failed' | 'indeterminate', + resultCode: string, + resultDigest: string | null, + ): Promise<'recoveredSucceeded' | 'recoveredFailed' | 'recoveredBlocked'> { + const execution = snapshot.execution; + if ( + execution.status !== 'executing' || + execution.leaseOwner === null || + execution.leaseToken === null || + (outcome === 'succeeded') !== (resultDigest !== null) + ) { + throw new PluginPackageKubernetesSecretActionControllerConflictError(); + } + const result = await this.options.executions.completeExecution({ + dispatchId: execution.dispatchId, + owner: execution.leaseOwner, + leaseToken: execution.leaseToken, + expectedVersion: execution.version, + resultMutationId: `k8s-recovery-${jobName}`, + outcome, + resultCode, + ...(resultDigest ? { resultDigest } : {}), + completedAtMs: Math.max(nowMs, execution.updatedAtMs), + }); + if (result.execution.status === 'succeeded') return 'recoveredSucceeded'; + if (result.execution.status === 'failed') return 'recoveredFailed'; + if (result.execution.status === 'blocked') return 'recoveredBlocked'; + throw new PluginPackageKubernetesSecretActionControllerConflictError(); + } + + async #blockBeforeStart( + snapshot: Readonly, + jobName: string, + nowMs: number, + resultCode: string, + ): Promise<'recoveredBlocked' | 'recoveryRequired'> { + const claimed = await this.options.executions.claimExecution({ + dispatchId: snapshot.execution.dispatchId, + owner: RECOVERY_OWNER, + leaseToken: `recovery-${jobName}`, + nowMs, + leaseDurationMs: RECOVERY_LEASE_DURATION_MS, + }); + if (claimed.status !== 'claimed') { + return claimed.status === 'blocked' ? 'recoveredBlocked' : 'recoveryRequired'; + } + const released = await this.options.executions.releaseExecutionBeforeStart({ + dispatchId: claimed.snapshot.execution.dispatchId, + owner: RECOVERY_OWNER, + leaseToken: claimed.snapshot.execution.leaseToken!, + expectedVersion: claimed.snapshot.execution.version, + resultMutationId: `k8s-recovery-${jobName}`, + resultCode, + atMs: Math.max(nowMs, claimed.snapshot.execution.updatedAtMs), + }); + if (released.execution.status !== 'blocked') { + throw new PluginPackageKubernetesSecretActionControllerConflictError(); + } + return 'recoveredBlocked'; + } + + async #durableResultDigest( + plan: Readonly, + startedAtMs: number, + ): Promise { + if ('bindingPlan' in plan) { + const normalized = normalizePluginPackageSecretBindingApprovalPlan(plan); + const expected = createPluginPackageSecretBindingFromApprovalPlan( + normalized, + startedAtMs, + ); + const stored = await this.options.bindings.find( + normalized.bindingPlan.target.generationDigest, + ); + if (!stored) return null; + if (JSON.stringify(stored) !== JSON.stringify(expected)) { + throw new PluginPackageKubernetesSecretActionControllerConflictError(); + } + return stored.bindingDigest; + } + const normalized = + normalizePluginPackageSecretBindingTransitionApprovalPlan(plan); + const binding = createPluginPackageSecretBindingFromTransitionPlan( + normalized.transitionPlan, + 'approved-action-execution', + normalized.approvalPlanDigest, + startedAtMs, + ); + const expected = createPluginPackageSecretBindingTransitionReceipt({ + transitionPlan: normalized.transitionPlan, + authority: Object.freeze({ + kind: 'approved-action-execution', + evidenceDigest: normalized.approvalPlanDigest, + }), + binding, + committedAtMs: startedAtMs, + }); + const stored = await this.options.transitionReceipts.find( + normalized.transitionPlan.nextTarget.generationDigest, + ); + if (!stored) return null; + if (JSON.stringify(stored) !== JSON.stringify(expected)) { + throw new PluginPackageKubernetesSecretActionControllerConflictError(); + } + return stored.receiptDigest; + } + #plan( snapshot: Readonly, ): Promise | null> { diff --git a/packages/ql3-cluster-admin/test/pluginPackageExecutorProcess.test.cjs b/packages/ql3-cluster-admin/test/pluginPackageExecutorProcess.test.cjs index 0a206fc6..1a334300 100644 --- a/packages/ql3-cluster-admin/test/pluginPackageExecutorProcess.test.cjs +++ b/packages/ql3-cluster-admin/test/pluginPackageExecutorProcess.test.cjs @@ -253,6 +253,9 @@ test('batch mode consumes approvals before reconciling exact Secret action Jobs' }, async createSecretActionController(options) { assert.equal(options.job.image.endsWith('c'.repeat(64)), true); + assert.equal(typeof options.executions.completeExecution, 'function'); + assert.equal(typeof options.bindings.find, 'function'); + assert.equal(typeof options.transitionReceipts.find, 'function'); calls.push('controller-open'); return { controller: { @@ -263,6 +266,9 @@ test('batch mode consumes approvals before reconciling exact Secret action Jobs' created: 1, existing: 0, active: 0, + recoveredSucceeded: 0, + recoveredFailed: 0, + recoveredBlocked: 0, recoveryRequired: 0, unavailable: 0, truncated: false, diff --git a/packages/ql3-cluster-admin/test/pluginPackageKubernetesSecretActionController.test.cjs b/packages/ql3-cluster-admin/test/pluginPackageKubernetesSecretActionController.test.cjs index 5ca63755..1c50c977 100644 --- a/packages/ql3-cluster-admin/test/pluginPackageKubernetesSecretActionController.test.cjs +++ b/packages/ql3-cluster-admin/test/pluginPackageKubernetesSecretActionController.test.cjs @@ -9,24 +9,44 @@ const { decideApprovalRequest, } = require('@qinglong/runtime-core/approved-action'); const { + claimApprovedActionExecution, + completeApprovedActionExecution, createApprovedActionExecution, + releaseApprovedActionExecutionBeforeStart, + startApprovedActionExecution, } = require('@qinglong/runtime-core/approved-action-execution'); const { createPluginPackageResourceGeneration, } = require('@qinglong/runtime-core/plugin-package-resource-generation'); const { + createPluginPackageSecretBindingFromApprovalPlan, createPluginPackageSecretBindingApprovalPlan, } = require('@qinglong/runtime-core/plugin-package-secret-binding-approval-plan'); const { createPluginPackageSecretBindingPlan, } = require('@qinglong/runtime-core/plugin-package-secret-binding-plan'); const { - createSecretRef, -} = require('@qinglong/runtime-core/secret-reference'); + createPluginPackageSecretBinding, +} = require('@qinglong/runtime-core/plugin-package-secret-binding'); +const { + createPluginPackageSecretBindingTransitionApprovalPlan, + pluginPackageSecretBindingTransitionApprovedAction, +} = require('@qinglong/runtime-core/plugin-package-secret-binding-transition-approval-plan'); +const { + createPluginPackageSecretBindingTransitionPlan, +} = require('@qinglong/runtime-core/plugin-package-secret-binding-transition-plan'); +const { + createPluginPackageSecretBindingFromTransitionPlan, + createPluginPackageSecretBindingTransitionReceipt, +} = require('@qinglong/runtime-core/plugin-package-secret-binding-transition-receipt'); +const { createSecretRef } = require('@qinglong/runtime-core/secret-reference'); const { PluginPackageKubernetesSecretActionController, PluginPackageKubernetesSecretActionControllerConflictError, } = require('@qinglong/cluster-admin/plugin-package-kubernetes-secret-action-controller'); +const { + createPluginPackageKubernetesSecretActionJob, +} = require('@qinglong/cluster-admin/plugin-package-kubernetes-secret-action-job'); const REQUESTER = Object.freeze({ type: 'user', id: 'cluster-owner' }); const REVIEWER = Object.freeze({ type: 'user', id: 'security-reviewer' }); @@ -97,8 +117,10 @@ function fixture() { requestedBy: REQUESTER, expiresAtMs: 1_000, }); - const action = require('@qinglong/runtime-core/plugin-package-secret-binding-approval-plan') - .pluginPackageSecretBindingApprovedAction(plan); + const action = + require('@qinglong/runtime-core/plugin-package-secret-binding-approval-plan').pluginPackageSecretBindingApprovedAction( + plan, + ); const approved = decideApprovalRequest( createApprovalRequest({ id: 'approval-controller-1', @@ -167,12 +189,197 @@ function jobOptions() { }; } +function transitionFixture() { + const manifest = (version) => ({ + apiVersion: 'qinglong.io/v1alpha1', + kind: 'Package', + metadata: { + name: 'controller-transition', + displayName: 'Controller Transition', + version, + description: 'Secret action transition controller fixture', + license: 'Apache-2.0', + }, + spec: { + compatibility: { + qinglong: '>=3.0.0-0 <4.0.0', + architectures: ['arm64'], + deploymentProfiles: ['cluster-control'], + }, + runtimes: [], + resources: { + memory: { recommended: '32Mi' }, + disk: { install: '4Mi', working: '8Mi' }, + }, + permissions: { + network: { allowedHosts: [] }, + secrets: [{ name: 'TOKEN', required: true }], + tools: ['secret.use'], + }, + contents: { tasks: [], workflows: [], prompts: [], tools: [] }, + }, + }); + const previousManifest = manifest('1.0.0'); + const previousGeneration = createPluginPackageResourceGeneration({ + installationId: 'install-controller-transition-v1', + projectId: 'project-1', + packageName: 'controller-transition', + lockDigest: '1'.repeat(64), + generation: 1, + previousActiveLockDigest: null, + contentDigest: '2'.repeat(64), + contents: previousManifest.spec.contents, + }); + const previousBinding = createPluginPackageSecretBinding({ + generation: previousGeneration, + manifest: previousManifest, + assignments: [ + { + name: 'TOKEN', + secretRef: createSecretRef({ + projectId: 'project-1', + name: 'runtime-token', + version: 1, + }), + }, + ], + authority: { + kind: 'approved-action-execution', + evidenceDigest: '3'.repeat(64), + }, + boundAtMs: 80, + }); + const nextManifest = manifest('2.0.0'); + const transitionPlan = createPluginPackageSecretBindingTransitionPlan({ + previousTarget: previousBinding.target, + previousBinding, + previousAttemptGeneration: 1, + nextGeneration: createPluginPackageResourceGeneration({ + installationId: 'install-controller-transition-v2', + projectId: 'project-1', + packageName: 'controller-transition', + lockDigest: '4'.repeat(64), + generation: 2, + previousActiveLockDigest: '1'.repeat(64), + contentDigest: '5'.repeat(64), + contents: nextManifest.spec.contents, + }), + nextManifest, + assignments: [ + { + name: 'TOKEN', + secretRef: createSecretRef({ + projectId: 'project-1', + name: 'runtime-token', + version: 2, + }), + }, + ], + plannedAtMs: 100, + }); + const plan = createPluginPackageSecretBindingTransitionApprovalPlan({ + actionRef: 'secret-transition:controller-v2', + transitionPlan, + requestedBy: REQUESTER, + plannedAtMs: 100, + expiresAtMs: 1_000, + }); + const action = pluginPackageSecretBindingTransitionApprovedAction(plan); + const approved = decideApprovalRequest( + createApprovalRequest({ + id: 'approval-controller-transition', + projectId: 'project-1', + action, + risk: 'high', + decisionMode: 'separation_of_duty', + requestedBy: REQUESTER, + requestedAtMs: 110, + expiresAtMs: 900, + requestFence: FENCE, + }), + { + expectedVersion: 1, + decisionId: 'decision-controller-transition', + decision: 'approved', + reasonCode: 'reviewed', + principal: { + subject: REVIEWER, + authenticationId: 'auth-transition-reviewer', + authenticatedAtMs: 100, + expiresAtMs: 800, + assurance: 'multi_factor', + }, + decidedAtMs: 120, + authorizationFence: FENCE, + }, + ); + const dispatch = consumeApprovalRequest(approved, { + expectedVersion: 2, + consumptionId: 'consume-controller-transition', + dispatchId: 'dispatch-controller-transition', + action, + requestedBy: REQUESTER, + consumedBy: CONSUMER, + consumedAtMs: 130, + authorizationFence: FENCE, + }).dispatch; + return { + plan, + snapshot: Object.freeze({ + dispatch, + execution: createApprovedActionExecution(dispatch), + }), + }; +} + function apiError(code) { return Object.assign(new Error(`Kubernetes ${code}`), { code }); } -function controller({ read, create, now = 200, snapshot: snapshotOverride }) { +function executingSnapshot(snapshot, startedAtMs = 150) { + const leased = claimApprovedActionExecution(snapshot.execution, { + owner: 'secret-action-job', + leaseToken: 'secret-action-lease', + nowMs: 140, + leaseDurationMs: 100, + }); + return Object.freeze({ + dispatch: snapshot.dispatch, + execution: startApprovedActionExecution( + { dispatch: snapshot.dispatch, execution: leased }, + { + dispatchId: snapshot.dispatch.id, + approvalRequestId: snapshot.dispatch.approvalRequestId, + actionDigest: snapshot.dispatch.action.actionDigest, + owner: leased.leaseOwner, + leaseToken: leased.leaseToken, + expectedVersion: leased.version, + startedAtMs, + }, + ), + }); +} + +function desiredJob(plan, snapshot) { + return createPluginPackageKubernetesSecretActionJob({ + dispatch: snapshot.dispatch, + approvalPlan: plan, + options: jobOptions(), + }); +} + +function controller({ + read, + create, + complete, + findBinding, + findTransitionReceipt, + now = 200, + plan: planOverride, + snapshot: snapshotOverride, +}) { const { plan, snapshot } = fixture(); + const selectedPlan = planOverride ?? plan; return new PluginPackageKubernetesSecretActionController({ executions: { async listReconciliableExecutions(query) { @@ -185,16 +392,72 @@ function controller({ read, create, now = 200, snapshot: snapshotOverride }) { truncated: false, }; }, + async completeExecution(command) { + if (!complete) throw new Error('completion must not run'); + return complete(command, snapshotOverride ?? snapshot); + }, + async claimExecution(command) { + const current = snapshotOverride ?? snapshot; + const execution = claimApprovedActionExecution(current.execution, { + owner: command.owner, + leaseToken: command.leaseToken, + nowMs: command.nowMs, + leaseDurationMs: command.leaseDurationMs, + }); + return { + status: 'claimed', + snapshot: { dispatch: current.dispatch, execution }, + }; + }, + async releaseExecutionBeforeStart(command) { + const current = snapshotOverride ?? snapshot; + const claimed = claimApprovedActionExecution(current.execution, { + owner: command.owner, + leaseToken: command.leaseToken, + nowMs: command.atMs, + leaseDurationMs: 60_000, + }); + return { + dispatch: current.dispatch, + execution: releaseApprovedActionExecutionBeforeStart(claimed, { + owner: command.owner, + leaseToken: command.leaseToken, + expectedVersion: command.expectedVersion, + resultMutationId: command.resultMutationId, + resultCode: command.resultCode, + atMs: command.atMs, + }), + }; + }, }, bindingPlans: { async findByActionRef(actionRef) { - assert.equal(actionRef, plan.actionRef); - return plan; + if (!('bindingPlan' in selectedPlan)) { + throw new Error('binding reader must not run'); + } + assert.equal(actionRef, selectedPlan.actionRef); + return selectedPlan; }, }, transitionPlans: { - async findByActionRef() { - throw new Error('transition reader must not run'); + async findByActionRef(actionRef) { + if ('bindingPlan' in selectedPlan) { + throw new Error('transition reader must not run'); + } + assert.equal(actionRef, selectedPlan.actionRef); + return selectedPlan; + }, + }, + bindings: { + async find(generationDigest) { + return findBinding ? findBinding(generationDigest) : null; + }, + }, + transitionReceipts: { + async find(generationDigest) { + return findTransitionReceipt + ? findTransitionReceipt(generationDigest) + : null; }, }, jobs: { @@ -228,7 +491,10 @@ test('creates one Strict deterministic Job without claiming the execution', asyn calls[1][1].fieldManager, 'qinglong-plugin-package-secret-action-controller', ); - assert.match(calls[1][1].body.metadata.name, /^ql3-package-secret-[0-9a-f]{32}$/); + assert.match( + calls[1][1].body.metadata.name, + /^ql3-package-secret-[0-9a-f]{32}$/, + ); }); test('converges a concurrent create through one exact get', async () => { @@ -284,10 +550,7 @@ test('does not recreate a missing Job for an already executing action', async () const { snapshot } = fixture(); let creates = 0; const subject = controller({ - snapshot: { - ...snapshot, - execution: { ...snapshot.execution, status: 'executing' }, - }, + snapshot: executingSnapshot(snapshot), async read() { throw apiError(404); }, @@ -301,7 +564,47 @@ test('does not recreate a missing Job for an already executing action', async () assert.equal(creates, 0); }); -test('does not create a missing Job after the approval plan expires', async () => { +test('recovers a missing executing Job when its exact durable binding exists', async () => { + const { plan, snapshot } = fixture(); + const executing = executingSnapshot(snapshot); + const durable = createPluginPackageSecretBindingFromApprovalPlan( + plan, + executing.execution.startedAtMs, + ); + const subject = controller({ + snapshot: executing, + async read() { + throw apiError(404); + }, + async create() { + throw new Error('must not create'); + }, + async findBinding() { + return durable; + }, + async complete(command) { + assert.equal(command.outcome, 'succeeded'); + return { + dispatch: executing.dispatch, + execution: completeApprovedActionExecution(executing.execution, { + owner: command.owner, + leaseToken: command.leaseToken, + expectedVersion: command.expectedVersion, + resultMutationId: command.resultMutationId, + outcome: command.outcome, + resultCode: command.resultCode, + resultDigest: command.resultDigest, + completedAtMs: command.completedAtMs, + }), + }; + }, + }); + const result = await subject.reconcile(); + assert.equal(result.recoveredSucceeded, 1); + assert.equal(result.recoveryRequired, 0); +}); + +test('blocks a missing Job after the approval plan expires before start', async () => { let creates = 0; const subject = controller({ now: 1_001, @@ -314,11 +617,12 @@ test('does not create a missing Job after the approval plan expires', async () = }, }); const result = await subject.reconcile(); - assert.equal(result.recoveryRequired, 1); + assert.equal(result.recoveredBlocked, 1); + assert.equal(result.recoveryRequired, 0); assert.equal(creates, 0); }); -test('marks a terminal Job with a nonterminal execution as recovery required', async () => { +test('blocks a terminal Job before the execution start barrier', async () => { let desired; const subject = controller({ async read() { @@ -336,10 +640,222 @@ test('marks a terminal Job with a nonterminal execution as recovery required', a }, }); const result = await subject.reconcile(); - assert.equal(result.recoveryRequired, 1); + assert.equal(result.recoveredBlocked, 1); + assert.equal(result.recoveryRequired, 0); assert.equal(result.created, 0); }); +test('recovers an executing terminal Job as succeeded from its exact durable binding', async () => { + const { plan, snapshot } = fixture(); + const executing = executingSnapshot(snapshot); + const desired = desiredJob(plan, executing); + const durable = createPluginPackageSecretBindingFromApprovalPlan( + plan, + executing.execution.startedAtMs, + ); + let completion; + const subject = controller({ + snapshot: executing, + async read() { + return { + ...desired, + status: { conditions: [{ type: 'Complete', status: 'True' }] }, + }; + }, + async create() { + throw new Error('must not create'); + }, + async findBinding(generationDigest) { + assert.equal(generationDigest, plan.bindingPlan.target.generationDigest); + return durable; + }, + async complete(command) { + completion = command; + return { + dispatch: executing.dispatch, + execution: completeApprovedActionExecution(executing.execution, { + owner: command.owner, + leaseToken: command.leaseToken, + expectedVersion: command.expectedVersion, + resultMutationId: command.resultMutationId, + outcome: command.outcome, + resultCode: command.resultCode, + resultDigest: command.resultDigest, + completedAtMs: command.completedAtMs, + }), + }; + }, + }); + const result = await subject.reconcile(); + assert.equal(result.recoveredSucceeded, 1); + assert.equal(result.recoveryRequired, 0); + assert.equal(completion.outcome, 'succeeded'); + assert.equal(completion.resultDigest, durable.bindingDigest); + assert.equal(completion.resultCode, 'package_secret_binding_job_recovered'); +}); + +test('recovers a failed executing Job without a durable mutation as failed', async () => { + const { plan, snapshot } = fixture(); + const executing = executingSnapshot(snapshot); + const desired = desiredJob(plan, executing); + const subject = controller({ + snapshot: executing, + async read() { + return { + ...desired, + status: { conditions: [{ type: 'Failed', status: 'True' }] }, + }; + }, + async create() { + throw new Error('must not create'); + }, + async complete(command) { + assert.equal(command.outcome, 'failed'); + return { + dispatch: executing.dispatch, + execution: completeApprovedActionExecution(executing.execution, { + owner: command.owner, + leaseToken: command.leaseToken, + expectedVersion: command.expectedVersion, + resultMutationId: command.resultMutationId, + outcome: command.outcome, + resultCode: command.resultCode, + completedAtMs: command.completedAtMs, + }), + }; + }, + }); + const result = await subject.reconcile(); + assert.equal(result.recoveredFailed, 1); + assert.equal(result.recoveryRequired, 0); +}); + +test('recovers a terminal transition Job from its exact durable receipt', async () => { + const { plan, snapshot } = transitionFixture(); + const executing = executingSnapshot(snapshot); + const desired = desiredJob(plan, executing); + const binding = createPluginPackageSecretBindingFromTransitionPlan( + plan.transitionPlan, + 'approved-action-execution', + plan.approvalPlanDigest, + executing.execution.startedAtMs, + ); + const receipt = createPluginPackageSecretBindingTransitionReceipt({ + transitionPlan: plan.transitionPlan, + authority: { + kind: 'approved-action-execution', + evidenceDigest: plan.approvalPlanDigest, + }, + binding, + committedAtMs: executing.execution.startedAtMs, + }); + const subject = controller({ + plan, + snapshot: executing, + async read() { + return { + ...desired, + status: { conditions: [{ type: 'Failed', status: 'True' }] }, + }; + }, + async create() { + throw new Error('must not create'); + }, + async findTransitionReceipt(generationDigest) { + assert.equal( + generationDigest, + plan.transitionPlan.nextTarget.generationDigest, + ); + return receipt; + }, + async complete(command) { + assert.equal(command.outcome, 'succeeded'); + assert.equal(command.resultDigest, receipt.receiptDigest); + return { + dispatch: executing.dispatch, + execution: completeApprovedActionExecution(executing.execution, { + owner: command.owner, + leaseToken: command.leaseToken, + expectedVersion: command.expectedVersion, + resultMutationId: command.resultMutationId, + outcome: command.outcome, + resultCode: command.resultCode, + resultDigest: command.resultDigest, + completedAtMs: command.completedAtMs, + }), + }; + }, + }); + const result = await subject.reconcile(); + assert.equal(result.recoveredSucceeded, 1); + assert.equal(result.recoveryRequired, 0); +}); + +test('blocks a completed executing Job when its durable receipt is missing', async () => { + const { plan, snapshot } = fixture(); + const executing = executingSnapshot(snapshot); + const desired = desiredJob(plan, executing); + const subject = controller({ + snapshot: executing, + async read() { + return { + ...desired, + status: { conditions: [{ type: 'Complete', status: 'True' }] }, + }; + }, + async create() { + throw new Error('must not create'); + }, + async complete(command) { + assert.equal(command.outcome, 'indeterminate'); + return { + dispatch: executing.dispatch, + execution: completeApprovedActionExecution(executing.execution, { + owner: command.owner, + leaseToken: command.leaseToken, + expectedVersion: command.expectedVersion, + resultMutationId: command.resultMutationId, + outcome: command.outcome, + resultCode: command.resultCode, + completedAtMs: command.completedAtMs, + }), + }; + }, + }); + const result = await subject.reconcile(); + assert.equal(result.recoveredBlocked, 1); + assert.equal(result.recoveryRequired, 0); +}); + +test('fails closed when the durable binding differs from the approved result', async () => { + const { plan, snapshot } = fixture(); + const executing = executingSnapshot(snapshot); + const desired = desiredJob(plan, executing); + const durable = createPluginPackageSecretBindingFromApprovalPlan( + plan, + executing.execution.startedAtMs, + ); + const subject = controller({ + snapshot: executing, + async read() { + return { + ...desired, + status: { conditions: [{ type: 'Complete', status: 'True' }] }, + }; + }, + async create() { + throw new Error('must not create'); + }, + async findBinding() { + return { ...durable, bindingDigest: 'f'.repeat(64) }; + }, + }); + await assert.rejects( + () => subject.reconcile(), + PluginPackageKubernetesSecretActionControllerConflictError, + ); +}); + test('fails closed when a deterministic Job name contains another contract', async () => { const subject = controller({ async read() {