diff --git a/docs/QINGLONG_3_0_ARCHITECTURE_RFC.md b/docs/QINGLONG_3_0_ARCHITECTURE_RFC.md index db0b8ab2..b35d6f1d 100644 --- a/docs/QINGLONG_3_0_ARCHITECTURE_RFC.md +++ b/docs/QINGLONG_3_0_ARCHITECTURE_RFC.md @@ -11,7 +11,23 @@ 最新增量证据(2026-08-19): -- D-367/ADR-0460(已接受;有界 blocked drill-down 与 Console 可选接入待完成):在现有 `ql3 run`/`ql3-run-client` 增加一次性 +- D-368/ADR-0461(已接受;Console 可选接入待完成):在既有 Run management plane 与 `ql3 run` 增加 + `run.cancellation.blocked.list`/`blocked` 一次性 drill-down。服务端固定 16 项、只查询 17 行,客户端不能提供 limit 或自动翻页;页面只含 + `{runId,blockedAtMs}`,以 PostgreSQL 首屏时间和 `(blockedAtMs,runId)` 组成稳定 cursor,严格 oldest-first、快照后新增项不进入后续页。强认证 User 继续使用 + `run.read`,Policy fence、数据库时间、键集读取与 allowed audit 在同一 5 秒 SERIALIZABLE 短事务;CLI cursor 是有版本、精确字段、大小受限的 canonical + base64url token,非法输入在网络前失败。`pg-0068-cancellation-dispatch-project-keyset`/capability v67 为 dispatch 增加一次性回填的 `project_id`、 + `runs(project_id,id)` 唯一键、复合外键和仅 blocked 的 `(project_id,updated_at_ms,run_id)` partial index,避免跨租户扫描;runtime claim 从已锁定 Run + 原子写入 Project。没有新增 package、依赖、binary、服务、端口、timer、queue、cache 或 Kubernetes 对象;代码只进入既有领域子目录,`src` 根仍只保留入口。 + Edge/Standalone 闭包不包含 Cluster 管理能力,低配路由器维持零新增常驻开销。Cluster Admin 全量 `413 total / 410 pass / 3 conditional skip / 0 fail`, + Cluster PostgreSQL 全量 `347 total / 344 pass / 3 conditional skip / 0 fail`,backend 为 `1,489 total / 1,487 pass / 2 conditional skip / 0 fail`, + 18-package clean build/顺序测试单次退出 0。package layout 聚焦审计 `10/10`,四项架构审计全部 compatible;workspace 保持 18 packages、无 + single/shallow package,Cluster Admin 为 `124 source / 123 nested`,Cluster PostgreSQL 为 `173 source / 172 nested`。`14/14` Local artifact audit + 全部 compatible;基础 Edge/Standalone 保持 `2,589,998 / 2,590,076` bytes,Application+AI 为 `4,493,151 / 4,493,283` bytes,MCP 为 + `7,315,930 / 7,316,038` bytes。PostgreSQL 18.6 arm64 HA `146/146`、timeline `1→2`,真实证明 v67 migration、Project partial-index list、同事务 audit、 + rearm、生产交付、WAL 与 promotion;报告 SHA-256 为 `1fbd58c5bb32bbf83b6c1970a594f7879c33d63057a3b9c34f13ec9917ff5c44`,独立 evidence audit + compatible 且零 finding。 + +- D-367/ADR-0460(已接受;有界 blocked drill-down 由 D-368 完成,Console 可选接入待完成):在现有 `ql3 run`/`ql3-run-client` 增加一次性 `status --config=... --assertion=... --project=... [--format=text|json]` 产品入口。它在内存生成固定 `run.cancellation.summary` 命令,仍经同一 exact codec、TLS 1.3、Run 专用 mTLS/OIDC、固定 management route 和响应交叉不变量校验;原 `--command` 私有文件模式保持兼容。默认 text 是无 ANSI 的确定性 Project 状态卡,JSON 使用 `qinglong/run-cancellation-status@v1`;两者只含 D-366 低敏计数与结论。告警映射固定为 diff --git a/docs/adr/ADR-0005-durable-cancellation-dispatch.md b/docs/adr/ADR-0005-durable-cancellation-dispatch.md index 3546836f..102d38f2 100644 --- a/docs/adr/ADR-0005-durable-cancellation-dispatch.md +++ b/docs/adr/ADR-0005-durable-cancellation-dispatch.md @@ -157,9 +157,11 @@ ADR-0459 进一步增加 Project-scoped、caller-driven summary,以固定计 ADR-0460 已在现有 `ql3 run` 增加 one-shot status 产品入口:同一强认证 summary 被投影为低敏 text/JSON 状态卡,并以 `0/10/20` 区分 clear、converging 与 attention-required;没有新增轮询器、临时 command 文件或常驻 authority。 +ADR-0461 已增加 Project-scoped blocked drill-down:固定 16 项、只返回 Run ID/blocked time,以数据库快照和 `(blockedAtMs,runId)` 键集显式逐页;v67 将 Project identity 固化到 dispatch 并用复合外键与 partial index 避免跨租户扫描。该入口仍是 one-shot,不自动翻页或引入后台 scanner。 + HTTP worker 已通过默认关闭的 manual-only manifest bootstrap 接入 Local Supervisor:只有 accepted 且全部 gate 通过时才启动,失败或 shutdown 时有界停止。以下工作仍未完成,因此它仍只允许显式 canary,不得扩大到默认生产流量: -- Copilot Console 的可选 Project 状态卡与有界 blocked drill-down;Project 聚合出口、一次性产品状态卡、告警退出码和私有 operator 处置协议已经实现。 +- Copilot Console 的可选 Project 状态卡与用户触发式 blocked drill-down;Project 聚合出口、一次性产品状态卡、告警退出码、有界发现和私有 operator 处置协议已经实现。 - 固定 edge 设备的数据库写放大、RSS、时延和磁盘基准。 - 固定 Project allowlist 的外部时序指标适配;数据库事实驱动的按需 availability/blocked 汇总已经实现。 - 首次真实目标实例完整激活/回滚仪式与共享 config 多写者 authority。 @@ -183,3 +185,4 @@ HTTP worker 已通过默认关闭的 manual-only manifest bootstrap 接入 Local 15. rearm、RunEvent 与 allowed audit 原子提交,数据库时间决定 retry due;stale fence、授权漂移与 mutation drift 均失败关闭。 16. Project summary 的五态与 blocking-result 计数必须交叉守恒,blocked 只产生 `attention_required` 告警而不撤回全局 readiness,且响应不包含 Run/Attempt/lease identity。 17. `ql3 run status` 必须只发送一次固定 summary,text/JSON 事实一致,`0/10/20` 与三态严格映射;原 command-file 模式和 Edge/Standalone 闭包不得变化。 +18. `ql3 run blocked` 每次只能读取一个固定 16 项页面;cursor 必须保持首屏数据库快照和严格键集顺序,响应只含 Run identity/blocked time,并由 Project 复合外键与 partial index 证明无跨租户扫描。 diff --git a/docs/adr/ADR-0459-project-scoped-cancellation-availability-summary.md b/docs/adr/ADR-0459-project-scoped-cancellation-availability-summary.md index 80c9c025..99542f70 100644 --- a/docs/adr/ADR-0459-project-scoped-cancellation-availability-summary.md +++ b/docs/adr/ADR-0459-project-scoped-cancellation-availability-summary.md @@ -58,4 +58,4 @@ QingLong 同时面对小型路由设备和多副本 Cluster。Local/Edge 不应 ## 后续 -ADR-0460 已用现有 `ql3 run status` 完成一次性 Project 状态卡、稳定 JSON 和 `0/10/20` 告警退出码,不增加常驻组件或混合 Copilot Console authority。未知 blocked Run 的有界发现与 Console 可选接入仍需独立设计。CloudNativePG live failover、多副本容量压力、固定 Linux x64/arm64 和物理 Edge 资源证据仍是独立发布门;本 ADR 不把 Docker HA 或按需汇总冒充这些现场证据。 +ADR-0460 已用现有 `ql3 run status` 完成一次性 Project 状态卡、稳定 JSON 和 `0/10/20` 告警退出码;ADR-0461 又以固定 16 项、快照键集和 Project partial index 完成未知 blocked Run 的有界发现。两者都不增加常驻组件或混合 Copilot Console authority。Console 可选接入、CloudNativePG live failover、多副本容量压力、固定 Linux x64/arm64 和物理 Edge 资源证据仍是独立发布门;本 ADR 不把 Docker HA 或按需汇总冒充这些现场证据。 diff --git a/docs/adr/ADR-0460-one-shot-cancellation-status-product-entry.md b/docs/adr/ADR-0460-one-shot-cancellation-status-product-entry.md index ac058d82..6f569432 100644 --- a/docs/adr/ADR-0460-one-shot-cancellation-status-product-entry.md +++ b/docs/adr/ADR-0460-one-shot-cancellation-status-product-entry.md @@ -59,4 +59,4 @@ QingLong 的部署跨度很大。Cluster operator 需要可读状态卡和可供 ## 后续 -D-368 可设计 Project-scoped blocked drill-down:固定小页、稳定数据库 cursor、只返回继续 inspect 所需的最低 identity,并证明索引、事务一致性、Policy/audit 和多副本 HA。Copilot Console 状态卡应复用本 ADR 的 projection 和显式 Run management authority,但不得默认持有该 authority 或建立轮询。 +ADR-0461/D-368 已完成 Project-scoped blocked drill-down:固定 16 项、数据库快照键集、最低 Run identity、Project partial index、Policy/audit 与一次性 `ql3 run blocked`。Copilot Console 状态卡和 drill-down 应复用这些 projection 与显式 Run management authority,但不得默认持有该 authority、建立轮询或自动翻页。 diff --git a/docs/adr/ADR-0461-project-scoped-blocked-cancellation-keyset.md b/docs/adr/ADR-0461-project-scoped-blocked-cancellation-keyset.md new file mode 100644 index 00000000..92052489 --- /dev/null +++ b/docs/adr/ADR-0461-project-scoped-blocked-cancellation-keyset.md @@ -0,0 +1,69 @@ +# ADR-0461:Project-scoped Blocked Cancellation 固定键集分页 + +- 状态:Accepted +- 日期:2026-08-19 +- 关联 RFC:QL-RFC-0001 D-368、PR-5、PR-7 +- 关联 ADR:ADR-0005、ADR-0458、ADR-0459、ADR-0460 +- Amends:ADR-0460 的未知 blocked Run drill-down 边界 + +## 上下文 + +`ql3 run status` 已能低成本判断一个 Project 是否存在 blocked cancellation,但 operator 若事先不知道 Run ID,仍无法进入 inspect/rearm。把 Run 列表直接并入 summary 会让固定低基数状态卡变成分页协议,也会把调用次数、状态筛选和数据量控制权交给客户端。 + +QingLong 必须同时覆盖低配路由设备和集群节点。该能力只属于显式安装的 Cluster 管理面;不能为 Edge/Standalone 增加 PostgreSQL、管理凭据、后台扫描器或新的包。Cluster 侧则必须避免跨租户扫描,并在并发 rearm、failover 和翻页期间保持可解释的快照边界。 + +## 决策 + +1. 在既有 Run management 协议增加 `run.cancellation.blocked.list`,并在既有 `ql3 run` 增加一次性 `blocked --config=... --assertion=... --project=... [--cursor=...] [--format=text|json]`。不新增 package、binary、服务、端口、timer、queue、cache、连接池或 Kubernetes 对象。 +2. 服务端页大小固定为 16,内部只读取 `limit + 1`。客户端不能提交 limit、状态、排序、时间窗口或任意路径,也不自动翻页;每次命令只发送一个请求并退出。继续读取必须由 operator 显式传回上一页的不透明 cursor。 +3. 页面只返回 `{runId, blockedAtMs}`,按 `(blockedAtMs, runId)` 严格升序;不返回 Attempt、blocking result、dispatch version/count、lease owner/token/digest、PID、Worker、命令、环境、Secret、日志或错误原文。具体诊断继续走既有单 Run inspect,处置继续走 exact-CAS rearm。 +4. 首页以 PostgreSQL `transaction_timestamp()` 固定 `snapshotAtMs`;后续 cursor 精确包含 `{snapshotAtMs, blockedAtMs, runId}`。所有页都要求 `updated_at_ms <= snapshotAtMs` 并从上一键之后继续,因此翻页不会吸收快照之后新 blocked 的记录。已被 rearm 的行可以从后续页消失;该列表是有界运维发现视图,不是历史审计账本。 +5. cursor 在产品 CLI 中编码为版本化 canonical base64url token,大小和字段精确受限。服务端、transport、client 和 CLI 分别验证 exact shape、整数范围、Project/request 绑定、严格排序、快照上限、16 项上限、`truncated/nextCursor` 一致性;非法 token 在网络 I/O 前以 usage 64 失败关闭。 +6. 读取要求强认证 User 与 `run.read`,复用 Run 专用 mTLS/OIDC 和固定 management route。Policy fence 确认、数据库时间、键集读取及 allowed audit 位于同一个最长 5 秒的 SERIALIZABLE 短事务;denied audit 保留既有事务外低敏失败路径。 +7. `pg-0068-cancellation-dispatch-project-keyset` 把 capability 提升至 v67:为 dispatch 持久化 `project_id`,由既有 Run 关系一次性回填并设为非空;增加 `runs(project_id,id)` 唯一键、dispatch `(project_id,run_id)` 复合外键,以及仅覆盖 blocked 行的 `(project_id,updated_at_ms,run_id)` partial index。运行时 claim 从已锁定 Run 取得 Project,并与 dispatch 原子写入,禁止跨 Project 身份漂移。 +8. migration 的一次性回填和索引建立属于 3.0 孵化期 schema 修正。升级前必须按现有 migration ceremony 评估表大小和锁窗口;运行时查询不得在缺少 v67 capability 时降级为跨租户扫描。 +9. 实现继续内聚在现有 `cluster-postgres/run-management`、`cluster-postgres/run`、`cluster-admin/run-management` 子域。`src` 根只保留受审入口;不为一个分页查询拆单文件微包,也不把实现重新平铺到 package 根。 +10. Edge/Standalone 及其 Application、AI、MCP 制品继续禁止依赖 `cluster-admin`、`cluster-postgres`、`pg` 或该命令。低配路由器默认零新增进程、连接、timer、磁盘 schema 和安装字节;只有 Cluster operator 的显式调用承担一次短事务和一次短 TLS 连接。 + +## 被拒绝的替代方案 + +### 在 summary 中附带前 N 个 Run + +拒绝。summary 的固定计数契约会被分页状态污染,告警调用也会无条件读取 identity;N 之外仍无法处置,而且无法表达稳定 continuation。 + +### 允许客户端选择 limit、排序或 blocking result + +拒绝。它扩大查询形态、索引组合和响应预算,并可能被用作高基数枚举。固定 16 项足以驱动人工 drill-down,规模化消费应另行设计受控导出。 + +### 只在 dispatch 上按 status/updated_at 建全局索引 + +拒绝。Project 查询必须先跨租户读取候选再 join/filter,既浪费资源也削弱租户边界的数据库证明。Project identity 必须进入 dispatch durable row、外键和 partial index。 + +### offset 分页或客户端自动翻到结束 + +拒绝。offset 在并发 rearm 下会跳项/重复且成本随页数增长;自动翻页会把一次性产品命令变成隐藏的无界循环。快照键集和显式逐页调用保持成本可见且有上限。 + +### 为列表新增 scanner、缓存或常驻 Console authority + +拒绝。durable dispatch 已是数据库事实源;第二份缓存会在重启/failover 后漂移。常驻高权限 authority 与后台扫描也不符合小型 Cluster 和可选 Console 的资源/安全边界。 + +## 资源、安全与部署影响 + +- 每次调用最多读取 17 个 partial-index entry、返回 16 个低敏 identity,并写一个安全审计事件;没有 N+1、后台 cadence 或隐藏重试。 +- v67 migration 会对已有 dispatch 行做一次 Project 回填并建立新索引;这是明确的升级维护窗口成本,不是运行时常驻成本。 +- rearm 可使同一快照中的尚未读取项消失,因此 cursor 保证“不会读入快照后新增项”和“顺序单调”,不承诺历史集合冻结。需要历史证明时使用 Security Audit/RunEvent,而不是列表页面。 +- 文本与 JSON 共用已验证 projection;cursor 不含凭据或 capability,但仍按不透明 continuation 处理,不写入 operator context。 +- workspace 维持 18 个职责包;新增源文件全部位于已有领域目录,package boundary 审计继续禁止不合理单文件包、浅层实现和 `src` 根实现增长。 + +## 验证 + +- migration/schema lockstep、repository、service、transport、client、产品 cursor/card 与真实 TLS CLI 测试覆盖 v67 checksum、Project 复合外键、partial index、固定 16+1、snapshot continuation、原子 audit、viewer `run.read`、低敏页面和网络前 cursor 拒绝。 +- Cluster Admin 全量 `413 total / 410 pass / 3 conditional skip / 0 fail`;Cluster PostgreSQL 全量 `347 total / 344 pass / 3 conditional skip / 0 fail`;package layout 聚焦审计 `10/10`。 +- 完整 backend `1,489 total / 1,487 pass / 2 conditional skip / 0 fail`;18-package clean build/顺序测试单次退出 0。依赖审计曾发现 Admin client 为页大小常量导入 PostgreSQL run-manager authority,最终把 `16` 收敛为 runtime-core profile-neutral protocol constant,不放宽审计;修复后相关三包重建与 26 项聚焦门通过。 +- package boundary、Cluster dependency、Edge import、Cluster deployment 四项审计全部 compatible;workspace 维持 18 包且 `singleSourcePackages=[]`、`shallowSourcePackages=[]`,Cluster Admin 为 `124 source / 123 nested`,Cluster PostgreSQL 为 `173 source / 172 nested`。 +- `14/14` Local artifact audit 全部 compatible;基础 Edge/Standalone 为 `2,589,998 / 2,590,076` bytes,Application+AI 为 `4,493,151 / 4,493,283` bytes,MCP 为 `7,315,930 / 7,316,038` bytes,证明 Cluster-only 分页和 v67 schema 未进入低配设备闭包。 +- PostgreSQL 18.6 arm64 HA `146/146`、timeline `1→2`;真实 migration、blocked list、Project partial-index plan、同事务 allowed audit、rearm、production delivery、WAL replay 和 promotion 后读取全部通过。私有报告 SHA-256 为 `1fbd58c5bb32bbf83b6c1970a594f7879c33d63057a3b9c34f13ec9917ff5c44`,独立 evidence audit compatible 且零 finding。 + +## 后续 + +D-369 可设计 Copilot Console 的显式可选 Run management authority 与用户触发式 status→blocked→inspect 导航,但不能默认持有管理凭据、后台轮询、自动翻页或把 rearm 变成只读 Console 能力。CloudNativePG live failover、固定 Linux x64/arm64 容量和物理 Edge 资源证据仍是独立发布门。 diff --git a/docs/adr/README.md b/docs/adr/README.md index de7d1cf7..60864019 100644 --- a/docs/adr/README.md +++ b/docs/adr/README.md @@ -459,6 +459,12 @@ | [ADR-0453](./ADR-0453-origin-scoped-legacy-shadow-capture-authority-and-primary-gate.md) | Origin-scoped Legacy Shadow 捕获权威与 Primary 门禁 | Accepted(首次真实目标实例 manual canary 待执行) | | [ADR-0454](./ADR-0454-target-instance-manual-primary-canary-ceremony.md) | 目标实例 Manual Primary Canary 与显式回滚仪式 | Accepted(首次真实用户目标实例执行待运维) | | [ADR-0455](./ADR-0455-profile-bounded-manual-primary-runtime-activation-receipt.md) | Profile-bounded Manual Primary 运行态激活凭据 | Accepted(首次真实目标实例执行待运维) | +| [ADR-0456](./ADR-0456-database-timed-postgresql-cancellation-dispatch.md) | PostgreSQL 数据库时间驱动的 CancellationDispatch | Accepted | +| [ADR-0457](./ADR-0457-worker-pull-cluster-cancellation-delivery.md) | Worker Pull Cluster Cancellation 生产交付 | Accepted | +| [ADR-0458](./ADR-0458-least-privilege-cancellation-diagnostics-and-rearm.md) | 最小权限 Cancellation 诊断与精确 Rearm | Accepted | +| [ADR-0459](./ADR-0459-project-scoped-cancellation-availability-summary.md) | Project-scoped Cancellation 可用性汇总 | Accepted | +| [ADR-0460](./ADR-0460-one-shot-cancellation-status-product-entry.md) | 一次性 Cancellation 状态产品入口 | Accepted | +| [ADR-0461](./ADR-0461-project-scoped-blocked-cancellation-keyset.md) | Project-scoped Blocked Cancellation 固定键集分页 | Accepted | ## 规则 diff --git a/packages/ql3-cluster-admin/src/run-management/runCancellationBlockedList.ts b/packages/ql3-cluster-admin/src/run-management/runCancellationBlockedList.ts new file mode 100644 index 00000000..89496ac2 --- /dev/null +++ b/packages/ql3-cluster-admin/src/run-management/runCancellationBlockedList.ts @@ -0,0 +1,191 @@ +import { randomUUID } from 'node:crypto'; + +import type { ClusterRunManagementClientResult } from './runManagementClient'; +import { + RUN_CANCELLATION_DISPATCH_BLOCKED_LIST_REQUEST_SCHEMA, + normalizeClusterRunManagementCommand, + type ClusterRunManagementCancellationBlockedListCommand, + type ClusterRunManagementCancellationBlockedListTransportResult, +} from './runManagementTransport'; + +export const RUN_CANCELLATION_BLOCKED_LIST_SCHEMA = + 'qinglong/run-cancellation-blocked-list@v1' as const; + +const CURSOR_PREFIX = 'v1.'; +const CURSOR_PAYLOAD_PATTERN = /^[A-Za-z0-9_-]{1,512}$/; +const IDENTIFIER_PATTERN = /^[A-Za-z0-9][A-Za-z0-9._:-]{0,127}$/; + +type BlockedCursor = NonNullable< + ClusterRunManagementCancellationBlockedListCommand['request']['body']['after'] +>; + +export interface RunCancellationBlockedListObservation { + readonly schemaVersion: 1; + readonly schema: typeof RUN_CANCELLATION_BLOCKED_LIST_SCHEMA; + readonly component: 'qinglong3-run-management-client'; + readonly event: 'cancellation_blocked_list_observed'; + readonly requestId: string; + readonly projectId: string; + readonly snapshotAtMs: number; + readonly observedAtMs: number; + readonly items: readonly Readonly<{ + runId: string; + blockedAtMs: number; + }>[]; + readonly truncated: boolean; + readonly nextCursor?: string; +} + +function invalidCursor(): never { + throw new TypeError('Run cancellation blocked cursor is invalid'); +} + +function exact(value: unknown, keys: readonly string[]): Record { + if (!value || typeof value !== 'object' || Array.isArray(value)) { + invalidCursor(); + } + const actual = Object.keys(value as object).sort(); + const expected = [...keys].sort(); + if ( + actual.length !== expected.length || + actual.some((key, index) => key !== expected[index]) + ) { + invalidCursor(); + } + return value as Record; +} + +function normalizeCursor(value: unknown): Readonly { + const cursor = exact(value, ['snapshotAtMs', 'blockedAtMs', 'runId']); + if ( + typeof cursor.snapshotAtMs !== 'number' || + !Number.isSafeInteger(cursor.snapshotAtMs) || + cursor.snapshotAtMs < 0 || + typeof cursor.blockedAtMs !== 'number' || + !Number.isSafeInteger(cursor.blockedAtMs) || + cursor.blockedAtMs < 0 || + cursor.blockedAtMs > cursor.snapshotAtMs || + typeof cursor.runId !== 'string' || + !IDENTIFIER_PATTERN.test(cursor.runId) + ) { + invalidCursor(); + } + return Object.freeze({ + snapshotAtMs: cursor.snapshotAtMs, + blockedAtMs: cursor.blockedAtMs, + runId: cursor.runId, + }); +} + +export function decodeRunCancellationBlockedCursor( + token: string, +): Readonly { + if (typeof token !== 'string' || !token.startsWith(CURSOR_PREFIX)) { + invalidCursor(); + } + const encoded = token.slice(CURSOR_PREFIX.length); + if (!CURSOR_PAYLOAD_PATTERN.test(encoded)) invalidCursor(); + let bytes: Buffer; + let parsed: unknown; + try { + bytes = Buffer.from(encoded, 'base64url'); + if ( + bytes.length < 2 || + bytes.length > 384 || + bytes.toString('base64url') !== encoded + ) { + invalidCursor(); + } + parsed = JSON.parse(bytes.toString('utf8')); + } catch { + invalidCursor(); + } + return normalizeCursor(parsed); +} + +export function encodeRunCancellationBlockedCursor( + value: Readonly, +): string { + const cursor = normalizeCursor(value); + return `${CURSOR_PREFIX}${Buffer.from(JSON.stringify(cursor)).toString( + 'base64url', + )}`; +} + +export function createRunCancellationBlockedListCommand( + projectId: string, + cursorToken?: string, + createUuid: () => string = randomUUID, +): Readonly { + const requestId = createUuid(); + const auditEventId = createUuid(); + let failureAuditEventId = createUuid(); + for ( + let attempts = 0; + failureAuditEventId === auditEventId && attempts < 3; + attempts += 1 + ) { + failureAuditEventId = createUuid(); + } + return normalizeClusterRunManagementCommand({ + schemaVersion: 1, + operation: 'run.cancellation.blocked.list', + request: { + projectId, + requestId, + auditEventId, + failureAuditEventId, + body: { + schema: RUN_CANCELLATION_DISPATCH_BLOCKED_LIST_REQUEST_SCHEMA, + after: + cursorToken === undefined + ? null + : decodeRunCancellationBlockedCursor(cursorToken), + }, + }, + }) as Readonly; +} + +export function projectRunCancellationBlockedList( + result: Readonly, +): Readonly { + if (result.result.operation !== 'run.cancellation.blocked.list') { + throw new TypeError('Run cancellation blocked list requires a list result'); + } + const page = ( + result.result as ClusterRunManagementCancellationBlockedListTransportResult + ).page; + return Object.freeze({ + schemaVersion: 1, + schema: RUN_CANCELLATION_BLOCKED_LIST_SCHEMA, + component: 'qinglong3-run-management-client', + event: 'cancellation_blocked_list_observed', + requestId: result.requestId, + projectId: page.projectId, + snapshotAtMs: page.snapshotAtMs, + observedAtMs: page.observedAtMs, + items: page.items, + truncated: page.truncated, + ...(page.nextCursor === undefined + ? {} + : { nextCursor: encodeRunCancellationBlockedCursor(page.nextCursor) }), + }); +} + +export function formatRunCancellationBlockedListCard( + observation: Readonly, +): string { + return [ + 'QingLong 3.0 / Blocked Cancellations', + `PROJECT ${observation.projectId}`, + `SNAPSHOT ${new Date(observation.snapshotAtMs).toISOString()}`, + `OBSERVED ${new Date(observation.observedAtMs).toISOString()}`, + `ITEMS ${observation.items.length}`, + ...observation.items.map( + (item) => + `BLOCKED ${new Date(item.blockedAtMs).toISOString()} ${item.runId}`, + ), + `NEXT_CURSOR ${observation.nextCursor ?? '-'}`, + `REQUEST ${observation.requestId}`, + ].join('\n'); +} diff --git a/packages/ql3-cluster-admin/src/run-management/runManagement.ts b/packages/ql3-cluster-admin/src/run-management/runManagement.ts index 0ed00efd..4697cd66 100644 --- a/packages/ql3-cluster-admin/src/run-management/runManagement.ts +++ b/packages/ql3-cluster-admin/src/run-management/runManagement.ts @@ -11,6 +11,8 @@ import { RunCancellationDispatchManagementNotFoundError, RunCancellationDispatchManagementUnavailableError, type BlockingCancellationDispatchResult, + type RunCancellationDispatchBlockedCursor, + type RunCancellationDispatchBlockedPage, type RunCancellationDispatchDiagnostic, type RunCancellationDispatchRearmReceipt, type RunCancellationDispatchSummary, @@ -84,6 +86,11 @@ export interface ClusterRunManagementCancellationSummaryRequest { readonly principal: Readonly; } +export interface ClusterRunManagementCancellationBlockedListRequest + extends ClusterRunManagementCancellationSummaryRequest { + readonly after?: Readonly; +} + export interface ClusterRunManagementCancellationRearmRequest extends ClusterRunManagementCancellationInspectRequest { readonly mutationId: string; @@ -102,6 +109,9 @@ export interface ClusterRunManagementService { summarizeCancellation( request: Readonly, ): Promise>; + listBlockedCancellations( + request: Readonly, + ): Promise>; inspectCancellation( request: Readonly, ): Promise>; @@ -267,6 +277,34 @@ function exactCancellationSummaryRequest( } } +function exactCancellationBlockedListRequest( + value: unknown, +): asserts value is Readonly { + const hasAfter = + value !== null && + typeof value === 'object' && + !Array.isArray(value) && + Object.hasOwn(value, 'after'); + if ( + !value || + typeof value !== 'object' || + Array.isArray(value) || + Object.keys(value).sort().join('\0') !== + [ + 'auditEventId', + 'failureAuditEventId', + 'principal', + 'projectId', + 'requestId', + ...(hasAfter ? ['after'] : []), + ] + .sort() + .join('\0') + ) { + throw new ClusterRunManagementRequestError(); + } +} + function exactCancellationRearmRequest( value: unknown, ): asserts value is Readonly { @@ -615,6 +653,99 @@ export function createClusterRunManagementService( throw new ClusterRunManagementUnavailableError({ cause: error }); } }, + async listBlockedCancellations( + requestValue: Readonly, + ) { + exactCancellationBlockedListRequest(requestValue); + const observedAtMs = now(); + let principal: Readonly; + const after = requestValue.after; + if ( + !Number.isSafeInteger(observedAtMs) || + observedAtMs < 0 || + !IDENTIFIER_PATTERN.test(requestValue.projectId) || + !IDENTIFIER_PATTERN.test(requestValue.requestId) || + !validUuid(requestValue.auditEventId) || + !validUuid(requestValue.failureAuditEventId) || + requestValue.auditEventId === requestValue.failureAuditEventId || + (after !== undefined && + (!after || + typeof after !== 'object' || + Array.isArray(after) || + Object.keys(after).sort().join('\0') !== + ['blockedAtMs', 'runId', 'snapshotAtMs'].join('\0') || + !Number.isSafeInteger(after.snapshotAtMs) || + after.snapshotAtMs < 0 || + !Number.isSafeInteger(after.blockedAtMs) || + after.blockedAtMs < 0 || + after.blockedAtMs > after.snapshotAtMs || + !IDENTIFIER_PATTERN.test(after.runId))) + ) { + throw new ClusterRunManagementRequestError(); + } + try { + principal = normalizeSecurityPrincipal( + requestValue.principal, + observedAtMs, + ); + } catch { + throw new ClusterRunManagementRequestError(); + } + let fence: Readonly | null = null; + try { + const decision = await policy.authorize( + principal, + requestValue.projectId, + 'run.read', + ); + fence = decision.fence; + if ( + decision.effect !== 'allow' || + !fence || + fence.bindingVersion === null + ) { + throw new ClusterRunManagementAuthorizationError(); + } + return await cancellationDispatches.listBlocked({ + projectId: requestValue.projectId, + requestId: requestValue.requestId, + auditEventId: requestValue.auditEventId, + principal, + policyFence: fence, + ...(after === undefined ? {} : { after }), + }); + } catch (error) { + try { + await audit.record( + normalizeSecurityAuditRecord({ + eventId: requestValue.failureAuditEventId, + requestId: requestValue.requestId, + operationId: 'run.cancellation.blocked.list', + projectId: requestValue.projectId, + subject: principal.subject, + authenticationId: principal.authenticationId, + outcome: 'denied', + reasons: [failureReason(error)], + fence, + occurredAtMs: observedAtMs, + }), + ); + } catch (auditError) { + throw new ClusterRunManagementUnavailableError({ cause: auditError }); + } + if (error instanceof ClusterRunManagementAuthorizationError) throw error; + if (error instanceof InvalidRunCancellationDispatchManagementError) { + throw new ClusterRunManagementRequestError(); + } + if (error instanceof RunCancellationDispatchManagementConflictError) { + throw new ClusterRunManagementConflictError(); + } + if (error instanceof RunCancellationDispatchManagementUnavailableError) { + throw new ClusterRunManagementUnavailableError({ cause: error }); + } + throw new ClusterRunManagementUnavailableError({ cause: error }); + } + }, async inspectCancellation( requestValue: Readonly, ) { diff --git a/packages/ql3-cluster-admin/src/run-management/runManagementClient.ts b/packages/ql3-cluster-admin/src/run-management/runManagementClient.ts index 2db2490b..a03a4243 100644 --- a/packages/ql3-cluster-admin/src/run-management/runManagementClient.ts +++ b/packages/ql3-cluster-admin/src/run-management/runManagementClient.ts @@ -8,6 +8,7 @@ import { } from '@qinglong/runtime-core/run-cancellation'; import { RUN_STATUSES } from '@qinglong/runtime-core/run'; import { + CANCELLATION_DISPATCH_BLOCKED_PAGE_LIMIT, CANCELLATION_DISPATCH_RESULTS, CANCELLATION_DISPATCH_STATUSES, } from '@qinglong/runtime-core/cancellation-dispatch'; @@ -21,6 +22,7 @@ import { } from '../management-support/pluginPackageManagementClient'; import { RUN_CANCELLATION_DISPATCH_DIAGNOSTIC_SCHEMA, + RUN_CANCELLATION_DISPATCH_BLOCKED_PAGE_SCHEMA, RUN_CANCELLATION_DISPATCH_REARM_RECEIPT_SCHEMA, RUN_CANCELLATION_DISPATCH_SUMMARY_SCHEMA, normalizeClusterRunManagementCommand, @@ -29,6 +31,7 @@ import { } from './runManagementTransport'; const MANAGEMENT_PATH = '/api/v3/runs/management'; +const IDENTIFIER_PATTERN = /^[A-Za-z0-9][A-Za-z0-9._:-]{0,127}$/; export type ClusterRunManagementClientPaths = ClusterPluginPackageManagementClientPaths; @@ -231,6 +234,86 @@ export function validateClusterRunManagementClientResult( envelope as unknown as ClusterRunManagementTransportResult, ); } + if (command.operation === 'run.cancellation.blocked.list') { + const envelope = exact(value, ['schemaVersion', 'operation', 'page']); + if ( + envelope.schemaVersion !== 1 || + envelope.operation !== command.operation + ) { + invalid(); + } + const page = exact(envelope.page, [ + 'schema', + 'projectId', + 'snapshotAtMs', + 'observedAtMs', + 'items', + 'truncated', + ...(Object.hasOwn(envelope.page as object, 'nextCursor') + ? ['nextCursor'] + : []), + ]); + if ( + page.schema !== RUN_CANCELLATION_DISPATCH_BLOCKED_PAGE_SCHEMA || + page.projectId !== command.request.projectId || + !safeInteger(page.snapshotAtMs) || + !safeInteger(page.observedAtMs) || + (page.snapshotAtMs as number) > (page.observedAtMs as number) || + (command.request.body.after === null && + page.snapshotAtMs !== page.observedAtMs) || + (command.request.body.after !== null && + page.snapshotAtMs !== command.request.body.after.snapshotAtMs) || + !Array.isArray(page.items) || + page.items.length > CANCELLATION_DISPATCH_BLOCKED_PAGE_LIMIT || + typeof page.truncated !== 'boolean' || + page.truncated !== Object.hasOwn(page, 'nextCursor') || + (page.truncated && + page.items.length !== CANCELLATION_DISPATCH_BLOCKED_PAGE_LIMIT) + ) { + invalid(); + } + const items = page.items as unknown[]; + let previous = command.request.body.after; + for (const value of items) { + const item = exact(value, ['runId', 'blockedAtMs']); + if ( + typeof item.runId !== 'string' || + !IDENTIFIER_PATTERN.test(item.runId) || + !safeInteger(item.blockedAtMs) || + (item.blockedAtMs as number) > (page.snapshotAtMs as number) || + (previous !== null && + ((item.blockedAtMs as number) < previous.blockedAtMs || + ((item.blockedAtMs as number) === previous.blockedAtMs && + item.runId <= previous.runId))) + ) { + invalid(); + } + previous = { + snapshotAtMs: page.snapshotAtMs as number, + blockedAtMs: item.blockedAtMs as number, + runId: item.runId, + }; + } + if (page.truncated) { + const cursor = exact(page.nextCursor, [ + 'snapshotAtMs', + 'blockedAtMs', + 'runId', + ]); + const last = previous; + if ( + last === null || + cursor.snapshotAtMs !== page.snapshotAtMs || + cursor.blockedAtMs !== last.blockedAtMs || + cursor.runId !== last.runId + ) { + invalid(); + } + } + return Object.freeze( + envelope as unknown as ClusterRunManagementTransportResult, + ); + } if (command.operation === 'run.cancellation.inspect') { const envelope = exact(value, [ 'schemaVersion', diff --git a/packages/ql3-cluster-admin/src/run-management/runManagementClientCli.ts b/packages/ql3-cluster-admin/src/run-management/runManagementClientCli.ts index ae507547..e01d24d9 100644 --- a/packages/ql3-cluster-admin/src/run-management/runManagementClientCli.ts +++ b/packages/ql3-cluster-admin/src/run-management/runManagementClientCli.ts @@ -5,6 +5,12 @@ import { executeClusterRunManagementClient, executeClusterRunManagementCommand, } from './runManagementClient'; +import { + createRunCancellationBlockedListCommand, + decodeRunCancellationBlockedCursor, + formatRunCancellationBlockedListCard, + projectRunCancellationBlockedList, +} from './runCancellationBlockedList'; import { createRunCancellationStatusCommand, formatRunCancellationStatusCard, @@ -14,6 +20,7 @@ import { const USAGE = [ 'Usage: ql3-run-client --config=/absolute/client.json --command=/absolute/command.json --assertion=/absolute/assertion.jwt', ' ql3-run-client status --config=/absolute/client.json --assertion=/absolute/assertion.jwt --project=PROJECT [--format=text|json]', + ' ql3-run-client blocked --config=/absolute/client.json --assertion=/absolute/assertion.jwt --project=PROJECT [--cursor=CURSOR] [--format=text|json]', '', 'Status exit codes: 0=clear, 10=converging, 20=attention_required.', ].join('\n'); @@ -32,18 +39,31 @@ type RunManagementClientArguments = assertionFile: string; projectId: string; format: 'text' | 'json'; + }> + | Readonly<{ + kind: 'blocked'; + configFile: string; + assertionFile: string; + projectId: string; + cursor?: string; + format: 'text' | 'json'; }>; function argumentsFrom( argv: readonly string[], ): Readonly | null { - const statusCount = argv.filter((argument) => argument === 'status').length; - if (statusCount > 0) { - if (statusCount !== 1 || argv.length < 4 || argv.length > 5) return null; + const modes = argv.filter( + (argument) => argument === 'status' || argument === 'blocked', + ); + if (modes.length > 0) { + if (modes.length !== 1 || argv.length < 4 || argv.length > 6) return null; + const kind = modes[0] as 'status' | 'blocked'; const values = new Map(); for (const argument of argv) { - if (argument === 'status') continue; - const match = /^--(config|assertion|project|format)=(.+)$/.exec(argument); + if (argument === kind) continue; + const match = /^--(config|assertion|project|cursor|format)=(.+)$/.exec( + argument, + ); if (!match || values.has(match[1]!)) return null; values.set(match[1]!, match[2]!); } @@ -54,17 +74,28 @@ function argumentsFrom( !values.get('assertion')!.startsWith('/') || !values.has('project') || !PROJECT_ID.test(values.get('project')!) || + (kind === 'status' && values.has('cursor')) || (values.has('format') && values.get('format') !== 'text' && values.get('format') !== 'json') ) { return null; } + if (kind === 'blocked' && values.has('cursor')) { + try { + decodeRunCancellationBlockedCursor(values.get('cursor')!); + } catch { + return null; + } + } return Object.freeze({ - kind: 'status', + kind, configFile: values.get('config')!, assertionFile: values.get('assertion')!, projectId: values.get('project')!, + ...(kind === 'blocked' && values.has('cursor') + ? { cursor: values.get('cursor')! } + : {}), format: (values.get('format') ?? 'text') as 'text' | 'json', }); } @@ -131,6 +162,23 @@ async function run(argv: readonly string[]): Promise { return; } try { + if (paths.kind === 'blocked') { + const result = await executeClusterRunManagementCommand({ + configFile: paths.configFile, + assertionFile: paths.assertionFile, + command: createRunCancellationBlockedListCommand( + paths.projectId, + paths.cursor, + ), + }); + const page = projectRunCancellationBlockedList(result); + process.stdout.write( + paths.format === 'json' + ? `${JSON.stringify(page)}\n` + : `${formatRunCancellationBlockedListCard(page)}\n`, + ); + return; + } if (paths.kind === 'status') { const result = await executeClusterRunManagementCommand({ configFile: paths.configFile, diff --git a/packages/ql3-cluster-admin/src/run-management/runManagementTransport.ts b/packages/ql3-cluster-admin/src/run-management/runManagementTransport.ts index 2aac3807..309c9b1c 100644 --- a/packages/ql3-cluster-admin/src/run-management/runManagementTransport.ts +++ b/packages/ql3-cluster-admin/src/run-management/runManagementTransport.ts @@ -26,6 +26,10 @@ export const RUN_CANCELLATION_DISPATCH_SUMMARY_REQUEST_SCHEMA = 'qinglong/run-cancellation-dispatch-summary-request@v1'; export const RUN_CANCELLATION_DISPATCH_SUMMARY_SCHEMA = 'qinglong/run-cancellation-dispatch-summary@v1'; +export const RUN_CANCELLATION_DISPATCH_BLOCKED_LIST_REQUEST_SCHEMA = + 'qinglong/run-cancellation-dispatch-blocked-list-request@v1'; +export const RUN_CANCELLATION_DISPATCH_BLOCKED_PAGE_SCHEMA = + 'qinglong/run-cancellation-dispatch-blocked-page@v1'; export const RUN_CANCELLATION_DISPATCH_DIAGNOSTIC_SCHEMA = 'qinglong/run-cancellation-dispatch-diagnostic@v1'; export const RUN_CANCELLATION_DISPATCH_REARM_REQUEST_SCHEMA = @@ -96,6 +100,25 @@ export type ClusterRunManagementCancellationSummaryCommand = Readonly<{ }>; }>; +export type ClusterRunManagementCancellationBlockedListCommand = Readonly<{ + schemaVersion: 1; + operation: 'run.cancellation.blocked.list'; + request: Readonly<{ + projectId: string; + requestId: string; + auditEventId: string; + failureAuditEventId: string; + body: Readonly<{ + schema: typeof RUN_CANCELLATION_DISPATCH_BLOCKED_LIST_REQUEST_SCHEMA; + after: Readonly<{ + snapshotAtMs: number; + blockedAtMs: number; + runId: string; + }> | null; + }>; + }>; +}>; + export type ClusterRunManagementCancellationRearmCommand = Readonly<{ schemaVersion: 1; operation: 'run.cancellation.rearm'; @@ -123,6 +146,7 @@ export type ClusterRunManagementCommand = | ClusterRunManagementRetryCommand | ClusterRunManagementStopCommand | ClusterRunManagementCancellationSummaryCommand + | ClusterRunManagementCancellationBlockedListCommand | ClusterRunManagementCancellationInspectCommand | ClusterRunManagementCancellationRearmCommand; @@ -158,6 +182,19 @@ export type ClusterRunManagementCancellationSummaryTransportResult = Readonly<{ >; }>; +export type ClusterRunManagementCancellationBlockedListTransportResult = + Readonly<{ + schemaVersion: 1; + operation: 'run.cancellation.blocked.list'; + page: Readonly< + Awaited< + ReturnType + > & { + schema: typeof RUN_CANCELLATION_DISPATCH_BLOCKED_PAGE_SCHEMA; + } + >; + }>; + export type ClusterRunManagementCancellationRearmTransportResult = Readonly<{ schemaVersion: 1; operation: 'run.cancellation.rearm'; @@ -172,6 +209,7 @@ export type ClusterRunManagementTransportResult = | ClusterRunManagementRetryTransportResult | ClusterRunManagementStopTransportResult | ClusterRunManagementCancellationSummaryTransportResult + | ClusterRunManagementCancellationBlockedListTransportResult | ClusterRunManagementCancellationInspectTransportResult | ClusterRunManagementCancellationRearmTransportResult; @@ -258,6 +296,7 @@ export function normalizeClusterRunManagementCommand( operation !== 'run.retry' && operation !== 'run.stop' && operation !== 'run.cancellation.summary' && + operation !== 'run.cancellation.blocked.list' && operation !== 'run.cancellation.inspect' && operation !== 'run.cancellation.rearm' ) { @@ -274,7 +313,8 @@ export function normalizeClusterRunManagementCommand( 'failureAuditEventId', 'body', ] - : operation === 'run.cancellation.summary' + : operation === 'run.cancellation.summary' || + operation === 'run.cancellation.blocked.list' ? [ 'projectId', 'requestId', @@ -333,6 +373,51 @@ export function normalizeClusterRunManagementCommand( }), }); } + if (operation === 'run.cancellation.blocked.list') { + const body = exact(request.body, ['schema', 'after']); + if (body.schema !== RUN_CANCELLATION_DISPATCH_BLOCKED_LIST_REQUEST_SCHEMA) { + invalid(); + } + let after: ClusterRunManagementCancellationBlockedListCommand['request']['body']['after'] = + null; + if (body.after !== null) { + const cursor = exact(body.after, [ + 'snapshotAtMs', + 'blockedAtMs', + 'runId', + ]); + if ( + typeof cursor.snapshotAtMs !== 'number' || + !Number.isSafeInteger(cursor.snapshotAtMs) || + cursor.snapshotAtMs < 0 || + typeof cursor.blockedAtMs !== 'number' || + !Number.isSafeInteger(cursor.blockedAtMs) || + cursor.blockedAtMs < 0 || + cursor.blockedAtMs > cursor.snapshotAtMs + ) { + invalid(); + } + after = Object.freeze({ + snapshotAtMs: cursor.snapshotAtMs, + blockedAtMs: cursor.blockedAtMs, + runId: identifier(cursor.runId), + }); + } + return Object.freeze({ + schemaVersion: 1, + operation, + request: Object.freeze({ + projectId: identifier(request.projectId), + requestId: identifier(request.requestId), + auditEventId, + failureAuditEventId, + body: Object.freeze({ + schema: RUN_CANCELLATION_DISPATCH_BLOCKED_LIST_REQUEST_SCHEMA, + after, + }), + }), + }); + } if (operation === 'run.cancellation.inspect') { const body = exact(request.body, ['schema']); if (body.schema !== RUN_CANCELLATION_DISPATCH_INSPECT_REQUEST_SCHEMA) { @@ -433,6 +518,7 @@ export function createClusterRunManagementTransport( typeof options.service.retry !== 'function' || typeof options.service.stop !== 'function' || typeof options.service.summarizeCancellation !== 'function' || + typeof options.service.listBlockedCancellations !== 'function' || typeof options.service.inspectCancellation !== 'function' || typeof options.service.rearmCancellation !== 'function' || (options.now !== undefined && typeof options.now !== 'function') @@ -511,6 +597,26 @@ export function createClusterRunManagementTransport( }), }); } + if (command.operation === 'run.cancellation.blocked.list') { + const result = await options.service.listBlockedCancellations({ + projectId: command.request.projectId, + requestId: command.request.requestId, + auditEventId: command.request.auditEventId, + failureAuditEventId: command.request.failureAuditEventId, + principal, + ...(command.request.body.after === null + ? {} + : { after: command.request.body.after }), + }); + return Object.freeze({ + schemaVersion: 1, + operation: command.operation, + page: Object.freeze({ + schema: RUN_CANCELLATION_DISPATCH_BLOCKED_PAGE_SCHEMA, + ...result, + }), + }); + } if (command.operation === 'run.cancellation.inspect') { const result = await options.service.inspectCancellation({ projectId: command.request.projectId, diff --git a/packages/ql3-cluster-admin/test/runCancellationBlockedList.test.cjs b/packages/ql3-cluster-admin/test/runCancellationBlockedList.test.cjs new file mode 100644 index 00000000..eeb4e4cf --- /dev/null +++ b/packages/ql3-cluster-admin/test/runCancellationBlockedList.test.cjs @@ -0,0 +1,104 @@ +'use strict'; + +const assert = require('node:assert/strict'); +const { test } = require('node:test'); + +const { + createRunCancellationBlockedListCommand, + decodeRunCancellationBlockedCursor, + encodeRunCancellationBlockedCursor, + formatRunCancellationBlockedListCard, + projectRunCancellationBlockedList, +} = require('../dist/run-management/runCancellationBlockedList.js'); + +const UUIDS = [ + '019fa000-0000-4000-8000-000000000001', + '019fa000-0000-4000-8000-000000000002', + '019fa000-0000-4000-8000-000000000003', +]; + +test('round-trips one bounded opaque blocked cursor', () => { + const cursor = { + snapshotAtMs: 1_700_000_000_000, + blockedAtMs: 1_699_999_999_000, + runId: 'run-16', + }; + const token = encodeRunCancellationBlockedCursor(cursor); + assert.match(token, /^v1\.[A-Za-z0-9_-]+$/); + assert.deepEqual(decodeRunCancellationBlockedCursor(token), cursor); + for (const invalid of [ + 'v2.abc', + 'v1.***', + `${token}=`, + 'v1.e30', + ]) { + assert.throws(() => decodeRunCancellationBlockedCursor(invalid)); + } +}); + +test('builds one fixed blocked list command with no caller limit', () => { + let index = 0; + const cursor = encodeRunCancellationBlockedCursor({ + snapshotAtMs: 1_700_000_000_000, + blockedAtMs: 1_699_999_999_000, + runId: 'run-16', + }); + const command = createRunCancellationBlockedListCommand( + 'project-1', + cursor, + () => UUIDS[index++], + ); + assert.equal(command.operation, 'run.cancellation.blocked.list'); + assert.equal(command.request.projectId, 'project-1'); + assert.equal(Object.hasOwn(command.request, 'runId'), false); + assert.equal(Object.hasOwn(command.request.body, 'limit'), false); + assert.deepEqual(command.request.body.after, { + snapshotAtMs: 1_700_000_000_000, + blockedAtMs: 1_699_999_999_000, + runId: 'run-16', + }); +}); + +test('projects a low-sensitive page and renders a deterministic card', () => { + const result = { + schemaVersion: 1, + requestId: 'request-blocked-1', + result: { + schemaVersion: 1, + operation: 'run.cancellation.blocked.list', + page: { + schema: 'qinglong/run-cancellation-dispatch-blocked-page@v1', + projectId: 'project-1', + snapshotAtMs: 1_700_000_000_000, + observedAtMs: 1_700_000_000_100, + items: Array.from({ length: 16 }, (_, index) => ({ + runId: `run-${index + 1}`, + blockedAtMs: 1_699_999_999_000 + index, + })), + truncated: true, + nextCursor: { + snapshotAtMs: 1_700_000_000_000, + blockedAtMs: 1_699_999_999_015, + runId: 'run-16', + }, + }, + }, + }; + const observation = projectRunCancellationBlockedList(result); + assert.equal(observation.schema, 'qinglong/run-cancellation-blocked-list@v1'); + assert.equal(observation.items.length, 16); + assert.match(observation.nextCursor, /^v1\./); + const card = formatRunCancellationBlockedListCard(observation); + assert.match(card, /Blocked Cancellations/); + assert.match(card, /run-1/); + assert.match(card, /NEXT_CURSOR\s+v1\./); + for (const forbidden of [ + 'attemptId', + 'lastResult', + 'leaseOwner', + 'leaseToken', + '\u001b', + ]) { + assert.equal(card.includes(forbidden), false); + } +}); diff --git a/packages/ql3-cluster-admin/test/runManagement.test.cjs b/packages/ql3-cluster-admin/test/runManagement.test.cjs index 6aeda9f4..556ec652 100644 --- a/packages/ql3-cluster-admin/test/runManagement.test.cjs +++ b/packages/ql3-cluster-admin/test/runManagement.test.cjs @@ -187,6 +187,18 @@ function fixture(role = 'operator', options = {}) { rowCount: 1, }; } + if ( + text.startsWith( + 'SELECT run_id AS "runId", updated_at_ms AS "blockedAtMs"', + ) + ) { + return { + rows: [ + { runId: 'run-1', blockedAtMs: String(NOW - 1_500) }, + ], + rowCount: 1, + }; + } if ( text.startsWith('SELECT attempt_id AS "attemptId"') && text.includes('FROM "ql3"."run_cancellation_dispatches"') && @@ -425,6 +437,35 @@ test('allows a viewer to summarize Project cancellation availability atomically' ); }); +test('allows a viewer to page blocked Run identities under run.read', async () => { + const { calls, service } = fixture('viewer'); + const listRequest = { + projectId: 'project-1', + requestId: 'request-blocked-1', + auditEventId: '019f9500-0000-4000-8000-000000000071', + failureAuditEventId: '019f9500-0000-4000-8000-000000000072', + principal: request().principal, + }; + const result = await service.listBlockedCancellations(listRequest); + assert.deepEqual(result.items, [ + { runId: 'run-1', blockedAtMs: NOW - 1_500 }, + ]); + assert.equal(result.snapshotAtMs, NOW); + assert.equal(result.truncated, false); + const list = calls.find(({ sql }) => + sql.startsWith( + 'SELECT run_id AS "runId", updated_at_ms AS "blockedAtMs"', + ), + ); + assert.deepEqual(list.params, ['project-1', NOW, null, '', 17]); + const audit = calls.find( + ({ sql, params }) => + sql.startsWith('INSERT INTO "ql3"."security_audit_events"') && + params[2] === 'run.cancellation.blocked.list', + ); + assert.equal(audit.params[0], listRequest.auditEventId); +}); + test('authorizes exact cancellation rearm and keeps the event identity server-side', async () => { const { calls, service } = fixture(); const rearmRequest = { diff --git a/packages/ql3-cluster-admin/test/runManagementClient.test.cjs b/packages/ql3-cluster-admin/test/runManagementClient.test.cjs index 3478fc93..8a77fb76 100644 --- a/packages/ql3-cluster-admin/test/runManagementClient.test.cjs +++ b/packages/ql3-cluster-admin/test/runManagementClient.test.cjs @@ -147,6 +147,21 @@ const summaryCommand = normalizeClusterRunManagementCommand({ }, }); +const blockedListCommand = normalizeClusterRunManagementCommand({ + schemaVersion: 1, + operation: 'run.cancellation.blocked.list', + request: { + projectId: 'project-1', + requestId: 'request-blocked-1', + auditEventId: '019f9400-0000-4000-8000-000000000061', + failureAuditEventId: '019f9400-0000-4000-8000-000000000062', + body: { + schema: 'qinglong/run-cancellation-dispatch-blocked-list-request@v1', + after: null, + }, + }, +}); + const rearmCommand = normalizeClusterRunManagementCommand({ schemaVersion: 1, operation: 'run.cancellation.rearm', @@ -412,6 +427,54 @@ test('validates the fixed low-sensitive Project cancellation summary', () => { } }); +test('validates one snapshot-bound low-sensitive blocked page', () => { + const value = { + schemaVersion: 1, + operation: 'run.cancellation.blocked.list', + page: { + schema: 'qinglong/run-cancellation-dispatch-blocked-page@v1', + projectId: 'project-1', + snapshotAtMs: 1_000_000, + observedAtMs: 1_000_000, + items: [ + { runId: 'run-1', blockedAtMs: 999_100 }, + { runId: 'run-2', blockedAtMs: 999_200 }, + ], + truncated: false, + }, + }; + assert.deepEqual( + validateClusterRunManagementClientResult(value, blockedListCommand), + value, + ); + for (const page of [ + { ...value.page, projectId: 'project-2' }, + { ...value.page, snapshotAtMs: 999_999 }, + { + ...value.page, + items: [value.page.items[1], value.page.items[0]], + }, + { + ...value.page, + items: [{ ...value.page.items[0], attemptId: 'attempt-1' }], + }, + { + ...value.page, + items: [{ ...value.page.items[0], lastResult: 'identity_mismatch' }], + }, + { ...value.page, truncated: true }, + ]) { + assert.throws( + () => + validateClusterRunManagementClientResult( + { ...value, page }, + blockedListCommand, + ), + ClusterPluginPackageManagementClientRequestError, + ); + } +}); + test('binds a rearm receipt to the exact dispatch version, result and delay fences', () => { const value = { schemaVersion: 1, diff --git a/packages/ql3-cluster-admin/test/runManagementClientCli.test.cjs b/packages/ql3-cluster-admin/test/runManagementClientCli.test.cjs index 68975432..cc006a1e 100644 --- a/packages/ql3-cluster-admin/test/runManagementClientCli.test.cjs +++ b/packages/ql3-cluster-admin/test/runManagementClientCli.test.cjs @@ -91,11 +91,27 @@ test('status mode calls only the summary operation and emits an alert exit', asy authorized: request.socket.authorized, command, }); - const bytes = Buffer.from( - JSON.stringify({ - schemaVersion: 1, - requestId: command.request.requestId, - result: { + const managementResult = + command.operation === 'run.cancellation.blocked.list' + ? { + schemaVersion: 1, + operation: 'run.cancellation.blocked.list', + page: { + schema: + 'qinglong/run-cancellation-dispatch-blocked-page@v1', + projectId: 'project-1', + snapshotAtMs: 1_700_000_000_000, + observedAtMs: 1_700_000_000_000, + items: [ + { + runId: 'run-1', + blockedAtMs: 1_699_999_999_000, + }, + ], + truncated: false, + }, + } + : { schemaVersion: 1, operation: 'run.cancellation.summary', summary: { @@ -121,7 +137,12 @@ test('status mode calls only the summary operation and emits an alert exit', asy }, oldestBlockedAtMs: 1_699_999_999_000, }, - }, + }; + const bytes = Buffer.from( + JSON.stringify({ + schemaVersion: 1, + requestId: command.request.requestId, + result: managementResult, }), 'utf8', ); @@ -192,12 +213,37 @@ test('status mode calls only the summary operation and emits an alert exit', asy schema: 'qinglong/run-cancellation-dispatch-summary-request@v1', }); assert.equal(Object.hasOwn(requests[0].command.request, 'runId'), false); + + const blocked = await runCli([ + 'blocked', + `--config=${configFile}`, + `--assertion=${assertionFile}`, + '--project=project-1', + '--format=json', + ]); + assert.equal(blocked.status, 0, blocked.stderr); + assert.equal(blocked.stderr, ''); + const blockedOutput = JSON.parse(blocked.stdout); + assert.equal( + blockedOutput.schema, + 'qinglong/run-cancellation-blocked-list@v1', + ); + assert.deepEqual(blockedOutput.items, [ + { runId: 'run-1', blockedAtMs: 1_699_999_999_000 }, + ]); + assert.equal(requests.length, 2); + assert.equal(requests[1].command.operation, 'run.cancellation.blocked.list'); + assert.deepEqual(requests[1].command.request.body, { + schema: 'qinglong/run-cancellation-dispatch-blocked-list-request@v1', + after: null, + }); }); test('help documents status routing and invalid projects fail before I/O', async () => { const help = await runCli(['--help']); assert.equal(help.status, 0); assert.match(help.stdout, /ql3-run-client status/); + assert.match(help.stdout, /ql3-run-client blocked/); assert.match(help.stdout, /0=clear, 10=converging, 20=attention_required/); const rejected = await runCli([ @@ -214,4 +260,14 @@ test('help documents status routing and invalid projects fail before I/O', async event: 'usage_invalid', code: 'QL3_RUN_MANAGEMENT_CLIENT_USAGE_INVALID', }); + + const invalidCursor = await runCli([ + 'blocked', + '--config=/private/client.json', + '--assertion=/private/assertion.jwt', + '--project=project-1', + '--cursor=v1.***', + ]); + assert.equal(invalidCursor.status, 64); + assert.equal(invalidCursor.stdout, ''); }); diff --git a/packages/ql3-cluster-admin/test/runManagementTransport.test.cjs b/packages/ql3-cluster-admin/test/runManagementTransport.test.cjs index 68e0837c..10c1f52f 100644 --- a/packages/ql3-cluster-admin/test/runManagementTransport.test.cjs +++ b/packages/ql3-cluster-admin/test/runManagementTransport.test.cjs @@ -144,6 +144,16 @@ function summaryResult() { }; } +function blockedListResult() { + return { + projectId: 'project-1', + snapshotAtMs: NOW, + observedAtMs: NOW, + items: [{ runId: 'run-1', blockedAtMs: NOW - 800 }], + truncated: false, + }; +} + function rearmResult() { return { status: 'rearmed', @@ -195,6 +205,24 @@ function summaryCommand(overrides = {}) { }; } +function blockedListCommand(overrides = {}) { + return { + schemaVersion: 1, + operation: 'run.cancellation.blocked.list', + request: { + projectId: 'project-1', + requestId: 'request-blocked-1', + auditEventId: '019f9300-0000-4000-8000-000000000061', + failureAuditEventId: '019f9300-0000-4000-8000-000000000062', + body: { + schema: 'qinglong/run-cancellation-dispatch-blocked-list-request@v1', + after: null, + }, + ...overrides, + }, + }; +} + function rearmCommand(overrides = {}) { return { schemaVersion: 1, @@ -232,6 +260,9 @@ test('routes one exact strong User retry and emits the shared response', async ( async summarizeCancellation() { return summaryResult(); }, + async listBlockedCancellations() { + return blockedListResult(); + }, async inspectCancellation() { return diagnosticResult(); }, @@ -271,6 +302,9 @@ test('routes one exact strong User stop and emits the shared response', async () async summarizeCancellation() { return summaryResult(); }, + async listBlockedCancellations() { + return blockedListResult(); + }, async inspectCancellation() { return diagnosticResult(); }, @@ -310,6 +344,9 @@ test('rejects weak or non-User identity before service authority', async () => { async summarizeCancellation() { return summaryResult(); }, + async listBlockedCancellations() { + return blockedListResult(); + }, async inspectCancellation() { return diagnosticResult(); }, @@ -342,6 +379,7 @@ test('routes bounded cancellation inspection without lease capability data', asy async retry() { return retryResult(); }, async stop() { return stopResult(); }, async summarizeCancellation() { return summaryResult(); }, + async listBlockedCancellations() { return blockedListResult(); }, async inspectCancellation(request) { calls.push(request); return diagnosticResult(); @@ -376,6 +414,7 @@ test('routes one Project-scoped cancellation summary without Run identity', asyn calls.push(request); return summaryResult(); }, + async listBlockedCancellations() { return blockedListResult(); }, async inspectCancellation() { return diagnosticResult(); }, async rearmCancellation() { return rearmResult(); }, }, @@ -398,6 +437,47 @@ test('routes one Project-scoped cancellation summary without Run identity', asyn assert.equal(JSON.stringify(result).includes('leaseOwner'), false); }); +test('routes one fixed Project blocked page without dispatch internals', async () => { + const calls = []; + const transport = createClusterRunManagementTransport({ + now: () => NOW, + service: { + async retry() { return retryResult(); }, + async stop() { return stopResult(); }, + async summarizeCancellation() { return summaryResult(); }, + async listBlockedCancellations(request) { + calls.push(request); + return blockedListResult(); + }, + async inspectCancellation() { return diagnosticResult(); }, + async rearmCancellation() { return rearmResult(); }, + }, + }); + const result = await transport.execute(blockedListCommand(), { + authenticate: async () => principal(), + }); + assert.equal(calls.length, 1); + assert.equal(calls[0].projectId, 'project-1'); + assert.equal(Object.hasOwn(calls[0], 'after'), false); + assert.deepEqual(result, { + schemaVersion: 1, + operation: 'run.cancellation.blocked.list', + page: { + schema: 'qinglong/run-cancellation-dispatch-blocked-page@v1', + ...blockedListResult(), + }, + }); + for (const forbidden of [ + 'attemptId', + 'dispatchVersion', + 'lastResult', + 'leaseOwner', + 'leaseToken', + ]) { + assert.equal(JSON.stringify(result).includes(forbidden), false); + } +}); + test('routes an exact blocked cancellation rearm receipt', async () => { const calls = []; const transport = createClusterRunManagementTransport({ @@ -406,6 +486,7 @@ test('routes an exact blocked cancellation rearm receipt', async () => { async retry() { return retryResult(); }, async stop() { return stopResult(); }, async summarizeCancellation() { return summaryResult(); }, + async listBlockedCancellations() { return blockedListResult(); }, async inspectCancellation() { return diagnosticResult(); }, async rearmCancellation(request) { calls.push(request); diff --git a/packages/ql3-cluster-postgres/src/entrypoints/runManager.ts b/packages/ql3-cluster-postgres/src/entrypoints/runManager.ts index dcf40d1a..21c1e182 100644 --- a/packages/ql3-cluster-postgres/src/entrypoints/runManager.ts +++ b/packages/ql3-cluster-postgres/src/entrypoints/runManager.ts @@ -5,11 +5,15 @@ export { RunCancellationDispatchManagementConflictError, RunCancellationDispatchManagementNotFoundError, RunCancellationDispatchManagementUnavailableError, + RUN_CANCELLATION_DISPATCH_BLOCKED_PAGE_LIMIT, type BlockingCancellationDispatchResult, + type PostgresRunCancellationDispatchBlockedListCommand, type PostgresRunCancellationDispatchInspectCommand, type PostgresRunCancellationDispatchRearmCommand, type PostgresRunCancellationDispatchSummaryCommand, type RunCancellationDispatchDiagnostic, + type RunCancellationDispatchBlockedCursor, + type RunCancellationDispatchBlockedPage, type RunCancellationDispatchRearmReceipt, type RunCancellationDispatchSummary, } from '../run-management/runCancellationDispatchManagementRepository'; diff --git a/packages/ql3-cluster-postgres/src/migration/migrationManifest.ts b/packages/ql3-cluster-postgres/src/migration/migrationManifest.ts index 8b6a4a2e..947c1c7f 100644 --- a/packages/ql3-cluster-postgres/src/migration/migrationManifest.ts +++ b/packages/ql3-cluster-postgres/src/migration/migrationManifest.ts @@ -343,5 +343,10 @@ export const postgresqlMainMigrationManifest: MigrationStreamManifest = checksum: 'e78e24a06dc4c4dbdd859685f28b4bc837a8cfb279eb3512e0a57dc6d27eaaaa', }), + Object.freeze({ + id: 'pg-0068-cancellation-dispatch-project-keyset', + checksum: + '2fcac38386581189db63faacff325356f11c4529a8db9cef6be1a1ca706aaf10', + }), ]), }); diff --git a/packages/ql3-cluster-postgres/src/migrations/index.ts b/packages/ql3-cluster-postgres/src/migrations/index.ts index bb231e12..64faae2a 100644 --- a/packages/ql3-cluster-postgres/src/migrations/index.ts +++ b/packages/ql3-cluster-postgres/src/migrations/index.ts @@ -70,6 +70,7 @@ import { pg0064PluginPackageSecretBindingTransitionApprovalPlansMigration } from import { pg0065ApprovedActionManualRecoveryMigration } from '../approved-action/pg-0065-approved-action-manual-recovery'; import { pg0066CancellationDispatchMigration } from '../run/migrations/pg-0066-cancellation-dispatch'; import { pg0067CancellationDispatchManagementMigration } from '../run-management/pg-0067-cancellation-dispatch-management'; +import { pg0068CancellationDispatchProjectKeysetMigration } from '../run-management/pg-0068-cancellation-dispatch-project-keyset'; export const postgresqlMainMigrationStream: MigrationStreamDefinition = Object.freeze({ @@ -145,5 +146,6 @@ export const postgresqlMainMigrationStream: MigrationStreamDefinition; +export const RUN_CANCELLATION_DISPATCH_BLOCKED_PAGE_LIMIT = + CANCELLATION_DISPATCH_BLOCKED_PAGE_LIMIT; + +export type RunCancellationDispatchBlockedCursor = Readonly<{ + snapshotAtMs: number; + blockedAtMs: number; + runId: string; +}>; + +export type RunCancellationDispatchBlockedPage = Readonly<{ + projectId: string; + snapshotAtMs: number; + observedAtMs: number; + items: readonly Readonly<{ + runId: string; + blockedAtMs: number; + }>[]; + truncated: boolean; + nextCursor?: Readonly; +}>; + interface ProjectManagementAuthority { readonly projectId: string; readonly requestId: string; @@ -98,6 +120,11 @@ interface ManagementAuthority extends ProjectManagementAuthority { export interface PostgresRunCancellationDispatchSummaryCommand extends ProjectManagementAuthority {} +export interface PostgresRunCancellationDispatchBlockedListCommand + extends ProjectManagementAuthority { + readonly after?: Readonly; +} + export interface PostgresRunCancellationDispatchInspectCommand extends ManagementAuthority {} @@ -347,6 +374,48 @@ function normalizeSummaryCommand( }); } +function normalizeBlockedCursor( + value: unknown, +): Readonly { + const cursor = exact(value, ['snapshotAtMs', 'blockedAtMs', 'runId']); + const snapshotAtMs = boundedInteger(cursor.snapshotAtMs, 0); + const blockedAtMs = boundedInteger(cursor.blockedAtMs, 0, snapshotAtMs); + return Object.freeze({ + snapshotAtMs, + blockedAtMs, + runId: identifier(cursor.runId), + }); +} + +function normalizeBlockedListCommand( + value: Readonly, +): Readonly { + const hasAfter = + value !== null && + typeof value === 'object' && + !Array.isArray(value) && + Object.hasOwn(value, 'after'); + const input = exact(value, [ + 'projectId', + 'requestId', + 'auditEventId', + 'principal', + 'policyFence', + ...(hasAfter ? ['after'] : []), + ]); + const authority = normalizeSummaryCommand({ + projectId: input.projectId as string, + requestId: input.requestId as string, + auditEventId: input.auditEventId as string, + principal: input.principal as SecurityPrincipal, + policyFence: input.policyFence as SecurityPolicyFence, + }); + return Object.freeze({ + ...authority, + ...(hasAfter ? { after: normalizeBlockedCursor(input.after) } : {}), + }); +} + function normalizeRearmCommand( value: Readonly, ): Readonly { @@ -474,6 +543,7 @@ async function recordAllowedAudit( command: Readonly, operationId: | 'run.cancellation.summary' + | 'run.cancellation.blocked.list' | 'run.cancellation.inspect' | 'run.cancellation.rearm', observedAtMs: number, @@ -603,6 +673,76 @@ function summaryProjection( }); } +function storedIdentifier(row: Row, key: string): string { + const value = text(row, key); + if (!IDENTIFIER_PATTERN.test(value)) { + throw new TypeError( + `PostgreSQL cancellation management ${key} is invalid`, + ); + } + return value; +} + +function blockedPageProjection( + command: Readonly, + observedAtMs: number, + snapshotAtMs: number, + rows: readonly Row[], +): Readonly { + if (rows.length > RUN_CANCELLATION_DISPATCH_BLOCKED_PAGE_LIMIT + 1) { + throw new TypeError( + 'PostgreSQL cancellation management blocked page is invalid', + ); + } + const projected = rows.map((row) => + Object.freeze({ + runId: storedIdentifier(row, 'runId'), + blockedAtMs: integer(row, 'blockedAtMs'), + }), + ); + let previous = command.after; + for (const item of projected) { + if ( + item.blockedAtMs > snapshotAtMs || + (previous !== undefined && + (item.blockedAtMs < previous.blockedAtMs || + (item.blockedAtMs === previous.blockedAtMs && + item.runId <= previous.runId))) + ) { + throw new TypeError( + 'PostgreSQL cancellation management blocked cursor order is invalid', + ); + } + previous = Object.freeze({ + snapshotAtMs, + blockedAtMs: item.blockedAtMs, + runId: item.runId, + }); + } + const truncated = + projected.length > RUN_CANCELLATION_DISPATCH_BLOCKED_PAGE_LIMIT; + const items = Object.freeze( + projected.slice(0, RUN_CANCELLATION_DISPATCH_BLOCKED_PAGE_LIMIT), + ); + const last = items.at(-1); + return Object.freeze({ + projectId: command.projectId, + snapshotAtMs, + observedAtMs, + items, + truncated, + ...(truncated && last + ? { + nextCursor: Object.freeze({ + snapshotAtMs, + blockedAtMs: last.blockedAtMs, + runId: last.runId, + }), + } + : {}), + }); +} + function runStatus(row: Row): RunStatus { const value = text(row, 'runStatus') as RunStatus; if (!RUN_STATUSES.includes(value)) { @@ -818,6 +958,53 @@ export class PostgresRunCancellationDispatchManagementRepository { }); } + listBlocked( + value: Readonly, + ): Promise> { + const command = normalizeBlockedListCommand(value); + return this.transaction(async (client) => { + const observedAtMs = await databaseNow(client); + const snapshotAtMs = command.after?.snapshotAtMs ?? observedAtMs; + if (snapshotAtMs > observedAtMs) { + throw new InvalidRunCancellationDispatchManagementError(); + } + const authorized = Object.freeze({ + ...command, + principal: strongPrincipal(command.principal, observedAtMs), + }); + await confirmAuthorization(client, authorized); + const result = await client.query( + `SELECT run_id AS "runId", updated_at_ms AS "blockedAtMs" + FROM "ql3"."run_cancellation_dispatches" + WHERE project_id = $1 AND status = 'blocked' + AND updated_at_ms <= $2 + AND ($3::bigint IS NULL OR + (updated_at_ms, run_id) > ($3::bigint, $4::varchar)) + ORDER BY updated_at_ms ASC, run_id ASC + LIMIT $5`, + [ + command.projectId, + snapshotAtMs, + command.after?.blockedAtMs ?? null, + command.after?.runId ?? '', + RUN_CANCELLATION_DISPATCH_BLOCKED_PAGE_LIMIT + 1, + ], + ); + await recordAllowedAudit( + client, + authorized, + 'run.cancellation.blocked.list', + observedAtMs, + ); + return blockedPageProjection( + command, + observedAtMs, + snapshotAtMs, + result.rows, + ); + }); + } + inspect( value: Readonly, ): Promise> { diff --git a/packages/ql3-cluster-postgres/src/run/cancellationDispatchRepository.ts b/packages/ql3-cluster-postgres/src/run/cancellationDispatchRepository.ts index 54180a61..54286110 100644 --- a/packages/ql3-cluster-postgres/src/run/cancellationDispatchRepository.ts +++ b/packages/ql3-cluster-postgres/src/run/cancellationDispatchRepository.ts @@ -180,7 +180,7 @@ export class PostgresCancellationDispatchRepository return this.transaction(async (client) => { const nowMs = await databaseNow(client); const run = await client.query( - `SELECT execution_owner AS "executionOwner", status, + `SELECT project_id AS "projectId", execution_owner AS "executionOwner", status, cancel_requested_at_ms AS "cancelRequestedAtMs" FROM "ql3"."runs" WHERE id = $1 FOR UPDATE`, [command.runId], @@ -222,11 +222,17 @@ export class PostgresCancellationDispatchRepository if (dispatchResult.rows.length === 0) { dispatchResult = await client.query( `INSERT INTO "ql3"."run_cancellation_dispatches" ( - run_id, attempt_id, status, version, dispatch_count, + project_id, run_id, attempt_id, status, version, dispatch_count, next_attempt_at_ms, created_at_ms, updated_at_ms - ) VALUES ($1, $2, 'pending', 0, 0, $3, $4, $4) + ) VALUES ($5, $1, $2, 'pending', 0, 0, $3, $4, $4) RETURNING ${DISPATCH_COLUMNS}`, - [command.runId, command.attemptId, command.requestedAtMs, nowMs], + [ + command.runId, + command.attemptId, + command.requestedAtMs, + nowMs, + text(runRow, 'projectId'), + ], ); } if (dispatchResult.rows.length !== 1) { diff --git a/packages/ql3-cluster-postgres/src/schema/schema.ts b/packages/ql3-cluster-postgres/src/schema/schema.ts index dee52227..55df24f0 100644 --- a/packages/ql3-cluster-postgres/src/schema/schema.ts +++ b/packages/ql3-cluster-postgres/src/schema/schema.ts @@ -3592,6 +3592,7 @@ export const runs = ql3Schema.table( uniqueIndex('ql3_runs_project_idempotency_uidx') .on(table.projectId, table.idempotencyKey) .where(sql`${table.idempotencyKey} is not null`), + uniqueIndex('ql3_runs_project_id_uidx').on(table.projectId, table.id), index('ql3_runs_project_created_idx').on( table.projectId, table.createdAtMs, @@ -4834,6 +4835,7 @@ export const runAttempts = ql3Schema.table( export const runCancellationDispatches = ql3Schema.table( 'run_cancellation_dispatches', { + projectId: varchar('project_id', { length: 128 }).notNull(), runId: varchar('run_id', { length: 36 }).primaryKey(), attemptId: varchar('attempt_id', { length: 36 }).notNull(), status: varchar('status', { length: 32 }).notNull(), @@ -4851,8 +4853,8 @@ export const runCancellationDispatches = ql3Schema.table( (table) => [ foreignKey({ name: 'ql3_run_cancellation_dispatch_run_fk', - columns: [table.runId], - foreignColumns: [runs.id], + columns: [table.projectId, table.runId], + foreignColumns: [runs.projectId, runs.id], }) .onDelete('cascade') .onUpdate('restrict'), @@ -4897,6 +4899,9 @@ export const runCancellationDispatches = ql3Schema.table( index('ql3_run_cancellation_dispatch_lease_expiry_idx') .on(table.leaseExpiresAtMs, table.runId) .where(sql`${table.status} = 'leased'`), + index('ql3_run_cancellation_dispatch_project_blocked_idx') + .on(table.projectId, table.updatedAtMs, table.runId) + .where(sql`${table.status} = 'blocked'`), ], ); diff --git a/packages/ql3-cluster-postgres/src/schema/schemaContract.ts b/packages/ql3-cluster-postgres/src/schema/schemaContract.ts index 7fd1e01d..e6ce6b5e 100644 --- a/packages/ql3-cluster-postgres/src/schema/schemaContract.ts +++ b/packages/ql3-cluster-postgres/src/schema/schemaContract.ts @@ -21,13 +21,14 @@ export interface PostgresSchemaContractTrigger { export interface PostgresSchemaContract { readonly schema: 'ql3'; readonly contractName: 'control-core'; - readonly contractVersion: 66; - readonly migrationId: 'pg-0067-cancellation-dispatch-management'; + readonly contractVersion: 67; + readonly migrationId: 'pg-0068-cancellation-dispatch-project-keyset'; readonly minimumServerMajor: 16; readonly maximumServerMajor: 18; readonly capabilities: Readonly<{ run_core: 1; run_cancellation_dispatch: 1; + run_cancellation_dispatch_blocked_list: 1; run_cancellation_dispatch_management: 1; run_attempt_log_retention: 1; run_management_boundary: 1; @@ -120,8 +121,8 @@ export const postgresqlControlSchemaContract: PostgresSchemaContract = Object.freeze({ schema: 'ql3', contractName: 'control-core', - contractVersion: 66, - migrationId: 'pg-0067-cancellation-dispatch-management', + contractVersion: 67, + migrationId: 'pg-0068-cancellation-dispatch-project-keyset', minimumServerMajor: 16, maximumServerMajor: 18, capabilities: Object.freeze({ @@ -170,6 +171,7 @@ export const postgresqlControlSchemaContract: PostgresSchemaContract = project_tool_definition_snapshot: 1, run_core: 1, run_cancellation_dispatch: 1, + run_cancellation_dispatch_blocked_list: 1, run_cancellation_dispatch_management: 1, run_attempt_log_retention: 1, run_management_boundary: 1, @@ -1294,6 +1296,7 @@ export const postgresqlControlSchemaContract: PostgresSchemaContract = 'error_summary', ]), table('run_cancellation_dispatches', [ + 'project_id', 'run_id', 'attempt_id', 'status', @@ -1718,6 +1721,7 @@ export const postgresqlControlSchemaContract: PostgresSchemaContract = 'ql3_api_credential_mutations_actor_idx', 'runs_pkey', 'ql3_runs_project_idempotency_uidx', + 'ql3_runs_project_id_uidx', 'ql3_runs_project_created_idx', 'ql3_runs_task_created_idx', 'ql3_runs_dispatch_candidates_idx', @@ -1796,6 +1800,7 @@ export const postgresqlControlSchemaContract: PostgresSchemaContract = 'run_cancellation_dispatches_pkey', 'ql3_run_cancellation_dispatch_due_idx', 'ql3_run_cancellation_dispatch_lease_expiry_idx', + 'ql3_run_cancellation_dispatch_project_blocked_idx', 'run_attempt_log_retention_controls_pkey', 'ql3_run_log_retention_control_artifact_key', 'ql3_run_log_retention_retry_idx', diff --git a/packages/ql3-cluster-postgres/test/postgresqlMigrationDefinitions.test.cjs b/packages/ql3-cluster-postgres/test/postgresqlMigrationDefinitions.test.cjs index dce51475..8fa5b217 100644 --- a/packages/ql3-cluster-postgres/test/postgresqlMigrationDefinitions.test.cjs +++ b/packages/ql3-cluster-postgres/test/postgresqlMigrationDefinitions.test.cjs @@ -118,6 +118,7 @@ test('defines the immutable PostgreSQL capability and Run core stream', async () 'pg-0065-approved-action-manual-recovery', 'pg-0066-cancellation-dispatch', 'pg-0067-cancellation-dispatch-management', + 'pg-0068-cancellation-dispatch-project-keyset', ], ); for (const migration of postgresqlMainMigrationStream.migrations) { @@ -591,6 +592,11 @@ test('freezes every published PostgreSQL migration checksum', () => { checksum: 'e78e24a06dc4c4dbdd859685f28b4bc837a8cfb279eb3512e0a57dc6d27eaaaa', }, + { + id: 'pg-0068-cancellation-dispatch-project-keyset', + checksum: + '2fcac38386581189db63faacff325356f11c4529a8db9cef6be1a1ca706aaf10', + }, ]; assert.deepEqual( postgresqlMainMigrationStream.migrations.map(({ id, checksum }) => ({ @@ -2373,3 +2379,45 @@ test('advances capability v66 with least-privilege cancellation diagnostics and assert.match(sql, /contract_version = 65/); assert.match(sql, /migration_id = 'pg-0066-cancellation-dispatch'/); }); + +test('advances capability v67 with a Project-scoped blocked keyset', async () => { + const migration = migrationById( + 'pg-0068-cancellation-dispatch-project-keyset', + ); + const statements = []; + await migration.up({ + async query(statement) { + statements.push(statement); + return { rows: [] }; + }, + }); + const sql = statements.join('\n'); + assert.match( + sql, + /ADD COLUMN project_id varchar\(128\)/, + ); + assert.match( + sql, + /SET project_id = run\.project_id FROM "ql3"\."runs" AS run/, + ); + assert.match(sql, /ALTER COLUMN project_id SET NOT NULL/); + assert.match( + sql, + /CREATE UNIQUE INDEX ql3_runs_project_id_uidx ON "ql3"\."runs" \(project_id, id\)/, + ); + assert.match( + sql, + /FOREIGN KEY \(project_id, run_id\) REFERENCES "ql3"\."runs" \(project_id, id\)/, + ); + assert.match( + sql, + /CREATE INDEX ql3_run_cancellation_dispatch_project_blocked_idx[\s\S]+\(project_id, updated_at_ms, run_id\) WHERE status = 'blocked'/, + ); + assert.match(sql, /contract_version = 67/); + assert.match(sql, /"run_cancellation_dispatch_blocked_list":1/); + assert.match(sql, /contract_version = 66/); + assert.match( + sql, + /migration_id = 'pg-0067-cancellation-dispatch-management'/, + ); +}); diff --git a/packages/ql3-cluster-postgres/test/postgresqlSchemaReadiness.test.cjs b/packages/ql3-cluster-postgres/test/postgresqlSchemaReadiness.test.cjs index 133be40d..d6b152a2 100644 --- a/packages/ql3-cluster-postgres/test/postgresqlSchemaReadiness.test.cjs +++ b/packages/ql3-cluster-postgres/test/postgresqlSchemaReadiness.test.cjs @@ -835,7 +835,7 @@ test('accepts the exact PostgreSQL control schema and least-privilege runtime ro serverMajor: 16, currentUser: 'ql3_runtime', contractName: 'control-core', - contractVersion: 66, + contractVersion: 67, migrationIds: [ 'pg-0001-schema-capability', 'pg-0002-run-core', @@ -904,6 +904,7 @@ test('accepts the exact PostgreSQL control schema and least-privilege runtime ro 'pg-0065-approved-action-manual-recovery', 'pg-0066-cancellation-dispatch', 'pg-0067-cancellation-dispatch-management', + 'pg-0068-cancellation-dispatch-project-keyset', ], }); }); @@ -934,10 +935,10 @@ test('accepts the exact schema and isolated least-privilege admin role', async ( }), ); assert.equal(report.currentUser, 'ql3_admin'); - assert.equal(report.contractVersion, 66); + assert.equal(report.contractVersion, 67); assert.equal( report.migrationIds.at(-1), - 'pg-0067-cancellation-dispatch-management', + 'pg-0068-cancellation-dispatch-project-keyset', ); }); @@ -950,10 +951,10 @@ test('accepts the isolated least-privilege automation manager role', async () => }), ); assert.equal(report.currentUser, 'ql3_automation_manager'); - assert.equal(report.contractVersion, 66); + assert.equal(report.contractVersion, 67); assert.equal( report.migrationIds.at(-1), - 'pg-0067-cancellation-dispatch-management', + 'pg-0068-cancellation-dispatch-project-keyset', ); const widened = automationManagerPrivileges(); @@ -982,10 +983,10 @@ test('accepts the isolated least-privilege human Approval manager role', async ( }), ); assert.equal(report.currentUser, 'ql3_approval_manager'); - assert.equal(report.contractVersion, 66); + assert.equal(report.contractVersion, 67); assert.equal( report.migrationIds.at(-1), - 'pg-0067-cancellation-dispatch-management', + 'pg-0068-cancellation-dispatch-project-keyset', ); const widened = approvalManagerPrivileges(); @@ -1016,10 +1017,10 @@ test('accepts the isolated least-privilege Run manager role', async () => { }), ); assert.equal(report.currentUser, 'ql3_run_manager'); - assert.equal(report.contractVersion, 66); + assert.equal(report.contractVersion, 67); assert.equal( report.migrationIds.at(-1), - 'pg-0067-cancellation-dispatch-management', + 'pg-0068-cancellation-dispatch-project-keyset', ); const widened = runManagerPrivileges(); @@ -1180,10 +1181,10 @@ test('accepts the exact schema and isolated Worker ingress role', async () => { }), ); assert.equal(report.currentUser, 'ql3_worker_ingress'); - assert.equal(report.contractVersion, 66); + assert.equal(report.contractVersion, 67); assert.equal( report.migrationIds.at(-1), - 'pg-0067-cancellation-dispatch-management', + 'pg-0068-cancellation-dispatch-project-keyset', ); }); diff --git a/packages/ql3-cluster-postgres/test/runCancellationDispatchManagementRepository.test.cjs b/packages/ql3-cluster-postgres/test/runCancellationDispatchManagementRepository.test.cjs index 8d55f322..7bb38e6b 100644 --- a/packages/ql3-cluster-postgres/test/runCancellationDispatchManagementRepository.test.cjs +++ b/packages/ql3-cluster-postgres/test/runCancellationDispatchManagementRepository.test.cjs @@ -47,6 +47,16 @@ function summaryCommand(overrides = {}) { return { ...authority, ...overrides }; } +function blockedListCommand(after) { + return { + ...summaryCommand({ + requestId: 'request-blocked-1', + auditEventId: '019f9600-0000-4000-8000-000000000021', + }), + ...(after === undefined ? {} : { after }), + }; +} + function runRow() { return { projectId: 'project-1', @@ -133,6 +143,13 @@ function fixture(options = {}) { rowCount: 1, }; } + if ( + text.startsWith( + 'SELECT run_id AS "runId", updated_at_ms AS "blockedAtMs"', + ) + ) { + return { rows: options.blockedRows ?? [], rowCount: 0 }; + } if ( text.startsWith('SELECT attempt_id AS "attemptId"') && !text.includes('dispatchStatus') && @@ -296,6 +313,70 @@ test('derives clear and converging assessments from fixed status counts', async assert.equal(result.operatorAction, 'wait'); }); +test('lists one fixed oldest-first blocked page with a snapshot cursor', async () => { + const blockedRows = Array.from({ length: 17 }, (_, index) => ({ + runId: `run-${String(index + 1).padStart(2, '0')}`, + blockedAtMs: String(NOW - 100 + index), + })); + const { calls, repository } = fixture({ blockedRows }); + const result = await repository.listBlocked(blockedListCommand()); + assert.equal(result.projectId, 'project-1'); + assert.equal(result.snapshotAtMs, NOW); + assert.equal(result.observedAtMs, NOW); + assert.equal(result.items.length, 16); + assert.equal(result.items[0].runId, 'run-01'); + assert.equal(result.items[15].runId, 'run-16'); + assert.equal(result.truncated, true); + assert.deepEqual(result.nextCursor, { + snapshotAtMs: NOW, + blockedAtMs: NOW - 85, + runId: 'run-16', + }); + const read = calls.find(({ sql }) => + sql.startsWith( + 'SELECT run_id AS "runId", updated_at_ms AS "blockedAtMs"', + ), + ); + assert.deepEqual(read.params, ['project-1', NOW, null, '', 17]); + assert.match( + read.sql, + /project_id = \$1 AND status = 'blocked'[\s\S]+ORDER BY updated_at_ms ASC, run_id ASC/, + ); + const audit = calls.find( + ({ sql, params }) => + sql.startsWith('INSERT INTO "ql3"."security_audit_events"') && + params[2] === 'run.cancellation.blocked.list', + ); + assert.equal(audit.params[0], blockedListCommand().auditEventId); +}); + +test('continues only inside the original blocked snapshot', async () => { + const after = { + snapshotAtMs: NOW - 50, + blockedAtMs: NOW - 80, + runId: 'run-03', + }; + const { calls, repository } = fixture({ + blockedRows: [{ runId: 'run-04', blockedAtMs: String(NOW - 79) }], + }); + const result = await repository.listBlocked(blockedListCommand(after)); + assert.equal(result.snapshotAtMs, NOW - 50); + assert.equal(result.truncated, false); + assert.equal(Object.hasOwn(result, 'nextCursor'), false); + const read = calls.find(({ sql }) => + sql.startsWith( + 'SELECT run_id AS "runId", updated_at_ms AS "blockedAtMs"', + ), + ); + assert.deepEqual(read.params, [ + 'project-1', + NOW - 50, + NOW - 80, + 'run-03', + 17, + ]); +}); + test('rearms an exact blocked dispatch with one event and allowed audit', async () => { const { calls, repository } = fixture(); const result = await repository.rearm(rearmCommand()); diff --git a/packages/ql3-runtime-core/src/run/cancellation-dispatch/cancellationDispatch.ts b/packages/ql3-runtime-core/src/run/cancellation-dispatch/cancellationDispatch.ts index b82828d2..24475f8a 100644 --- a/packages/ql3-runtime-core/src/run/cancellation-dispatch/cancellationDispatch.ts +++ b/packages/ql3-runtime-core/src/run/cancellation-dispatch/cancellationDispatch.ts @@ -49,6 +49,7 @@ const CANCELLATION_DISPATCH_RETRY_HISTORY_RESULTS = new Set< export const MAX_CANCELLATION_DISPATCH_LEASE_MS = 5 * 60_000; export const MAX_CANCELLATION_DISPATCH_RETRY_DELAY_MS = 24 * 60 * 60_000; +export const CANCELLATION_DISPATCH_BLOCKED_PAGE_LIMIT = 16; export interface CancellationDispatchRecord { readonly runId: string; diff --git a/scripts/ql3-postgres-ha-cancellation-dispatch-fixture.cjs b/scripts/ql3-postgres-ha-cancellation-dispatch-fixture.cjs index f3e01dd0..f4bffd0e 100644 --- a/scripts/ql3-postgres-ha-cancellation-dispatch-fixture.cjs +++ b/scripts/ql3-postgres-ha-cancellation-dispatch-fixture.cjs @@ -25,6 +25,7 @@ const FIXTURE = Object.freeze({ terminalEventId: 'ha-cancel-terminal-event-d363', settledEventId: 'ha-cancel-settled-event-d363', blockedEventId: 'ha-cancel-blocked-event-d365', + blockedListAuditEventId: '019f9700-0000-4000-8000-000000000006', summaryAuditEventId: '019f9700-0000-4000-8000-000000000005', inspectAuditEventId: '019f9700-0000-4000-8000-000000000001', rearmAuditEventId: '019f9700-0000-4000-8000-000000000002', @@ -263,6 +264,59 @@ async function persistCancellationDispatchHaFixture(options) { principal, policyFence: Object.freeze({ projectVersion: 1, bindingVersion: 1 }), }); + const blockedPage = await management.listBlocked({ + projectId: FIXTURE.projectId, + requestId: 'ha-cancel-blocked-list-d368', + auditEventId: FIXTURE.blockedListAuditEventId, + principal, + policyFence: authority.policyFence, + }); + assert.equal(blockedPage.projectId, FIXTURE.projectId); + assert.equal(blockedPage.snapshotAtMs, blockedPage.observedAtMs); + assert.deepEqual(blockedPage.items, [ + { + runId: FIXTURE.runId, + blockedAtMs: blocked.dispatch.updatedAtMs, + }, + ]); + assert.equal(blockedPage.truncated, false); + assert.equal(Object.hasOwn(blockedPage, 'nextCursor'), false); + assert.equal(JSON.stringify(blockedPage).includes('attemptId'), false); + assert.equal(JSON.stringify(blockedPage).includes('lastResult'), false); + assert.equal(JSON.stringify(blockedPage).includes('leaseOwner'), false); + assert.equal(JSON.stringify(blockedPage).includes('leaseToken'), false); + const blockedListAudit = await migrationPool.query( + `SELECT operation_id AS "operationId", outcome + FROM "ql3"."security_audit_events" WHERE event_id = $1`, + [FIXTURE.blockedListAuditEventId], + ); + assert.deepEqual(blockedListAudit.rows, [ + { + operationId: 'run.cancellation.blocked.list', + outcome: 'allowed', + }, + ]); + await migrationPool.query('BEGIN'); + let blockedListPlan; + try { + await migrationPool.query('SET LOCAL enable_seqscan = off'); + blockedListPlan = await migrationPool.query( + `EXPLAIN (FORMAT JSON, COSTS OFF) + SELECT run_id, updated_at_ms + FROM "ql3"."run_cancellation_dispatches" + WHERE project_id = $1 AND status = 'blocked' + AND updated_at_ms <= $2 + ORDER BY updated_at_ms ASC, run_id ASC + LIMIT 17`, + [FIXTURE.projectId, blockedPage.snapshotAtMs], + ); + } finally { + await migrationPool.query('ROLLBACK'); + } + assert.match( + JSON.stringify(blockedListPlan.rows), + /ql3_run_cancellation_dispatch_project_blocked_idx/, + ); const summary = await management.summary({ projectId: FIXTURE.projectId, requestId: 'ha-cancel-summary-d366', @@ -398,6 +452,8 @@ async function persistCancellationDispatchHaFixture(options) { retryDeferredUntilDue: true, operatorDiagnosticLowSensitive: true, operatorSummaryLowSensitiveAndActionable: true, + blockedListLowSensitiveAndSnapshotBound: true, + blockedListUsesProjectKeysetIndex: true, manualBlockedRearmExact: true, manualRearmDeferredUntilDue: true, productionDeliverySettledBeforeStop: true, diff --git a/scripts/ql3-postgres-ha-contract.cjs b/scripts/ql3-postgres-ha-contract.cjs index a46cd1ed..1b9c2ac8 100644 --- a/scripts/ql3-postgres-ha-contract.cjs +++ b/scripts/ql3-postgres-ha-contract.cjs @@ -13987,6 +13987,9 @@ async function main(argv = process.argv.slice(2)) { cancellationDispatch.retryDeferredUntilDue, cancellationDispatchSummaryIsLowSensitiveAndActionable: cancellationDispatch.operatorSummaryLowSensitiveAndActionable, + cancellationDispatchBlockedListIsBoundedAndProjectIndexed: + cancellationDispatch.blockedListLowSensitiveAndSnapshotBound && + cancellationDispatch.blockedListUsesProjectKeysetIndex, cancellationDispatchReplicatesAndSurvivesPromotion: cancellationDispatch.replicatedBeforePromotion && cancellationDispatch.survivedPromotion, diff --git a/test/back/ql3PackageBoundaryAudit.test.cjs b/test/back/ql3PackageBoundaryAudit.test.cjs index 5bdcb5d8..3db55e27 100644 --- a/test/back/ql3PackageBoundaryAudit.test.cjs +++ b/test/back/ql3PackageBoundaryAudit.test.cjs @@ -340,10 +340,10 @@ test('current QL3 workspace has exactly eighteen reviewed package boundaries', ( rootSourceFileRoles: clusterAdmin.rootSourceFileRoles, }, { - sourceFiles: 123, + sourceFiles: 124, rootSourceFiles: 1, rootSourceLines: 61, - nestedSourceFiles: 122, + nestedSourceFiles: 123, rootSourceFileRoles: { 'modelInvocationMigrationCli.ts': 'binary_entry', }, @@ -421,10 +421,10 @@ test('current QL3 workspace has exactly eighteen reviewed package boundaries', ( rootSourceFileRoles: clusterPostgres.rootSourceFileRoles, }, { - sourceFiles: 172, + sourceFiles: 173, rootSourceFiles: 1, rootSourceLines: 126, - nestedSourceFiles: 171, + nestedSourceFiles: 172, rootSourceFileRoles: { 'index.ts': 'public_export' }, }, );