From 5d3b40ce3b19a52eb0836298280016b22b61917a Mon Sep 17 00:00:00 2001 From: whyour Date: Wed, 12 Aug 2026 04:41:59 +0800 Subject: [PATCH] feat(ql3): complete cluster log retention lifecycle --- docs/QINGLONG_3_0_ARCHITECTURE_RFC.md | 29 +- ...-0379-cluster-run-attempt-log-retention.md | 35 +- docs/adr/README.md | 2 +- .../src/application-runtime/application.ts | 8 + .../clusterControlRuntime.ts | 208 +++++++- .../productionApplication.ts | 1 + .../src/artifact/s3ArtifactStore.ts | 134 ++++- .../src/artifact/workerArtifactBinding.ts | 8 +- .../src/production-process/config.ts | 101 ++++ .../production-process/processApplication.ts | 70 ++- .../run/runAttemptLogRetentionLifecycle.ts | 490 ++++++++++++++++++ .../test/bootstrap.test.cjs | 2 + .../ql3-cluster-control/test/config.test.cjs | 57 ++ .../test/processApplication.test.cjs | 41 +- .../test/productionApplication.test.cjs | 81 +++ .../runAttemptLogRetentionLifecycle.test.cjs | 281 ++++++++++ .../test/s3ArtifactStore.integration.test.cjs | 131 ++++- .../test/s3ArtifactStore.test.cjs | 162 ++++++ scripts/ql3-postgres-ha-contract.cjs | 270 +++++++++- test/back/ql3PackageBoundaryAudit.test.cjs | 4 +- 20 files changed, 2026 insertions(+), 89 deletions(-) create mode 100644 packages/ql3-cluster-control/src/run/runAttemptLogRetentionLifecycle.ts create mode 100644 packages/ql3-cluster-control/test/runAttemptLogRetentionLifecycle.test.cjs diff --git a/docs/QINGLONG_3_0_ARCHITECTURE_RFC.md b/docs/QINGLONG_3_0_ARCHITECTURE_RFC.md index 8748491f..9d603952 100644 --- a/docs/QINGLONG_3_0_ARCHITECTURE_RFC.md +++ b/docs/QINGLONG_3_0_ARCHITECTURE_RFC.md @@ -11,17 +11,24 @@ 最新增量证据(2026-08-12): -- D-291/ADR-0379(进行中,PostgreSQL authority 阶段完成) - Cluster Run Attempt 日志 retention 第一阶段已冻结并实现多副本权威边界,且没有新增 package:共享 claim contract 位于既有 - Runtime Core Run log-retention 目录,PostgreSQL v54 增加 durable control 与 immutable tombstone、terminal remote Worker - candidate index 和 `run_attempt_log_retention` capability。候选只允许 runtime-owned、Run/Attempt 双终态、非 lost、canonical - `remote_worker/wlog-*` 且超过 cutoff;每批最多 16 条,通过短 `READ COMMITTED` 事务和 `FOR UPDATE ... SKIP LOCKED` 获取 - owner/token/version/expiry fence。S3 网络调用明确位于数据库事务之外;finalize 在第二个短事务中重验完整 fence 与 durable - Run/Attempt identity,写 exact retirement record 后删除 control。retry/manual 持久化失败分类并提供最长 24 小时 backoff;lease - 过期可由其他副本安全接管。`ql3_runtime` 可完整维护 control,但 tombstone 只有 SELECT/INSERT,其他角色保持零权限;repository - 已提供 digest/identity 失败关闭的 retention state read,为 Cluster 410 wiring 提供权威来源。migration/schema/readiness/repository - 定向 63/63 通过。ADR 暂保持 Proposed;下一阶段继续完成 validated S3 HEAD、ETag/VersionId 条件删除、bounded lifecycle、MinIO - failure matrix 与 PostgreSQL HA failover,全部完整门通过后才接受。 +- D-291/ADR-0379(已接受) + Cluster Run Attempt 日志 retention 已完成多副本纵向闭环,且没有新增 package:共享 claim contract 位于 Runtime Core 既有 Run + log-retention 目录;PostgreSQL v54 提供 durable control、immutable tombstone、terminal remote Worker candidate index 与最小权限 + `ql3_runtime` authority。副本以短 `READ COMMITTED`/`SKIP LOCKED` 事务取得 owner/token/version/expiry fence,事务提交后才做 + validated S3 HEAD;versioned object 与 upload 临时对象按精确 VersionId 删除,unversioned object 按 ETag `If-Match` 删除,412 + identity drift 失败关闭,删除响应丢失由下一租约经 HEAD absent 收敛为 `already_absent`。第二个短事务重验数据库时钟、完整 claim 与 + Run/Attempt identity,原子写 exact tombstone 并清除 control。Cluster application 已接入有界 claim/delete、指数退避/manual、单 + `unref` timer、共享 wall-clock abort 和 reverse-stop drain;只有 Worker ingress/S3 激活时才装配。生产日志读取在对象存储前后检查 + tombstone,稳定返回 410,而普通 missing 保持 503。 + + 真实 MinIO 在强制 SSE-S3 下通过 versioning disabled/enabled 条件删除,versioned 路径最终零旧版本、零 delete marker;PostgreSQL + 18.4 arm64 physical HA 在 timeline `1→2` 下通过 113 gates,证明旧主 claim 已同步复制、旧 settlement 被 fenced、新主以 version 2 + 接管并把 control/tombstone 收敛为 0/1,报告 SHA-256 为 + `4be3053fc1af9ad6304715f5398292ba9a31ec5b3d49f64787510e2f2645ec5f`。最终 `cluster-control` 230 pass/2 外部条件 skip,完整 + 18-package clean build/test 退出 0,backend 1163 pass/2 skip/0 fail;workspace 保持 18 package/1054 source/1036 nested/18 + reviewed root entry,`cluster-control` 为 51 source/49 nested/2 root binary entry,无 single-source/shallow package。dependency、 + package、Edge import 边界零 finding,121 个 Edge 实际 imported module 不含 PostgreSQL、AWS SDK 或 Cluster package;14 档 Local + Profile artifact 与 Local image static audit 保持 compatible。 - D-290/ADR-0378(已接受) Local Run Attempt 日志 retention 已形成真实纵向切片:Runtime Core 增加精确 identity、canonical SHA-256 的 immutable retirement record、容量压力策略、有界 page/delete budget 与 durable cursor;日志读取在存储前检查 tombstone,并在 missing 后二次 diff --git a/docs/adr/ADR-0379-cluster-run-attempt-log-retention.md b/docs/adr/ADR-0379-cluster-run-attempt-log-retention.md index 877191b8..c2187a7c 100644 --- a/docs/adr/ADR-0379-cluster-run-attempt-log-retention.md +++ b/docs/adr/ADR-0379-cluster-run-attempt-log-retention.md @@ -1,6 +1,6 @@ # ADR-0379:Cluster Run Attempt 日志多副本保留与条件删除 -- 状态:Proposed(PostgreSQL authority 已实现,S3/lifecycle/HA 验收待完成) +- 状态:Accepted - 日期:2026-08-12 - 关联 RFC:QL-RFC-0001 D-291 - 前置决策:ADR-0026、ADR-0027、ADR-0377、ADR-0378 @@ -34,14 +34,17 @@ Cluster 还必须保持与低配 Local 部署的物理隔离:本能力不得 1. validated HEAD 校验 content type、metadata identity、checksum、byte length,并取得 ETag 与可用的 VersionId; 2. 对 versioned object 使用精确 VersionId,对未版本化对象使用 ETag `If-Match` 条件删除; -3. 412/对象身份变化失败关闭并进入 retry/manual,不得删除新对象; -4. 删除成功或 HEAD 已不存在后,开启第二个短 PostgreSQL 事务; -5. 事务重验 owner/token/version/expiry、terminal Run/Attempt 与 immutable identity,插入 exact tombstone 后删除 control; -6. 删除响应丢失时,lease 过期后的新 claim 以 HEAD absent 写入 `already_absent`,最终收敛。 +3. upload promotion 的临时对象也必须先取得精确 VersionId 或 ETag authority,再按同一规则清理;无法证明临时对象身份时宁可留下可诊断对象,不得执行无条件删除; +4. 412/对象身份变化失败关闭并进入 retry/manual,不得删除新对象; +5. 删除成功或 HEAD 已不存在后,开启第二个短 PostgreSQL 事务; +6. 事务重验 owner/token/version/expiry、terminal Run/Attempt 与 immutable identity,插入 exact tombstone 后删除 control; +7. 删除响应丢失时,lease 过期后的新 claim 以 HEAD absent 写入 `already_absent`,最终收敛。 ### 4. 有界调度与退避 -每轮 claim 不超过 16 条,lease 范围 5 秒至 5 分钟;retry delay 最大 24 小时。副本不得持有跨 sweep cursor,也不得为每个 Attempt 建 timer。调度复用 Cluster control application 既有 lifecycle,并受每轮 claim 数、删除数和 wall-clock budget 共同限制。`artifact_unavailable`、`artifact_integrity_mismatch`、`retirement_record_unavailable` 是持久化失败分类;达到策略阈值后转 manual,避免坏对象形成热循环。 +每轮 claim 不超过 16 条,lease 范围 5 秒至 5 分钟;retry delay 最大 24 小时。副本不得持有跨 sweep cursor,也不得为每个 Attempt 建 timer。`ClusterRunAttemptLogRetentionLifecycle` 只拥有一个 `unref` timer,重叠 tick 合并为同一轮,单轮共享一个 wall-clock `AbortSignal`;停机先 abort、再在上限内 drain,且位于 application reverse-stop 的最前端。`artifact_unavailable`、`artifact_integrity_mismatch`、`retirement_record_unavailable` 是持久化失败分类;达到策略阈值后转 manual,避免坏对象形成热循环。 + +生产配置使用显式 `QL3_CLUSTER_LOG_RETENTION_*` 边界控制 retention、claim、lease、cycle budget、retry、manual threshold、cadence 与 stop timeout。能力默认启用,但只有 Worker ingress/S3 store 已激活时才装配,不为没有远端日志的 Cluster 进程增加 timer。 ### 5. 读取收敛 @@ -49,20 +52,16 @@ PostgreSQL claim repository 同时实现 retention state reader。Cluster 日志 ## 阶段验收 -第一阶段已完成 PostgreSQL authority:共享 claim contract、v54 migration、typed Drizzle/schema/readiness、最小权限、短事务 claim、完整 lease fence、retry/manual settlement、exact tombstone finalize/replay、tombstone state read。定向 63 项 migration/schema/readiness/repository 门通过;`runtime-core` 498 项全通过,`cluster-postgres` 302 项通过、1 项条件跳过,`cluster-control` 216 项通过、2 项条件跳过。 +本 ADR 已完成从 PostgreSQL ownership 到对象删除、生产 lifecycle 和读取 410 的纵向闭环: -阶段收口还通过 18 个 QL3 workspace package 的完整 build/test 门,以及后端兼容回归 1163 项通过、2 项条件跳过、0 失败/取消。v54 readiness 引入的两张表已同步进入 `cluster-control` 测试数据库的最小权限 fixture,避免旧 v53 fixture 把正常启动误判为 `runtime_role_invalid`。 +- PostgreSQL v54 authority 已覆盖 typed schema/readiness、最小权限、短事务 claim、完整 lease fence、retry/manual settlement、exact tombstone finalize/replay 与 state read;`cluster-postgres` 为 302 pass/1 条件 skip。 +- S3 单元矩阵覆盖 versioned/unversioned 条件删除、412 identity drift、对象已不存在、删除响应丢失收敛、malformed HEAD,以及临时对象的 VersionId/ETag 精确清理;真实 MinIO 在强制 SSE-S3 下分别通过 versioning disabled/enabled,versioned 路径最终为零旧版本、零 delete marker。 +- production composition 已把 reader、retirement store、coordinator 与 lifecycle 接入唯一 Cluster application;生产 HTTP 读取可由 durable tombstone 返回 410。`cluster-control` 为 230 pass/2 外部条件 skip。 +- PostgreSQL 18.4 arm64 physical HA 在 timeline `1→2` 下通过 113 gates:旧主 claim 已 `remote_apply` 到 standby,提升且同步冗余恢复后旧 owner settlement 被 fenced,新主以 claim version 2 接管,随后原子写唯一 tombstone 并把 control count 收敛为 0。报告 SHA-256 为 `4be3053fc1af9ad6304715f5398292ba9a31ec5b3d49f64787510e2f2645ec5f`。 +- 18 个 QL3 workspace package 的 clean build/test 门退出 0;后端兼容回归为 1163 pass/2 条件 skip/0 fail。package boundary 保持 18 package、1054 source、1036 nested、18 个受审 root entry,`singleSourcePackages=[]`、`shallowSourcePackages=[]`;`cluster-control` 为 51 source,其中 49 个 nested、2 个 root binary entry。 +- dependency 与 Edge import audit 零 finding;Edge closure 的 121 个实际 imported module 不包含 PostgreSQL、AWS SDK 或 Cluster package。14 档 Local Profile artifact 与 Local image static audit 保持 compatible,因此本能力不会改变路由设备的默认数据库、连接、timer 或对象存储负担。 -结构门保持 18 个 workspace package,`singleSourcePackages=[]`、`shallowSourcePackages=[]`。新 repository 只从明确的 Cluster runtime entrypoint 导出,不扩大 package root;PostgreSQL migration append-only `ordered_ledger` 的 reviewed hard cap 随 pg-0055 由 57 精确推进到 58,不把版本账本伪拆为子目录。 - -ADR 保持 Proposed,只有以下剩余项全部完成后才转 Accepted: - -1. S3 validated HEAD + ETag/VersionId 条件删除和失败矩阵; -2. Cluster service 的 bounded claim/delete/backoff/manual 策略; -3. production composition、lifecycle drain 与读取 410; -4. MinIO versioned/unversioned 集成、响应丢失重放; -5. PostgreSQL 18 HA failover 中 lease takeover/tombstone 收敛; -6. 完整 package/backend/boundary/Profile/image gates,证明 Local closure 无 PostgreSQL/AWS SDK 回归。 +这些证据满足原六项收口条件,ADR 转为 Accepted。真实生产对象存储厂商矩阵、Kubernetes 多节点分区、基础设施 STONITH 与长期容量基准仍属于 Release Gate,不由本地 MinIO/PostgreSQL Docker 合约代替。 ## 被否决的替代方案 diff --git a/docs/adr/README.md b/docs/adr/README.md index 0a8e3d40..372c7b32 100644 --- a/docs/adr/README.md +++ b/docs/adr/README.md @@ -382,7 +382,7 @@ | [ADR-0376](./ADR-0376-policy-and-digest-fenced-task-start.md) | Policy 与 digest fenced 的 Task Start | Accepted | | [ADR-0377](./ADR-0377-profile-aware-run-attempt-log-range-read.md) | Profile-aware Run Attempt 日志 Range 读取 | Accepted | | [ADR-0378](./ADR-0378-local-run-attempt-log-retention-and-tombstones.md) | Local Run Attempt 日志有界保留与 durable tombstone | Accepted | -| [ADR-0379](./ADR-0379-cluster-run-attempt-log-retention.md) | Cluster Run Attempt 日志多副本保留与条件删除 | Proposed(PostgreSQL authority 已完成) | +| [ADR-0379](./ADR-0379-cluster-run-attempt-log-retention.md) | Cluster Run Attempt 日志多副本保留与条件删除 | Accepted | ## 规则 diff --git a/packages/ql3-cluster-control/src/application-runtime/application.ts b/packages/ql3-cluster-control/src/application-runtime/application.ts index 6895f3c9..8cd10aa0 100644 --- a/packages/ql3-cluster-control/src/application-runtime/application.ts +++ b/packages/ql3-cluster-control/src/application-runtime/application.ts @@ -13,6 +13,7 @@ import { type ClusterControlAssemblyInput, type ClusterControlRecoveryRuntimeOptions, type ClusterRunCancellationConvergenceRuntimeOptions, + type ClusterRunAttemptLogRetentionRuntimeOptions, type ClusterSchedulerRuntimeOptions, type ClusterWorkerRuntimeDependencies, } from './clusterControlRuntime'; @@ -39,6 +40,7 @@ export interface ClusterControlApplicationOptions { readonly recovery?: ClusterControlRecoveryRuntimeOptions; readonly scheduler?: ClusterSchedulerRuntimeOptions; readonly cancellationConvergence?: ClusterRunCancellationConvergenceRuntimeOptions; + readonly logRetention?: ClusterRunAttemptLogRetentionRuntimeOptions; readonly workerRuntime?: ClusterWorkerRuntimeDependencies; readonly openDatabase: OpenPostgresDatabase; readonly availability: ClusterControlAvailabilitySource; @@ -88,6 +90,9 @@ function inactiveBootstrap( ...(options.cancellationConvergence === undefined ? {} : { cancellationConvergence: options.cancellationConvergence }), + ...(options.logRetention === undefined + ? {} + : { logRetention: options.logRetention }), ...(options.workerRuntime === undefined ? {} : { workerRuntime: options.workerRuntime }), @@ -163,6 +168,9 @@ export async function startClusterControlApplication( ...(options.cancellationConvergence === undefined ? {} : { cancellationConvergence: options.cancellationConvergence }), + ...(options.logRetention === undefined + ? {} + : { logRetention: options.logRetention }), ...(options.workerRuntime === undefined ? {} : { workerRuntime: options.workerRuntime }), diff --git a/packages/ql3-cluster-control/src/application-runtime/clusterControlRuntime.ts b/packages/ql3-cluster-control/src/application-runtime/clusterControlRuntime.ts index db3a46af..1b81e316 100644 --- a/packages/ql3-cluster-control/src/application-runtime/clusterControlRuntime.ts +++ b/packages/ql3-cluster-control/src/application-runtime/clusterControlRuntime.ts @@ -30,6 +30,17 @@ import { type ClusterRunCancellationConvergenceCycleResult, } from '@qinglong/runtime-core'; import type { ClusterRunCancellationRepository } from '@qinglong/runtime-core/cluster-run-cancellation'; +import { + MAX_CLUSTER_RUN_ATTEMPT_LOG_RETENTION_CLAIMS, + MAX_CLUSTER_RUN_ATTEMPT_LOG_RETENTION_LEASE_MS, + MAX_CLUSTER_RUN_ATTEMPT_LOG_RETENTION_RETRY_DELAY_MS, + MIN_CLUSTER_RUN_ATTEMPT_LOG_RETENTION_LEASE_MS, +} from '@qinglong/runtime-core/cluster-run-attempt-log-retention'; +import { + MAX_RUN_ATTEMPT_LOG_RETENTION_MS, + MIN_RUN_ATTEMPT_LOG_RETENTION_MS, + type RunAttemptLogRetentionStateReader, +} from '@qinglong/runtime-core/run-attempt-log-retention'; import type { ProjectRunListReader } from '@qinglong/runtime-core/project-run-list'; import type { ClusterScheduleStore } from '@qinglong/runtime-core/cluster-scheduler'; import type { TaskDefinitionSource } from '@qinglong/runtime-core/task-definition'; @@ -60,6 +71,7 @@ import { PostgresSecurityAuditRepository, PostgresRunRepository, PostgresWorkerExecutionAttestationRepository, + PostgresRunAttemptLogRetentionClaimRepository, PostgresTaskDefinitionSource, PostgresTaskExecutionRevisionSource, PostgresTriggerSource, @@ -104,6 +116,12 @@ import { import { ClusterWorkflowSchedulerCoordinator } from '../scheduling/workflowScheduler'; import { ClusterRuntimeSchedulerCoordinator } from '../scheduling/runtimeScheduler'; import { ClusterRunCancellationConvergenceLifecycle } from '../run/runCancellationLifecycle'; +import { + ClusterRunAttemptLogRetentionCoordinator, + ClusterRunAttemptLogRetentionLifecycle, + type ClusterRunAttemptLogRetirementStore, + type ClusterRunAttemptLogRetentionCycleSummary, +} from '../run/runAttemptLogRetentionLifecycle'; import type { TaskStartRepository } from '@qinglong/runtime-core/task-start'; import { createClusterWorkerRuntimePort, @@ -131,6 +149,7 @@ export interface ClusterControlAssemblyInput { readonly evidence: ClusterControlReadinessEvidence; readonly policies: ProjectPolicyRepository; readonly runs: RunRepository & ProjectRunListReader; + readonly runAttemptLogRetention: RunAttemptLogRetentionStateReader; readonly runCancellation: ClusterRunCancellationRepository; readonly taskStart: TaskStartRepository; readonly taskDefinitions: TaskDefinitionSource; @@ -178,6 +197,24 @@ export interface ClusterRunCancellationConvergenceRuntimeOptions { ) => void | Promise; } +export interface ClusterRunAttemptLogRetentionRuntimeOptions { + readonly store: ClusterRunAttemptLogRetirementStore; + readonly ownerId?: string; + readonly retentionMs?: number; + readonly claimLimit?: number; + readonly leaseMs?: number; + readonly maximumCycleMs?: number; + readonly retryBaseMs?: number; + readonly retryMaximumMs?: number; + readonly maximumFailures?: number; + readonly intervalMs?: number; + readonly stopTimeoutMs?: number; + readonly onDiagnostic?: ( + error: unknown, + summary?: Readonly, + ) => void | Promise; +} + export interface ClusterControlBootstrapOptions { readonly enabled?: boolean; readonly profile: DeploymentProfile; @@ -185,6 +222,7 @@ export interface ClusterControlBootstrapOptions { readonly recovery?: ClusterControlRecoveryRuntimeOptions; readonly scheduler?: ClusterSchedulerRuntimeOptions; readonly cancellationConvergence?: ClusterRunCancellationConvergenceRuntimeOptions; + readonly logRetention?: ClusterRunAttemptLogRetentionRuntimeOptions; readonly workerRuntime?: ClusterWorkerRuntimeDependencies; readonly openDatabase: OpenPostgresDatabase; readonly create: ( @@ -223,6 +261,21 @@ interface PreparedCancellationConvergenceRuntime { readonly onDiagnostic?: ClusterRunCancellationConvergenceRuntimeOptions['onDiagnostic']; } +interface PreparedLogRetentionRuntime { + readonly store: ClusterRunAttemptLogRetirementStore; + readonly ownerId: string; + readonly retentionMs: number; + readonly claimLimit: number; + readonly leaseMs: number; + readonly maximumCycleMs: number; + readonly retryBaseMs: number; + readonly retryMaximumMs: number; + readonly maximumFailures: number; + readonly intervalMs: number; + readonly stopTimeoutMs: number; + readonly onDiagnostic?: ClusterRunAttemptLogRetentionRuntimeOptions['onDiagnostic']; +} + function boundedInteger( name: string, value: number | undefined, @@ -439,6 +492,116 @@ function prepareCancellationConvergenceRuntime( }); } +function prepareLogRetentionRuntime( + options: ClusterRunAttemptLogRetentionRuntimeOptions | undefined, + fallbackOwnerId: string, +): PreparedLogRetentionRuntime | undefined { + if (options === undefined) return undefined; + const allowedKeys = new Set([ + 'claimLimit', + 'intervalMs', + 'leaseMs', + 'maximumCycleMs', + 'maximumFailures', + 'onDiagnostic', + 'ownerId', + 'retentionMs', + 'retryBaseMs', + 'retryMaximumMs', + 'stopTimeoutMs', + 'store', + ]); + if ( + !options || + typeof options !== 'object' || + Array.isArray(options) || + Object.keys(options).some((key) => !allowedKeys.has(key)) || + typeof options.store?.retire !== 'function' || + (options.onDiagnostic !== undefined && + typeof options.onDiagnostic !== 'function') + ) { + throw new TypeError( + 'Cluster Run Attempt log retention configuration is invalid', + ); + } + const ownerId = options.ownerId ?? fallbackOwnerId; + if (!/^[A-Za-z0-9][A-Za-z0-9._:-]{0,127}$/.test(ownerId)) { + throw new TypeError('Cluster Run Attempt log retention ownerId is invalid'); + } + const leaseMs = boundedInteger( + 'Cluster Run Attempt log retention lease', + options.leaseMs, + 30_000, + MIN_CLUSTER_RUN_ATTEMPT_LOG_RETENTION_LEASE_MS, + MAX_CLUSTER_RUN_ATTEMPT_LOG_RETENTION_LEASE_MS, + ); + const retryBaseMs = boundedInteger( + 'Cluster Run Attempt log retention retry base', + options.retryBaseMs, + 5_000, + 0, + MAX_CLUSTER_RUN_ATTEMPT_LOG_RETENTION_RETRY_DELAY_MS, + ); + return Object.freeze({ + store: options.store, + ownerId, + retentionMs: boundedInteger( + 'Cluster Run Attempt log retention duration', + options.retentionMs, + 30 * 24 * 60 * 60_000, + MIN_RUN_ATTEMPT_LOG_RETENTION_MS, + MAX_RUN_ATTEMPT_LOG_RETENTION_MS, + ), + claimLimit: boundedInteger( + 'Cluster Run Attempt log retention claim limit', + options.claimLimit, + 4, + 1, + MAX_CLUSTER_RUN_ATTEMPT_LOG_RETENTION_CLAIMS, + ), + leaseMs, + maximumCycleMs: boundedInteger( + 'Cluster Run Attempt log retention cycle budget', + options.maximumCycleMs, + 10_000, + 100, + leaseMs - 500, + ), + retryBaseMs, + retryMaximumMs: boundedInteger( + 'Cluster Run Attempt log retention retry maximum', + options.retryMaximumMs, + 60 * 60_000, + retryBaseMs, + MAX_CLUSTER_RUN_ATTEMPT_LOG_RETENTION_RETRY_DELAY_MS, + ), + maximumFailures: boundedInteger( + 'Cluster Run Attempt log retention failure limit', + options.maximumFailures, + 8, + 1, + 32, + ), + intervalMs: boundedInteger( + 'Cluster Run Attempt log retention interval', + options.intervalMs, + 60_000, + 1_000, + 24 * 60 * 60_000, + ), + stopTimeoutMs: boundedInteger( + 'Cluster Run Attempt log retention stop timeout', + options.stopTimeoutMs, + 10_000, + 100, + 30_000, + ), + ...(options.onDiagnostic === undefined + ? {} + : { onDiagnostic: options.onDiagnostic }), + }); +} + function readinessEvidence( report: PostgresSchemaReadinessReport, ): ClusterControlReadinessEvidence { @@ -463,6 +626,7 @@ export async function bootstrapClusterControlRuntime( let cancellationConvergenceRuntime: | PreparedCancellationConvergenceRuntime | undefined; + let logRetentionRuntime: PreparedLogRetentionRuntime | undefined; let recoveryRegistry: ClusterControlRecoveryEvidenceRegistry | undefined; if ((options.enabled ?? false) && options.profile === 'cluster-control') { assertClusterControlApiCredentialPepper(options.apiCredentialPepper ?? ''); @@ -474,6 +638,10 @@ export async function bootstrapClusterControlRuntime( cancellationConvergenceRuntime = prepareCancellationConvergenceRuntime( options.cancellationConvergence, ); + logRetentionRuntime = prepareLogRetentionRuntime( + options.logRetention, + recoveryRuntime.ownerId, + ); } let database: PostgresDatabaseResource | undefined; let closePromise: Promise | undefined; @@ -576,6 +744,8 @@ export async function bootstrapClusterControlRuntime( ); const schedules = new PostgresClusterScheduleRepository(database.pool); const runs = new PostgresRunRepository(database.pool); + const runAttemptLogRetention = + new PostgresRunAttemptLogRetentionClaimRepository(database.pool); const trustedToolStorage: ClusterTrustedToolStorage = Object.freeze({ invocationArtifacts: new PostgresToolInvocationArtifactRepository( database.pool, @@ -653,6 +823,32 @@ export async function bootstrapClusterControlRuntime( }), }, ); + const logRetentionLifecycle = + logRetentionRuntime === undefined + ? undefined + : new ClusterRunAttemptLogRetentionLifecycle( + new ClusterRunAttemptLogRetentionCoordinator( + runAttemptLogRetention, + logRetentionRuntime.store, + { + ownerId: logRetentionRuntime.ownerId, + retentionMs: logRetentionRuntime.retentionMs, + claimLimit: logRetentionRuntime.claimLimit, + leaseMs: logRetentionRuntime.leaseMs, + maximumCycleMs: logRetentionRuntime.maximumCycleMs, + retryBaseMs: logRetentionRuntime.retryBaseMs, + retryMaximumMs: logRetentionRuntime.retryMaximumMs, + maximumFailures: logRetentionRuntime.maximumFailures, + }, + ), + { + intervalMs: logRetentionRuntime.intervalMs, + stopTimeoutMs: logRetentionRuntime.stopTimeoutMs, + ...(logRetentionRuntime.onDiagnostic === undefined + ? {} + : { onDiagnostic: logRetentionRuntime.onDiagnostic }), + }, + ); const runCancellation = new PostgresClusterRunCancellationRepository( database.pool, ); @@ -665,6 +861,7 @@ export async function bootstrapClusterControlRuntime( ), policies: new PostgresProjectPolicyRepository(database.pool), runs, + runAttemptLogRetention, runCancellation, taskStart, taskDefinitions: new PostgresTaskDefinitionSource(database.pool), @@ -739,6 +936,7 @@ export async function bootstrapClusterControlRuntime( if (!(await application.startLifecycles())) return false; schedulerLifecycle.start(); cancellationConvergenceLifecycle.start(); + logRetentionLifecycle?.start(); return true; }, installAdmission: () => application.installAdmission(), @@ -746,14 +944,21 @@ export async function bootstrapClusterControlRuntime( recoveryRegistry?.dispose(); let schedulerStatus: 'stopped' | 'timed_out' = 'stopped'; let cancellationStatus: 'stopped' | 'timed_out' = 'stopped'; + let logRetentionStatus: 'stopped' | 'timed_out' = 'stopped'; let applicationStatus: ClusterControlStopResult = 'stopped'; let primaryError: unknown; + try { + logRetentionStatus = + (await logRetentionLifecycle?.stopAndDrain()) ?? 'stopped'; + } catch (error) { + primaryError = error; + } try { cancellationStatus = ( await cancellationConvergenceLifecycle.stopAndDrain() ).status; } catch (error) { - primaryError = error; + primaryError ??= error; } try { schedulerStatus = (await schedulerLifecycle.stopAndDrain()) @@ -768,6 +973,7 @@ export async function bootstrapClusterControlRuntime( } if (primaryError) throw primaryError; return cancellationStatus === 'timed_out' || + logRetentionStatus === 'timed_out' || schedulerStatus === 'timed_out' || applicationStatus === 'timed_out' ? 'timed_out' diff --git a/packages/ql3-cluster-control/src/application-runtime/productionApplication.ts b/packages/ql3-cluster-control/src/application-runtime/productionApplication.ts index bcf2e2b8..916748c4 100644 --- a/packages/ql3-cluster-control/src/application-runtime/productionApplication.ts +++ b/packages/ql3-cluster-control/src/application-runtime/productionApplication.ts @@ -200,6 +200,7 @@ export function createProductionClusterControlApplicationStack( createClusterControlRunAttemptLogReadRoute( input.runs, input.workerRuntime?.runAttemptLogRead, + input.runAttemptLogRetention, ), createClusterControlRunCancellationRoute( input.runCancellation, diff --git a/packages/ql3-cluster-control/src/artifact/s3ArtifactStore.ts b/packages/ql3-cluster-control/src/artifact/s3ArtifactStore.ts index 9025f9f4..a9d0c54c 100644 --- a/packages/ql3-cluster-control/src/artifact/s3ArtifactStore.ts +++ b/packages/ql3-cluster-control/src/artifact/s3ArtifactStore.ts @@ -19,6 +19,12 @@ import { type RunAttemptLogReadIdentity, type RunAttemptLogReadRange, } from '@qinglong/runtime-core/run-attempt-log-read'; +import { + normalizeRunAttemptLogRetentionCandidate, + type RunAttemptLogRetentionCandidate, + type RunAttemptLogRetirementStore, + type RunAttemptLogRetirementStoreResult, +} from '@qinglong/runtime-core/run-attempt-log-retention'; import { MAX_REMOTE_WORKER_ARTIFACT_BYTES, REMOTE_WORKER_ARTIFACT_CONTENT_TYPE, @@ -124,6 +130,7 @@ type ArtifactAuthority = Readonly<{ type StoredArtifactHead = Readonly<{ receipt: Readonly; eTag?: string; + versionId?: string; }>; type NormalizedStorageCommand = ArtifactAuthority & @@ -484,6 +491,18 @@ function canonicalETag(value: unknown): string { return value; } +function canonicalVersionId(value: unknown): string { + if ( + typeof value !== 'string' || + value.length < 1 || + Buffer.byteLength(value, 'utf8') > 1024 || + /[\u0000-\u001f\u007f]/.test(value) + ) { + throw new S3ClusterRemoteWorkerArtifactStoreError('integrity_mismatch'); + } + return value; +} + function assertRangeMetadata( authority: ArtifactAuthority, receipt: Readonly, @@ -580,6 +599,20 @@ function isNotFound(error: unknown): boolean { ); } +function isPreconditionFailed(error: unknown): boolean { + if (!error || typeof error !== 'object') return false; + const value = error as { + name?: unknown; + Code?: unknown; + $metadata?: { httpStatusCode?: unknown }; + }; + return ( + value.name === 'PreconditionFailed' || + value.Code === 'PreconditionFailed' || + value.$metadata?.httpStatusCode === 412 + ); +} + function requestOptions( signal?: AbortSignal, ): { abortSignal: AbortSignal } | undefined { @@ -655,10 +688,11 @@ class ArtifactContentDigest { /** * Shared immutable S3 adapter. A unique temporary upload is checksummed first, * then promoted by one destination-conditional server-side copy. Permanent - * objects are never overwritten or deleted by this adapter. + * objects are never overwritten and are retired only after an identity-checked + * HEAD followed by a VersionId- or ETag-fenced delete. */ export class S3ClusterRemoteWorkerArtifactStore - implements ClusterRemoteWorkerArtifactStore + implements ClusterRemoteWorkerArtifactStore, RunAttemptLogRetirementStore { private readonly options: PreparedOptions; @@ -751,6 +785,75 @@ export class S3ClusterRemoteWorkerArtifactStore } } + async retire( + rawCandidate: Readonly, + signal?: AbortSignal, + ): Promise> { + const candidate = normalizeRunAttemptLogRetentionCandidate(rawCandidate); + if (candidate.executorType !== 'remote_worker') { + throw new S3ClusterRemoteWorkerArtifactStoreError('integrity_mismatch'); + } + const authority = Object.freeze({ + projectId: candidate.projectId, + runId: candidate.runId, + attemptId: candidate.attemptId, + logArtifactId: candidate.logArtifactId, + }); + const stored = await this.head(authority, signal); + if (!stored) { + return Object.freeze({ + disposition: 'already_absent' as const, + byteLength: 0, + truncation: Object.freeze({ truncated: 'unknown' as const }), + }); + } + const versionId = + stored.versionId === undefined + ? undefined + : canonicalVersionId(stored.versionId); + const eTag = + versionId === undefined ? canonicalETag(stored.eTag) : undefined; + try { + await this.options.client.send( + new DeleteObjectCommand({ + Bucket: this.options.bucket, + Key: finalObjectKey(this.options.prefix, authority), + ...(versionId === undefined + ? { IfMatch: eTag } + : { VersionId: versionId }), + ...(this.options.expectedBucketOwner === undefined + ? {} + : { ExpectedBucketOwner: this.options.expectedBucketOwner }), + }), + requestOptions(signal), + ); + } catch (error) { + if (isNotFound(error)) { + return Object.freeze({ + disposition: 'already_absent' as const, + byteLength: 0, + truncation: Object.freeze({ truncated: 'unknown' as const }), + }); + } + if (isPreconditionFailed(error)) { + throw new S3ClusterRemoteWorkerArtifactStoreError( + 'integrity_mismatch', + { cause: error }, + ); + } + throw new S3ClusterRemoteWorkerArtifactStoreError('unavailable', { + cause: error, + }); + } + return Object.freeze({ + disposition: 'deleted' as const, + byteLength: stored.receipt.byteLength, + truncation: Object.freeze({ + truncated: stored.receipt.truncated ?? ('unknown' as const), + }), + }); + } + private async head( authority: ArtifactAuthority, signal?: AbortSignal, @@ -771,6 +874,9 @@ export class S3ClusterRemoteWorkerArtifactStore return Object.freeze({ receipt: parseStoredReceipt(authority, output), ...(output.ETag === undefined ? {} : { eTag: output.ETag }), + ...(output.VersionId === undefined + ? {} + : { versionId: output.VersionId }), }); } catch (error) { if (isNotFound(error)) return undefined; @@ -811,7 +917,9 @@ export class S3ClusterRemoteWorkerArtifactStore this.options.createTemporaryId, ); const temporaryOwner = temporaryOwnershipDigest(); - let temporaryOwned = false; + let temporaryDeleteAuthority: + | Readonly<{ readonly eTag: string; readonly versionId?: string }> + | undefined; let result: Readonly | undefined; let primaryError: unknown; try { @@ -837,21 +945,19 @@ export class S3ClusterRemoteWorkerArtifactStore }), requestOptions(signal), ); - temporaryOwned = true; } catch (error) { if (!digest.isComplete()) throw error; } finally { body.destroy(); } const sha256 = digest.digest(); - await this.assertTemporaryObject( + temporaryDeleteAuthority = await this.assertTemporaryObject( temporaryKey, temporaryOwner, command.byteLength, sha256, signal, ); - temporaryOwned = true; let copied = false; try { @@ -898,12 +1004,15 @@ export class S3ClusterRemoteWorkerArtifactStore primaryError = error; } - if (temporaryOwned) { + if (temporaryDeleteAuthority !== undefined) { try { await this.options.client.send( new DeleteObjectCommand({ Bucket: this.options.bucket, Key: temporaryKey, + ...(temporaryDeleteAuthority.versionId === undefined + ? { IfMatch: temporaryDeleteAuthority.eTag } + : { VersionId: temporaryDeleteAuthority.versionId }), ...(this.options.expectedBucketOwner === undefined ? {} : { ExpectedBucketOwner: this.options.expectedBucketOwner }), @@ -937,7 +1046,7 @@ export class S3ClusterRemoteWorkerArtifactStore byteLength: number, sha256: string, signal?: AbortSignal, - ): Promise { + ): Promise> { let output; try { output = await this.options.client.send( @@ -965,5 +1074,14 @@ export class S3ClusterRemoteWorkerArtifactStore ) { throw new S3ClusterRemoteWorkerArtifactStoreError('integrity_mismatch'); } + const eTag = canonicalETag(output.ETag); + const versionId = + output.VersionId === undefined + ? undefined + : canonicalVersionId(output.VersionId); + return Object.freeze({ + eTag, + ...(versionId === undefined ? {} : { versionId }), + }); } } diff --git a/packages/ql3-cluster-control/src/artifact/workerArtifactBinding.ts b/packages/ql3-cluster-control/src/artifact/workerArtifactBinding.ts index 87bf8150..4251d57e 100644 --- a/packages/ql3-cluster-control/src/artifact/workerArtifactBinding.ts +++ b/packages/ql3-cluster-control/src/artifact/workerArtifactBinding.ts @@ -5,9 +5,11 @@ import { } from './s3ArtifactStore'; import type { ClusterRemoteWorkerArtifactStore } from '../remote-execution/remoteWorkerCompletionService'; import type { ClusterWorkerArtifactS3Config } from '../worker-ingress/workerIngressConfig'; +import type { ClusterRunAttemptLogRetirementStore } from '../run/runAttemptLogRetentionLifecycle'; export interface ClusterWorkerArtifactBinding { - readonly store: ClusterRemoteWorkerArtifactStore; + readonly store: ClusterRemoteWorkerArtifactStore & + ClusterRunAttemptLogRetirementStore; close(): Promise; } @@ -19,9 +21,7 @@ export function createClusterWorkerArtifactBinding( } const client = createS3ClusterRemoteWorkerArtifactClient({ region: config.region, - ...(config.endpoint === undefined - ? {} - : { endpoint: config.endpoint }), + ...(config.endpoint === undefined ? {} : { endpoint: config.endpoint }), forcePathStyle: config.forcePathStyle, }); const store = new S3ClusterRemoteWorkerArtifactStore({ diff --git a/packages/ql3-cluster-control/src/production-process/config.ts b/packages/ql3-cluster-control/src/production-process/config.ts index d2fca8e7..36a3c59a 100644 --- a/packages/ql3-cluster-control/src/production-process/config.ts +++ b/packages/ql3-cluster-control/src/production-process/config.ts @@ -12,6 +12,16 @@ import { } from '@qinglong/cluster-postgres/runtime'; import { ClusterControlAvailabilityFence } from '../database/availability'; import type { ClusterControlHttpSurfaceOptions } from '../transport/httpSurface'; +import { + MAX_CLUSTER_RUN_ATTEMPT_LOG_RETENTION_CLAIMS, + MAX_CLUSTER_RUN_ATTEMPT_LOG_RETENTION_LEASE_MS, + MAX_CLUSTER_RUN_ATTEMPT_LOG_RETENTION_RETRY_DELAY_MS, + MIN_CLUSTER_RUN_ATTEMPT_LOG_RETENTION_LEASE_MS, +} from '@qinglong/runtime-core/cluster-run-attempt-log-retention'; +import { + MAX_RUN_ATTEMPT_LOG_RETENTION_MS, + MIN_RUN_ATTEMPT_LOG_RETENTION_MS, +} from '@qinglong/runtime-core/run-attempt-log-retention'; export type ClusterControlEnvironment = Readonly< Record @@ -33,6 +43,20 @@ export interface EnabledClusterControlConfig { readonly security: Readonly<{ apiCredentialPepper: string; }>; + readonly logRetention: + | Readonly<{ readonly enabled: false }> + | Readonly<{ + readonly enabled: true; + readonly retentionMs: number; + readonly claimLimit: number; + readonly leaseMs: number; + readonly maximumCycleMs: number; + readonly retryBaseMs: number; + readonly retryMaximumMs: number; + readonly maximumFailures: number; + readonly intervalMs: number; + readonly stopTimeoutMs: number; + }>; } export type ClusterControlConfig = @@ -222,6 +246,82 @@ function apiCredentialPepper(environment: ClusterControlEnvironment): string { return value; } +function logRetentionConfig( + environment: ClusterControlEnvironment, +): EnabledClusterControlConfig['logRetention'] { + if (!booleanValue(environment, 'QL3_CLUSTER_LOG_RETENTION_ENABLED', true)) { + return Object.freeze({ enabled: false as const }); + } + const leaseMs = integerValue( + environment, + 'QL3_CLUSTER_LOG_RETENTION_LEASE_MS', + 30_000, + MIN_CLUSTER_RUN_ATTEMPT_LOG_RETENTION_LEASE_MS, + MAX_CLUSTER_RUN_ATTEMPT_LOG_RETENTION_LEASE_MS, + ); + const retryBaseMs = integerValue( + environment, + 'QL3_CLUSTER_LOG_RETENTION_RETRY_BASE_MS', + 5_000, + 0, + MAX_CLUSTER_RUN_ATTEMPT_LOG_RETENTION_RETRY_DELAY_MS, + ); + return Object.freeze({ + enabled: true as const, + retentionMs: integerValue( + environment, + 'QL3_CLUSTER_LOG_RETENTION_MS', + 30 * 24 * 60 * 60_000, + MIN_RUN_ATTEMPT_LOG_RETENTION_MS, + MAX_RUN_ATTEMPT_LOG_RETENTION_MS, + ), + claimLimit: integerValue( + environment, + 'QL3_CLUSTER_LOG_RETENTION_CLAIM_LIMIT', + 4, + 1, + MAX_CLUSTER_RUN_ATTEMPT_LOG_RETENTION_CLAIMS, + ), + leaseMs, + maximumCycleMs: integerValue( + environment, + 'QL3_CLUSTER_LOG_RETENTION_CYCLE_BUDGET_MS', + 10_000, + 100, + leaseMs - 500, + ), + retryBaseMs, + retryMaximumMs: integerValue( + environment, + 'QL3_CLUSTER_LOG_RETENTION_RETRY_MAX_MS', + 60 * 60_000, + retryBaseMs, + MAX_CLUSTER_RUN_ATTEMPT_LOG_RETENTION_RETRY_DELAY_MS, + ), + maximumFailures: integerValue( + environment, + 'QL3_CLUSTER_LOG_RETENTION_MAX_FAILURES', + 8, + 1, + 32, + ), + intervalMs: integerValue( + environment, + 'QL3_CLUSTER_LOG_RETENTION_INTERVAL_MS', + 60_000, + 1_000, + 24 * 60 * 60_000, + ), + stopTimeoutMs: integerValue( + environment, + 'QL3_CLUSTER_LOG_RETENTION_STOP_TIMEOUT_MS', + 10_000, + 100, + 30_000, + ), + }); +} + /** * Parses the profile gate before reading PostgreSQL configuration. A disabled * cluster-control therefore does not touch its runtime credential source. @@ -345,6 +445,7 @@ export function loadClusterControlConfig( security: Object.freeze({ apiCredentialPepper: apiCredentialPepper(environment), }), + logRetention: logRetentionConfig(environment), }; return Object.freeze(config); } diff --git a/packages/ql3-cluster-control/src/production-process/processApplication.ts b/packages/ql3-cluster-control/src/production-process/processApplication.ts index 49fe5197..fd372136 100644 --- a/packages/ql3-cluster-control/src/production-process/processApplication.ts +++ b/packages/ql3-cluster-control/src/production-process/processApplication.ts @@ -38,6 +38,7 @@ export interface ClusterControlProcessEvent { scope: | 'scheduler' | 'cancellation-convergence' + | 'log-retention' | 'database' | 'worker-ingress'; name: string; @@ -78,10 +79,7 @@ export class ClusterControlProcessError extends Error { | 'QL3_CLUSTER_CONTROL_PROCESS_CONFIG_INVALID' | 'QL3_CLUSTER_CONTROL_PROCESS_DISABLED'; - constructor( - code: ClusterControlProcessError['code'], - message: string, - ) { + constructor(code: ClusterControlProcessError['code'], message: string) { super(message); this.name = 'ClusterControlProcessError'; this.code = code; @@ -103,10 +101,7 @@ function processConfiguration(environment: ClusterControlEnvironment): { ); } const replicaId = environment.QL3_CLUSTER_REPLICA_ID; - if ( - typeof replicaId !== 'string' || - !REPLICA_ID_PATTERN.test(replicaId) - ) { + if (typeof replicaId !== 'string' || !REPLICA_ID_PATTERN.test(replicaId)) { throw new ClusterControlProcessError( 'QL3_CLUSTER_CONTROL_PROCESS_CONFIG_INVALID', 'QL3_CLUSTER_REPLICA_ID must be a stable safe identifier', @@ -205,9 +200,11 @@ export async function runProductionClusterControlProcess( let resolveSignal: | ((signal: ClusterControlProcessSignal) => void) | undefined; - const requestedSignal = new Promise((resolve) => { - resolveSignal = resolve; - }); + const requestedSignal = new Promise( + (resolve) => { + resolveSignal = resolve; + }, + ); let acceptedSignal = false; const unsubscribe = options.signals.subscribe((signal) => { if (acceptedSignal) return; @@ -230,6 +227,14 @@ export async function runProductionClusterControlProcess( ); } artifactBinding = await createBinding(workerIngress.artifact); + if ( + config.logRetention.enabled && + typeof artifactBinding?.store?.retire !== 'function' + ) { + throw new TypeError( + 'Cluster Worker Artifact binding has no log retirement capability', + ); + } if ( workerIngress.secret !== undefined && workerSecretProvider === undefined @@ -274,15 +279,40 @@ export async function runProductionClusterControlProcess( event(replicaId, { level: 'error', event: 'runtime_diagnostic', - diagnostic: diagnosticFact( - 'cancellation-convergence', - error, - ), + diagnostic: diagnosticFact('cancellation-convergence', error), }), ), ).catch(() => undefined); }, }, + ...(workerIngress !== undefined && config.logRetention.enabled + ? { + logRetention: { + store: artifactBinding!.store, + ownerId: replicaId, + retentionMs: config.logRetention.retentionMs, + claimLimit: config.logRetention.claimLimit, + leaseMs: config.logRetention.leaseMs, + maximumCycleMs: config.logRetention.maximumCycleMs, + retryBaseMs: config.logRetention.retryBaseMs, + retryMaximumMs: config.logRetention.retryMaximumMs, + maximumFailures: config.logRetention.maximumFailures, + intervalMs: config.logRetention.intervalMs, + stopTimeoutMs: config.logRetention.stopTimeoutMs, + onDiagnostic(error: unknown) { + void Promise.resolve( + options.emit( + event(replicaId, { + level: 'error', + event: 'runtime_diagnostic', + diagnostic: diagnosticFact('log-retention', error), + }), + ), + ).catch(() => undefined); + }, + }, + } + : {}), ...(workerIngress === undefined ? {} : { @@ -298,10 +328,7 @@ export async function runProductionClusterControlProcess( event(replicaId, { level: 'error', event: 'runtime_diagnostic', - diagnostic: diagnosticFact( - 'worker-ingress', - error, - ), + diagnostic: diagnosticFact('worker-ingress', error), }), ), ).catch(() => undefined); @@ -395,10 +422,7 @@ export async function runProductionClusterControlProcess( unsubscribe(); resolveSignal = undefined; let cleanupError: unknown; - if ( - application?.status === 'active' && - !applicationStopStarted - ) { + if (application?.status === 'active' && !applicationStopStarted) { try { applicationStopStarted = true; await application.stop(); diff --git a/packages/ql3-cluster-control/src/run/runAttemptLogRetentionLifecycle.ts b/packages/ql3-cluster-control/src/run/runAttemptLogRetentionLifecycle.ts new file mode 100644 index 00000000..4089e1d5 --- /dev/null +++ b/packages/ql3-cluster-control/src/run/runAttemptLogRetentionLifecycle.ts @@ -0,0 +1,490 @@ +// Run owns bounded multi-replica Cluster log retirement and one lifecycle timer. +import { + MAX_CLUSTER_RUN_ATTEMPT_LOG_RETENTION_CLAIMS, + MAX_CLUSTER_RUN_ATTEMPT_LOG_RETENTION_LEASE_MS, + MAX_CLUSTER_RUN_ATTEMPT_LOG_RETENTION_RETRY_DELAY_MS, + MIN_CLUSTER_RUN_ATTEMPT_LOG_RETENTION_LEASE_MS, + type ClusterRunAttemptLogRetentionClaim, + type ClusterRunAttemptLogRetentionClaimRepository, + type ClusterRunAttemptLogRetentionFailureCode, +} from '@qinglong/runtime-core/cluster-run-attempt-log-retention'; +import { + createRunAttemptLogRetirementRecord, + MAX_RUN_ATTEMPT_LOG_RETENTION_MS, + MIN_RUN_ATTEMPT_LOG_RETENTION_MS, + type RunAttemptLogRetentionCandidate, + type RunAttemptLogRetirementStore, + type RunAttemptLogRetirementStoreResult, +} from '@qinglong/runtime-core/run-attempt-log-retention'; + +const OWNER_PATTERN = /^[A-Za-z0-9][A-Za-z0-9._:-]{0,127}$/; +const MIN_SETTLEMENT_BUDGET_MS = 500; + +export interface ClusterRunAttemptLogRetirementStore + extends RunAttemptLogRetirementStore { + retire( + candidate: Readonly, + signal?: AbortSignal, + ): Promise>; +} + +export interface ClusterRunAttemptLogRetentionCoordinatorOptions { + readonly ownerId: string; + readonly retentionMs: number; + readonly claimLimit: number; + readonly leaseMs: number; + readonly maximumCycleMs: number; + readonly retryBaseMs: number; + readonly retryMaximumMs: number; + readonly maximumFailures: number; +} + +export interface ClusterRunAttemptLogRetentionCycleEntry { + readonly attemptId: string; + readonly outcome: + | 'deleted' + | 'already_absent' + | 'retry' + | 'manual' + | 'fenced'; +} + +export interface ClusterRunAttemptLogRetentionCycleSummary { + readonly status: 'complete' | 'saturated' | 'budget_exhausted'; + readonly claimed: number; + readonly attempted: number; + readonly retired: number; + readonly alreadyAbsent: number; + readonly retried: number; + readonly manual: number; + readonly fenced: number; + readonly hasMore: boolean; + readonly entries: readonly Readonly[]; +} + +function integer( + name: string, + value: number, + minimum: number, + maximum: number, +): number { + if (!Number.isSafeInteger(value) || value < minimum || value > maximum) { + throw new RangeError(`${name} must be between ${minimum} and ${maximum}`); + } + return value; +} + +function prepareOptions( + options: ClusterRunAttemptLogRetentionCoordinatorOptions, +): Readonly { + const allowed = new Set([ + 'claimLimit', + 'leaseMs', + 'maximumCycleMs', + 'maximumFailures', + 'ownerId', + 'retentionMs', + 'retryBaseMs', + 'retryMaximumMs', + ]); + if ( + !options || + typeof options !== 'object' || + Array.isArray(options) || + Object.keys(options).some((key) => !allowed.has(key)) || + !OWNER_PATTERN.test(options.ownerId) + ) { + throw new TypeError( + 'Cluster Run Attempt log retention options are invalid', + ); + } + const leaseMs = integer( + 'Cluster Run Attempt log retention lease', + options.leaseMs, + MIN_CLUSTER_RUN_ATTEMPT_LOG_RETENTION_LEASE_MS, + MAX_CLUSTER_RUN_ATTEMPT_LOG_RETENTION_LEASE_MS, + ); + const maximumCycleMs = integer( + 'Cluster Run Attempt log retention cycle budget', + options.maximumCycleMs, + 100, + leaseMs - MIN_SETTLEMENT_BUDGET_MS, + ); + const retryBaseMs = integer( + 'Cluster Run Attempt log retention retry base', + options.retryBaseMs, + 0, + MAX_CLUSTER_RUN_ATTEMPT_LOG_RETENTION_RETRY_DELAY_MS, + ); + const retryMaximumMs = integer( + 'Cluster Run Attempt log retention retry maximum', + options.retryMaximumMs, + retryBaseMs, + MAX_CLUSTER_RUN_ATTEMPT_LOG_RETENTION_RETRY_DELAY_MS, + ); + return Object.freeze({ + ownerId: options.ownerId, + retentionMs: integer( + 'Cluster Run Attempt log retention duration', + options.retentionMs, + MIN_RUN_ATTEMPT_LOG_RETENTION_MS, + MAX_RUN_ATTEMPT_LOG_RETENTION_MS, + ), + claimLimit: integer( + 'Cluster Run Attempt log retention claim limit', + options.claimLimit, + 1, + MAX_CLUSTER_RUN_ATTEMPT_LOG_RETENTION_CLAIMS, + ), + leaseMs, + maximumCycleMs, + retryBaseMs, + retryMaximumMs, + maximumFailures: integer( + 'Cluster Run Attempt log retention failure limit', + options.maximumFailures, + 1, + 32, + ), + }); +} + +function failureCode(error: unknown): ClusterRunAttemptLogRetentionFailureCode { + return (error as { readonly reason?: unknown })?.reason === + 'integrity_mismatch' + ? 'artifact_integrity_mismatch' + : 'artifact_unavailable'; +} + +function retryDelay( + failureCount: number, + options: Readonly, +): number { + const multiplier = 2 ** Math.min(failureCount, 30); + return Math.min(options.retryMaximumMs, options.retryBaseMs * multiplier); +} + +function isAbortSignal(value: unknown): value is AbortSignal { + return ( + !!value && + typeof value === 'object' && + typeof (value as AbortSignal).aborted === 'boolean' && + typeof (value as AbortSignal).addEventListener === 'function' && + typeof (value as AbortSignal).removeEventListener === 'function' + ); +} + +export class ClusterRunAttemptLogRetentionCoordinator { + private readonly options: Readonly; + + constructor( + private readonly repository: ClusterRunAttemptLogRetentionClaimRepository, + private readonly store: ClusterRunAttemptLogRetirementStore, + options: ClusterRunAttemptLogRetentionCoordinatorOptions, + ) { + if ( + !repository || + typeof repository.claim !== 'function' || + typeof repository.settle !== 'function' || + !store || + typeof store.retire !== 'function' + ) { + throw new TypeError( + 'Cluster Run Attempt log retention dependencies are invalid', + ); + } + this.options = prepareOptions(options); + } + + async runOnce( + externalSignal?: AbortSignal, + ): Promise> { + if (externalSignal !== undefined && !isAbortSignal(externalSignal)) { + throw new TypeError( + 'Cluster Run Attempt log retention signal is invalid', + ); + } + if (externalSignal?.aborted) throw externalSignal.reason; + const controller = new AbortController(); + const forwardAbort = () => controller.abort(externalSignal?.reason); + externalSignal?.addEventListener('abort', forwardAbort, { once: true }); + const timeout = setTimeout( + () => + controller.abort( + new Error('Cluster log retention cycle budget expired'), + ), + this.options.maximumCycleMs, + ); + timeout.unref?.(); + try { + const page = await this.repository.claim({ + ownerId: this.options.ownerId, + retentionMs: this.options.retentionMs, + limit: this.options.claimLimit, + leaseMs: this.options.leaseMs, + }); + const entries: ClusterRunAttemptLogRetentionCycleEntry[] = []; + let attempted = 0; + let retired = 0; + let alreadyAbsent = 0; + let retried = 0; + let manual = 0; + let fenced = 0; + for (const claim of page.claims) { + if (controller.signal.aborted) break; + attempted += 1; + let result: Readonly; + try { + result = await this.store.retire(claim.candidate, controller.signal); + } catch (error) { + if (controller.signal.aborted) break; + const settled = await this.settleFailure(claim, failureCode(error)); + if (settled === 'fenced') { + fenced += 1; + entries.push({ + attemptId: claim.candidate.attemptId, + outcome: 'fenced', + }); + } else if (claim.failureCount + 1 >= this.options.maximumFailures) { + manual += 1; + entries.push({ + attemptId: claim.candidate.attemptId, + outcome: 'manual', + }); + } else { + retried += 1; + entries.push({ + attemptId: claim.candidate.attemptId, + outcome: 'retry', + }); + } + continue; + } + if (controller.signal.aborted) break; + let record; + try { + record = createRunAttemptLogRetirementRecord({ + ...claim.candidate, + eligibleAtMs: claim.eligibleAtMs, + retiredAtMs: claim.observedAtMs, + ...result, + }); + } catch { + const settled = await this.settleFailure( + claim, + 'retirement_record_unavailable', + ); + if (settled === 'fenced') { + fenced += 1; + entries.push({ + attemptId: claim.candidate.attemptId, + outcome: 'fenced', + }); + } else if (claim.failureCount + 1 >= this.options.maximumFailures) { + manual += 1; + entries.push({ + attemptId: claim.candidate.attemptId, + outcome: 'manual', + }); + } else { + retried += 1; + entries.push({ + attemptId: claim.candidate.attemptId, + outcome: 'retry', + }); + } + continue; + } + const settled = await this.repository.settle(claim, { + status: 'retired', + record, + }); + if (settled === 'fenced') { + fenced += 1; + entries.push({ + attemptId: claim.candidate.attemptId, + outcome: 'fenced', + }); + continue; + } + if (record.disposition === 'already_absent') alreadyAbsent += 1; + else retired += 1; + entries.push({ + attemptId: claim.candidate.attemptId, + outcome: record.disposition, + }); + } + const budgetExhausted = controller.signal.aborted; + return Object.freeze({ + status: budgetExhausted + ? ('budget_exhausted' as const) + : page.hasMore + ? ('saturated' as const) + : ('complete' as const), + claimed: page.claims.length, + attempted, + retired, + alreadyAbsent, + retried, + manual, + fenced, + hasMore: page.hasMore, + entries: Object.freeze(entries.map((entry) => Object.freeze(entry))), + }); + } finally { + clearTimeout(timeout); + externalSignal?.removeEventListener('abort', forwardAbort); + } + } + + private settleFailure( + claim: Readonly, + code: ClusterRunAttemptLogRetentionFailureCode, + ): Promise<'settled' | 'fenced'> { + if (claim.failureCount + 1 >= this.options.maximumFailures) { + return this.repository.settle(claim, { + status: 'manual', + failureCode: code, + }); + } + return this.repository.settle(claim, { + status: 'retry', + delayMs: retryDelay(claim.failureCount, this.options), + failureCode: code, + }); + } +} + +export interface ClusterRunAttemptLogRetentionLifecycleOptions { + readonly intervalMs: number; + readonly stopTimeoutMs: number; + readonly onDiagnostic?: ( + error: unknown, + summary?: Readonly, + ) => void | Promise; +} + +export class ClusterRunAttemptLogRetentionLifecycle { + private timer: NodeJS.Timeout | undefined; + private inFlight: + | Promise> + | undefined; + private controller: AbortController | undefined; + private stopPromise: Promise<'stopped' | 'timed_out'> | undefined; + private running = false; + private stopping = false; + + constructor( + private readonly coordinator: Pick< + ClusterRunAttemptLogRetentionCoordinator, + 'runOnce' + >, + private readonly options: ClusterRunAttemptLogRetentionLifecycleOptions, + ) { + if ( + !coordinator || + typeof coordinator.runOnce !== 'function' || + !options || + typeof options !== 'object' || + Array.isArray(options) || + !Number.isSafeInteger(options.intervalMs) || + options.intervalMs < 1_000 || + options.intervalMs > 24 * 60 * 60_000 || + !Number.isSafeInteger(options.stopTimeoutMs) || + options.stopTimeoutMs < 100 || + options.stopTimeoutMs > 30_000 || + (options.onDiagnostic !== undefined && + typeof options.onDiagnostic !== 'function') + ) { + throw new TypeError( + 'Cluster Run Attempt log retention lifecycle options are invalid', + ); + } + } + + start(): 'started' { + if (!this.running && !this.stopping) { + this.running = true; + this.schedule(); + } + return 'started'; + } + + runOnce(): Promise> { + if (this.stopping) { + return Promise.reject( + new Error('Cluster Run Attempt log retention lifecycle is stopping'), + ); + } + if (this.inFlight) return this.inFlight; + const controller = new AbortController(); + this.controller = controller; + const work = this.coordinator.runOnce(controller.signal).finally(() => { + if (this.inFlight === work) { + this.inFlight = undefined; + this.controller = undefined; + } + }); + this.inFlight = work; + return work; + } + + stopAndDrain(): Promise<'stopped' | 'timed_out'> { + if (this.stopPromise) return this.stopPromise; + this.stopping = true; + this.running = false; + if (this.timer) clearTimeout(this.timer); + this.timer = undefined; + this.controller?.abort( + new Error('Cluster Run Attempt log retention lifecycle is stopping'), + ); + this.stopPromise = (async () => { + const work = this.inFlight; + if (!work) return 'stopped' as const; + let timeout: NodeJS.Timeout | undefined; + try { + return await Promise.race([ + work.then( + () => 'stopped' as const, + () => 'stopped' as const, + ), + new Promise<'timed_out'>((resolve) => { + timeout = setTimeout( + () => resolve('timed_out' as const), + this.options.stopTimeoutMs, + ); + timeout.unref?.(); + }), + ]); + } finally { + if (timeout) clearTimeout(timeout); + } + })(); + return this.stopPromise; + } + + private schedule(): void { + if (!this.running || this.timer) return; + this.timer = setTimeout(() => { + this.timer = undefined; + if (!this.running) return; + void this.runOnce() + .then((summary) => this.diagnostic(undefined, summary)) + .catch((error) => this.diagnostic(error)) + .finally(() => this.schedule()); + }, this.options.intervalMs); + this.timer.unref?.(); + } + + private async diagnostic( + error: unknown, + summary?: Readonly, + ): Promise { + if (this.stopping) return; + try { + await this.options.onDiagnostic?.(error, summary); + } catch { + // Diagnostics cannot own or stop retention. + } + } +} diff --git a/packages/ql3-cluster-control/test/bootstrap.test.cjs b/packages/ql3-cluster-control/test/bootstrap.test.cjs index a34fbf4e..e3bbbb72 100644 --- a/packages/ql3-cluster-control/test/bootstrap.test.cjs +++ b/packages/ql3-cluster-control/test/bootstrap.test.cjs @@ -388,6 +388,7 @@ function bootstrapOptions(events, overrides = {}) { authenticator, policies, runs, + runAttemptLogRetention, runCancellation, taskDefinitions, taskExecutionRevisions, @@ -408,6 +409,7 @@ function bootstrapOptions(events, overrides = {}) { assert.equal(typeof authenticator.authenticate, 'function'); assert.equal(typeof policies.resolve, 'function'); assert.equal(typeof runs.transaction, 'function'); + assert.equal(typeof runAttemptLogRetention.inspect, 'function'); assert.equal(typeof runCancellation.requestUserCancellation, 'function'); assert.equal( typeof taskDefinitions.findCurrentTaskDefinition, diff --git a/packages/ql3-cluster-control/test/config.test.cjs b/packages/ql3-cluster-control/test/config.test.cjs index 44b27ada..22ac717a 100644 --- a/packages/ql3-cluster-control/test/config.test.cjs +++ b/packages/ql3-cluster-control/test/config.test.cjs @@ -114,6 +114,18 @@ test('builds an exact runtime-only TLS-verified Pool configuration', async () => assert.deepEqual(config.security, { apiCredentialPepper: BASE_ENV.QL3_API_CREDENTIAL_PEPPER, }); + assert.deepEqual(config.logRetention, { + enabled: true, + retentionMs: 30 * 24 * 60 * 60_000, + claimLimit: 4, + leaseMs: 30_000, + maximumCycleMs: 10_000, + retryBaseMs: 5_000, + retryMaximumMs: 60 * 60_000, + maximumFailures: 8, + intervalMs: 60_000, + stopTimeoutMs: 10_000, + }); const binding = createClusterControlDatabaseBinding(config); assert.equal(binding.availability.status, 'available'); @@ -192,6 +204,17 @@ test('rejects TLS query overrides, missing credentials and unbounded values', () { ...BASE_ENV, QL3_CLUSTER_AUTH_RATE_GLOBAL: '1000001' }, { ...BASE_ENV, QL3_CLUSTER_AUTH_RATE_MAX_PEERS: '65537' }, { ...BASE_ENV, QL3_API_CREDENTIAL_PEPPER: 'weak' }, + { ...BASE_ENV, QL3_CLUSTER_LOG_RETENTION_CLAIM_LIMIT: '17' }, + { + ...BASE_ENV, + QL3_CLUSTER_LOG_RETENTION_LEASE_MS: '5000', + QL3_CLUSTER_LOG_RETENTION_CYCLE_BUDGET_MS: '4501', + }, + { + ...BASE_ENV, + QL3_CLUSTER_LOG_RETENTION_RETRY_BASE_MS: '5000', + QL3_CLUSTER_LOG_RETENTION_RETRY_MAX_MS: '4999', + }, ]) { assert.throws( () => loadClusterControlConfig(environment), @@ -199,3 +222,37 @@ test('rejects TLS query overrides, missing credentials and unbounded values', () ); } }); + +test('loads bounded Cluster log retention policy and permits explicit disable', () => { + const disabled = loadClusterControlConfig({ + ...BASE_ENV, + QL3_CLUSTER_LOG_RETENTION_ENABLED: 'false', + QL3_CLUSTER_LOG_RETENTION_CLAIM_LIMIT: '999', + }); + assert.deepEqual(disabled.logRetention, { enabled: false }); + + const configured = loadClusterControlConfig({ + ...BASE_ENV, + QL3_CLUSTER_LOG_RETENTION_MS: '60000', + QL3_CLUSTER_LOG_RETENTION_CLAIM_LIMIT: '2', + QL3_CLUSTER_LOG_RETENTION_LEASE_MS: '5000', + QL3_CLUSTER_LOG_RETENTION_CYCLE_BUDGET_MS: '4000', + QL3_CLUSTER_LOG_RETENTION_RETRY_BASE_MS: '250', + QL3_CLUSTER_LOG_RETENTION_RETRY_MAX_MS: '1000', + QL3_CLUSTER_LOG_RETENTION_MAX_FAILURES: '3', + QL3_CLUSTER_LOG_RETENTION_INTERVAL_MS: '2000', + QL3_CLUSTER_LOG_RETENTION_STOP_TIMEOUT_MS: '500', + }); + assert.deepEqual(configured.logRetention, { + enabled: true, + retentionMs: 60_000, + claimLimit: 2, + leaseMs: 5_000, + maximumCycleMs: 4_000, + retryBaseMs: 250, + retryMaximumMs: 1_000, + maximumFailures: 3, + intervalMs: 2_000, + stopTimeoutMs: 500, + }); +}); diff --git a/packages/ql3-cluster-control/test/processApplication.test.cjs b/packages/ql3-cluster-control/test/processApplication.test.cjs index 5683305d..ac070ad9 100644 --- a/packages/ql3-cluster-control/test/processApplication.test.cjs +++ b/packages/ql3-cluster-control/test/processApplication.test.cjs @@ -75,8 +75,14 @@ test('runs one production replica and drains it on the first signal', async () = assert.equal(result, 'stopped'); assert.deepEqual(events, ['subscribe', 'start', 'stop', 'unsubscribe']); - assert.equal(facts.some((fact) => fact.event === 'activation'), true); - assert.equal(facts.some((fact) => fact.event === 'listening'), true); + assert.equal( + facts.some((fact) => fact.event === 'activation'), + true, + ); + assert.equal( + facts.some((fact) => fact.event === 'listening'), + true, + ); assert.equal( facts.some( (fact) => @@ -105,7 +111,11 @@ test('fails closed before startup for a disabled profile or invalid replica id', await assert.rejects( runProductionClusterControlProcess({ environment, - signals: { subscribe() { return () => {}; } }, + signals: { + subscribe() { + return () => {}; + }, + }, emit() {}, async start() { starts += 1; @@ -154,6 +164,7 @@ test('starts the optional Worker listener and closes its lazy Artifact binding', const artifactStore = { async put() {}, async inspect() {}, + async retire() {}, }; const environment = { ...BASE_ENV, @@ -187,6 +198,14 @@ test('starts the optional Worker listener and closes its lazy Artifact binding', events.push('start'); assert.equal(options.workerIngress.config.enabled, true); assert.equal(options.workerIngress.artifactStore, artifactStore); + assert.equal(options.logRetention.store, artifactStore); + assert.equal(options.logRetention.ownerId, 'cluster-control-0'); + assert.equal(options.logRetention.claimLimit, 4); + options.logRetention.onDiagnostic( + Object.assign(new Error('must-not-be-logged'), { + code: 'S3Unavailable', + }), + ); return { status: 'active', address: { host: '0.0.0.0', port: 5800 }, @@ -219,6 +238,16 @@ test('starts the optional Worker listener and closes its lazy Artifact binding', facts.some((fact) => fact.event === 'worker_ingress_listening'), true, ); + assert.equal( + facts.some( + (fact) => + fact.event === 'runtime_diagnostic' && + fact.diagnostic.scope === 'log-retention' && + fact.diagnostic.code === 'S3Unavailable' && + JSON.stringify(fact).includes('must-not-be-logged') === false, + ), + true, + ); }); test('creates the configured mounted Secret provider before Worker activation', async () => { @@ -226,6 +255,7 @@ test('creates the configured mounted Secret provider before Worker activation', const artifactStore = { async put() {}, async inspect() {}, + async retire() {}, }; const provider = { async resolve() {} }; const result = await runProductionClusterControlProcess({ @@ -353,7 +383,10 @@ test('fails the process after a database fence drains the active application', a ), true, ); - assert.equal(facts.some(({ event }) => event === 'shutdown_requested'), false); + assert.equal( + facts.some(({ event }) => event === 'shutdown_requested'), + false, + ); assert.equal(facts.at(-1).event, 'stopped'); assert.equal( JSON.stringify(facts).includes('must-not-escape-database-detail'), diff --git a/packages/ql3-cluster-control/test/productionApplication.test.cjs b/packages/ql3-cluster-control/test/productionApplication.test.cjs index 46976ccf..f29dddb2 100644 --- a/packages/ql3-cluster-control/test/productionApplication.test.cjs +++ b/packages/ql3-cluster-control/test/productionApplication.test.cjs @@ -617,6 +617,87 @@ test('wires the production Worker object reader into the Project-scoped log rout assert.equal(Buffer.from(result.body.content, 'base64').toString(), 'prod'); }); +test('wires durable retirement authority into the production log route', async () => { + const { + createRunAttemptLogRetirementRecord, + } = require('@qinglong/runtime-core/run-attempt-log-retention'); + const { input } = fixture(); + const run = await input.runs.findRunById('run-1'); + const logArtifactId = `wlog-${'b'.repeat(30)}`; + let objectReads = 0; + const stack = createProductionClusterControlApplicationStack({ + ...input, + runs: { + ...input.runs, + async findRunById() { + return { ...run, status: 'succeeded', finishedAtMs: 10 }; + }, + async findAttemptById() { + return { + id: 'attempt-1', + runId: 'run-1', + attempt: 1, + status: 'succeeded', + executorType: 'remote_worker', + logArtifactId, + callbackSequence: 0, + createdAtMs: 1, + finishedAtMs: 10, + }; + }, + }, + runAttemptLogRetention: { + async inspect(identity) { + assert.deepEqual(identity, { + projectId: 'project-1', + runId: 'run-1', + attemptId: 'attempt-1', + logArtifactId, + }); + return { + status: 'retired', + record: createRunAttemptLogRetirementRecord({ + ...identity, + executorType: 'remote_worker', + finishedAtMs: 10, + eligibleAtMs: 20, + retiredAtMs: 30, + disposition: 'deleted', + byteLength: 64, + truncation: { truncated: 'unknown' }, + }), + }; + }, + }, + workerRuntime: { + offers: { claimNext() {} }, + activation: { + acknowledgeStarting() {}, + acknowledgeRunning() {}, + failStart() {}, + }, + artifacts: { upload() {} }, + completion: { complete() {} }, + leaseControl: { control() {} }, + runAttemptLogRead: { + async read() { + objectReads += 1; + return { status: 'missing' }; + }, + }, + }, + }); + const result = await invoke( + stack, + metadata('/api/v3/projects/project-1/runs/run-1/attempts/attempt-1/log'), + ); + assert.equal(result.statusCode, 410); + assert.equal(result.body.status, 'retired'); + assert.equal(result.body.retiredAtMs, 30); + assert.equal(result.body.byteLength, 64); + assert.equal(objectReads, 0); +}); + test('optionally exposes Prompt execution behind shared admission and policy', async () => { const { events, input } = fixture(); let command; diff --git a/packages/ql3-cluster-control/test/runAttemptLogRetentionLifecycle.test.cjs b/packages/ql3-cluster-control/test/runAttemptLogRetentionLifecycle.test.cjs new file mode 100644 index 00000000..55addf0e --- /dev/null +++ b/packages/ql3-cluster-control/test/runAttemptLogRetentionLifecycle.test.cjs @@ -0,0 +1,281 @@ +'use strict'; + +const assert = require('node:assert/strict'); +const { test } = require('node:test'); +const { + ClusterRunAttemptLogRetentionCoordinator, + ClusterRunAttemptLogRetentionLifecycle, +} = require('../dist/run/runAttemptLogRetentionLifecycle'); + +function claim(overrides = {}) { + return Object.freeze({ + candidate: Object.freeze({ + projectId: 'project-1', + runId: 'run-1', + attemptId: 'attempt-1', + logArtifactId: `wlog-${'a'.repeat(30)}`, + executorType: 'remote_worker', + finishedAtMs: 1_000, + }), + eligibleAtMs: 61_000, + observedAtMs: 70_000, + ownerId: 'replica-a', + token: '00000000-0000-4000-8000-000000000055', + version: 1, + expiresAtMs: 100_000, + failureCount: 0, + ...overrides, + }); +} + +function options(overrides = {}) { + return { + ownerId: 'replica-a', + retentionMs: 60_000, + claimLimit: 4, + leaseMs: 30_000, + maximumCycleMs: 10_000, + retryBaseMs: 1_000, + retryMaximumMs: 8_000, + maximumFailures: 3, + ...overrides, + }; +} + +function coordinator({ + claims = [claim()], + hasMore = false, + retire, + settle, +} = {}) { + const calls = []; + const value = new ClusterRunAttemptLogRetentionCoordinator( + { + async claim(input) { + calls.push(['claim', input]); + return { claims, hasMore }; + }, + async settle(current, settlement) { + calls.push(['settle', current, settlement]); + return (await settle?.(current, settlement)) ?? 'settled'; + }, + async inspect() { + return { status: 'active' }; + }, + }, + { + async retire(candidate, signal) { + calls.push(['retire', candidate, signal]); + return ( + (await retire?.(candidate, signal)) ?? { + disposition: 'deleted', + byteLength: 11, + truncation: { truncated: false }, + } + ); + }, + }, + options(), + ); + return { calls, coordinator: value }; +} + +test('claims one bounded page and records exact DB-clock retirement evidence', async () => { + const { calls, coordinator: value } = coordinator({ + claims: [ + claim(), + claim({ + candidate: Object.freeze({ + ...claim().candidate, + attemptId: 'attempt-2', + logArtifactId: `wlog-${'b'.repeat(30)}`, + }), + }), + ], + hasMore: true, + retire(candidate) { + return candidate.attemptId === 'attempt-1' + ? { + disposition: 'deleted', + byteLength: 11, + truncation: { truncated: false }, + } + : { + disposition: 'already_absent', + byteLength: 0, + truncation: { truncated: 'unknown' }, + }; + }, + }); + + const summary = await value.runOnce(); + assert.deepEqual(summary, { + status: 'saturated', + claimed: 2, + attempted: 2, + retired: 1, + alreadyAbsent: 1, + retried: 0, + manual: 0, + fenced: 0, + hasMore: true, + entries: [ + { attemptId: 'attempt-1', outcome: 'deleted' }, + { attemptId: 'attempt-2', outcome: 'already_absent' }, + ], + }); + assert.deepEqual(calls[0][1], { + ownerId: 'replica-a', + retentionMs: 60_000, + limit: 4, + leaseMs: 30_000, + }); + const records = calls + .filter(([kind]) => kind === 'settle') + .map(([, , settlement]) => settlement.record); + assert.equal(records[0].retiredAtMs, 70_000); + assert.equal(records[0].recordDigest.length, 64); + assert.equal(records[1].byteLength, 0); +}); + +test('uses bounded exponential retry then moves repeated failures to manual', async () => { + const first = claim({ failureCount: 2 }); + const second = claim({ + candidate: Object.freeze({ + ...claim().candidate, + attemptId: 'attempt-2', + logArtifactId: `wlog-${'b'.repeat(30)}`, + }), + failureCount: 1, + }); + const { calls, coordinator: value } = coordinator({ + claims: [first, second], + retire(candidate) { + const error = new Error('object drift'); + if (candidate.attemptId === 'attempt-1') + error.reason = 'integrity_mismatch'; + throw error; + }, + }); + + const summary = await value.runOnce(); + assert.equal(summary.manual, 1); + assert.equal(summary.retried, 1); + const settlements = calls + .filter(([kind]) => kind === 'settle') + .map(([, , settlement]) => settlement); + assert.deepEqual(settlements, [ + { status: 'manual', failureCode: 'artifact_integrity_mismatch' }, + { + status: 'retry', + delayMs: 2_000, + failureCode: 'artifact_unavailable', + }, + ]); +}); + +test('classifies malformed retirement evidence and preserves a fenced settlement', async () => { + const { coordinator: value } = coordinator({ + retire() { + return { + disposition: 'already_absent', + byteLength: 5, + truncation: { truncated: 'unknown' }, + }; + }, + settle() { + return 'fenced'; + }, + }); + const summary = await value.runOnce(); + assert.equal(summary.fenced, 1); + assert.deepEqual(summary.entries, [ + { attemptId: 'attempt-1', outcome: 'fenced' }, + ]); +}); + +test('cycle budget aborts object work and leaves the durable claim for takeover', async () => { + let settlements = 0; + const value = new ClusterRunAttemptLogRetentionCoordinator( + { + async claim() { + return { claims: [claim()], hasMore: false }; + }, + async settle() { + settlements += 1; + return 'settled'; + }, + async inspect() { + return { status: 'active' }; + }, + }, + { + retire(_candidate, signal) { + return new Promise((_, reject) => { + signal.addEventListener('abort', () => reject(signal.reason), { + once: true, + }); + }); + }, + }, + options({ maximumCycleMs: 100 }), + ); + const summary = await value.runOnce(); + assert.equal(summary.status, 'budget_exhausted'); + assert.equal(summary.attempted, 1); + assert.equal(settlements, 0); +}); + +test('lifecycle coalesces cycles and aborts one in-flight object call on drain', async () => { + let calls = 0; + let observedAbort = false; + const lifecycle = new ClusterRunAttemptLogRetentionLifecycle( + { + runOnce(signal) { + calls += 1; + return new Promise((resolve) => { + signal.addEventListener( + 'abort', + () => { + observedAbort = true; + resolve({ status: 'budget_exhausted' }); + }, + { once: true }, + ); + }); + }, + }, + { intervalMs: 60_000, stopTimeoutMs: 1_000 }, + ); + const first = lifecycle.runOnce(); + assert.equal(lifecycle.runOnce(), first); + assert.equal(await lifecycle.stopAndDrain(), 'stopped'); + await first; + assert.equal(calls, 1); + assert.equal(observedAbort, true); + await assert.rejects(lifecycle.runOnce(), /is stopping/); +}); + +test('rejects configurations that can outlive the lease settlement budget', () => { + const dependencies = [ + { claim() {}, settle() {}, inspect() {} }, + { retire() {} }, + ]; + assert.throws( + () => + new ClusterRunAttemptLogRetentionCoordinator( + dependencies[0], + dependencies[1], + options({ leaseMs: 5_000, maximumCycleMs: 4_501 }), + ), + /cycle budget/, + ); + assert.throws( + () => + new ClusterRunAttemptLogRetentionLifecycle( + { runOnce() {} }, + { intervalMs: 999, stopTimeoutMs: 1_000 }, + ), + /lifecycle options/, + ); +}); diff --git a/packages/ql3-cluster-control/test/s3ArtifactStore.integration.test.cjs b/packages/ql3-cluster-control/test/s3ArtifactStore.integration.test.cjs index 480e523e..b22b543c 100644 --- a/packages/ql3-cluster-control/test/s3ArtifactStore.integration.test.cjs +++ b/packages/ql3-cluster-control/test/s3ArtifactStore.integration.test.cjs @@ -7,7 +7,9 @@ const { CreateBucketCommand, DeleteBucketCommand, DeleteObjectsCommand, + ListObjectVersionsCommand, ListObjectsV2Command, + PutBucketVersioningCommand, S3Client, } = require('@aws-sdk/client-s3'); const { @@ -35,6 +37,7 @@ test( credentials: { accessKeyId, secretAccessKey }, }); const bucket = `ql3-artifact-${process.pid}-${Date.now()}`.slice(0, 63); + const versionedBucket = `${bucket}-v`.slice(0, 63); const command = Object.freeze({ projectId: 'project-s3-integration', runId: 'run-s3-integration', @@ -101,27 +104,123 @@ test( ); assert.equal(objects.KeyCount, 1); assert.match(objects.Contents[0].Key, /\/objects\//); - } finally { - try { - const objects = await client.send( - new ListObjectsV2Command({ - Bucket: bucket, - }), - ); - if (objects.Contents?.length) { + const retired = await store.retire({ + projectId: command.projectId, + runId: command.runId, + attemptId: command.attemptId, + logArtifactId: command.logArtifactId, + executorType: 'remote_worker', + finishedAtMs: 1, + }); + assert.deepEqual(retired, { + disposition: 'deleted', + byteLength: content.byteLength, + truncation: { truncated: true }, + }); + assert.equal( + ( await client.send( - new DeleteObjectsCommand({ + new ListObjectsV2Command({ Bucket: bucket, - Delete: { - Objects: objects.Contents.map(({ Key }) => ({ Key })), - Quiet: true, - }, + Prefix: 'qinglong/integration/', }), + ) + ).KeyCount, + 0, + ); + assert.deepEqual( + await store.retire({ + projectId: command.projectId, + runId: command.runId, + attemptId: command.attemptId, + logArtifactId: command.logArtifactId, + executorType: 'remote_worker', + finishedAtMs: 1, + }), + { + disposition: 'already_absent', + byteLength: 0, + truncation: { truncated: 'unknown' }, + }, + ); + + await client.send(new CreateBucketCommand({ Bucket: versionedBucket })); + await client.send( + new PutBucketVersioningCommand({ + Bucket: versionedBucket, + VersioningConfiguration: { Status: 'Enabled' }, + }), + ); + const versionedStore = new S3ClusterRemoteWorkerArtifactStore({ + client, + bucket: versionedBucket, + prefix: 'qinglong/integration', + encryption: { mode: 's3' }, + }); + assert.equal( + (await versionedStore.put(command, body(content))).status, + 'stored', + ); + const beforeVersionedRetirement = await client.send( + new ListObjectVersionsCommand({ Bucket: versionedBucket }), + ); + assert.equal(beforeVersionedRetirement.Versions?.length, 1); + assert.equal(beforeVersionedRetirement.DeleteMarkers?.length ?? 0, 0); + assert.match(beforeVersionedRetirement.Versions[0].Key, /\/objects\//); + assert.equal( + ( + await versionedStore.retire({ + projectId: command.projectId, + runId: command.runId, + attemptId: command.attemptId, + logArtifactId: command.logArtifactId, + executorType: 'remote_worker', + finishedAtMs: 1, + }) + ).disposition, + 'deleted', + ); + const afterVersionedRetirement = await client.send( + new ListObjectVersionsCommand({ Bucket: versionedBucket }), + ); + assert.equal(afterVersionedRetirement.Versions?.length ?? 0, 0); + assert.equal(afterVersionedRetirement.DeleteMarkers?.length ?? 0, 0); + } finally { + for (const cleanupBucket of [versionedBucket, bucket]) { + try { + const versions = await client.send( + new ListObjectVersionsCommand({ Bucket: cleanupBucket }), ); + const versionedObjects = [ + ...(versions.Versions ?? []), + ...(versions.DeleteMarkers ?? []), + ].map(({ Key, VersionId }) => ({ Key, VersionId })); + if (versionedObjects.length) { + await client.send( + new DeleteObjectsCommand({ + Bucket: cleanupBucket, + Delete: { Objects: versionedObjects, Quiet: true }, + }), + ); + } + const objects = await client.send( + new ListObjectsV2Command({ Bucket: cleanupBucket }), + ); + if (objects.Contents?.length) { + await client.send( + new DeleteObjectsCommand({ + Bucket: cleanupBucket, + Delete: { + Objects: objects.Contents.map(({ Key }) => ({ Key })), + Quiet: true, + }, + }), + ); + } + await client.send(new DeleteBucketCommand({ Bucket: cleanupBucket })); + } catch { + // Preserve the integration assertion; the ephemeral container is removed. } - await client.send(new DeleteBucketCommand({ Bucket: bucket })); - } catch { - // Preserve the integration assertion; the ephemeral container is removed. } client.destroy(); } diff --git a/packages/ql3-cluster-control/test/s3ArtifactStore.test.cjs b/packages/ql3-cluster-control/test/s3ArtifactStore.test.cjs index 1ee8a3b6..06dc5150 100644 --- a/packages/ql3-cluster-control/test/s3ArtifactStore.test.cjs +++ b/packages/ql3-cluster-control/test/s3ArtifactStore.test.cjs @@ -73,6 +73,14 @@ class MemoryS3Client { ? Buffer.alloc(32, 9).toString('base64') : checksum(object.content), Metadata: metadata, + ...(input.Key.includes('/objects/') && + this.options.headVersionId !== undefined + ? { VersionId: this.options.headVersionId } + : {}), + ...(input.Key.includes('/temporary/') && + this.options.temporaryHeadVersionId !== undefined + ? { VersionId: this.options.temporaryHeadVersionId } + : {}), }; } if (command instanceof GetObjectCommand) { @@ -168,7 +176,33 @@ class MemoryS3Client { } if (command instanceof DeleteObjectCommand) { if (this.options.failDelete) throw new Error('delete unavailable'); + if (input.Key.includes('/objects/')) { + if (this.options.permanentDeletePreconditionFailure) { + const error = new Error('precondition failed'); + error.name = 'PreconditionFailed'; + error.$metadata = { httpStatusCode: 412 }; + throw error; + } + const object = this.objects.get(input.Key); + if (!object) throw notFound(); + if (this.options.headVersionId === undefined) { + assert.equal( + input.IfMatch, + `"${checksum(object.content).slice(0, 32)}"`, + ); + assert.equal(input.VersionId, undefined); + } else { + assert.equal(input.IfMatch, undefined); + assert.equal(input.VersionId, this.options.headVersionId); + } + } this.objects.delete(input.Key); + if ( + input.Key.includes('/objects/') && + this.options.throwAfterPermanentDelete + ) { + throw new Error('lost delete response'); + } return {}; } throw new Error(`unexpected command: ${command.constructor.name}`); @@ -201,6 +235,15 @@ function permanentKey(client) { return [...client.objects.keys()].find((key) => key.includes('/objects/')); } +function retentionCandidate(overrides = {}) { + return Object.freeze({ + ...LOOKUP, + executorType: 'remote_worker', + finishedAtMs: 1_000, + ...overrides, + }); +} + test('streams to a checksummed temporary object then conditionally promotes it', async () => { const client = new MemoryS3Client(); const adapter = store(client); @@ -233,6 +276,13 @@ test('streams to a checksummed temporary object then conditionally promotes it', const copy = client.commands.find( (command) => command instanceof CopyObjectCommand, ); + const cleanup = client.commands.find( + (command) => + command instanceof DeleteObjectCommand && + command.input.Key.includes('/temporary/'), + ); + assert.match(cleanup.input.IfMatch, /^"[A-Za-z0-9+/=]+"$/); + assert.equal(cleanup.input.VersionId, undefined); assert.equal(copy.input.Metadata['ql3-content-sha256'], CONTENT_SHA256); assert.equal( JSON.stringify(copy.input.Metadata).includes(COMMAND.projectId), @@ -247,6 +297,20 @@ test('streams to a checksummed temporary object then conditionally promotes it', assert.deepEqual(inspected, { ...receipt, status: 'already_stored' }); }); +test('cleans one exact temporary object version after validated HEAD', async () => { + const client = new MemoryS3Client({ + temporaryHeadVersionId: 'temporary/version+1=', + }); + await store(client).put(COMMAND, chunks()); + const cleanup = client.commands.find( + (command) => + command instanceof DeleteObjectCommand && + command.input.Key.includes('/temporary/'), + ); + assert.equal(cleanup.input.VersionId, 'temporary/version+1='); + assert.equal(cleanup.input.IfMatch, undefined); +}); + test('exact replay consumes and hashes the whole body without another write', async () => { const client = new MemoryS3Client(); const adapter = store(client); @@ -510,3 +574,101 @@ test('a pre-aborted request performs no object-store operation', async () => { ); assert.equal(client.commands.length, 0); }); + +test('retires an unversioned Artifact only with its validated ETag', async () => { + const client = new MemoryS3Client(); + const adapter = store(client, { expectedBucketOwner: '123456789012' }); + await adapter.put(COMMAND, chunks()); + client.commands.length = 0; + + assert.deepEqual(await adapter.retire(retentionCandidate()), { + disposition: 'deleted', + byteLength: CONTENT.byteLength, + truncation: { truncated: false }, + }); + assert.deepEqual( + client.commands.map((command) => command.constructor.name), + ['HeadObjectCommand', 'DeleteObjectCommand'], + ); + assert.equal(client.commands[1].input.ExpectedBucketOwner, '123456789012'); + assert.equal(permanentKey(client), undefined); +}); + +test('retires one exact version when HEAD returns an opaque VersionId', async () => { + const client = new MemoryS3Client({ headVersionId: 'version/opaque+1=' }); + const adapter = store(client); + await adapter.put(COMMAND, chunks()); + client.commands.length = 0; + + const result = await adapter.retire(retentionCandidate()); + assert.equal(result.disposition, 'deleted'); + assert.equal(client.commands[1].input.VersionId, 'version/opaque+1='); + assert.equal(client.commands[1].input.IfMatch, undefined); + assert.equal(permanentKey(client), undefined); +}); + +test('returns durable absent evidence without issuing a delete', async () => { + const client = new MemoryS3Client(); + const adapter = store(client); + + assert.deepEqual(await adapter.retire(retentionCandidate()), { + disposition: 'already_absent', + byteLength: 0, + truncation: { truncated: 'unknown' }, + }); + assert.deepEqual( + client.commands.map((command) => command.constructor.name), + ['HeadObjectCommand'], + ); +}); + +test('fails closed on conditional-delete drift and malformed version authority', async () => { + for (const options of [ + { permanentDeletePreconditionFailure: true }, + { headVersionId: 'invalid\nversion' }, + ]) { + const client = new MemoryS3Client(options); + const adapter = store(client); + await adapter.put(COMMAND, chunks()); + client.commands.length = 0; + + await assert.rejects( + adapter.retire(retentionCandidate()), + (error) => + error instanceof S3ClusterRemoteWorkerArtifactStoreError && + error.reason === 'integrity_mismatch', + ); + assert.notEqual(permanentKey(client), undefined); + } + + const wrongExecutor = new MemoryS3Client(); + await assert.rejects( + store(wrongExecutor).retire( + retentionCandidate({ executorType: 'local_process' }), + ), + (error) => + error instanceof S3ClusterRemoteWorkerArtifactStoreError && + error.reason === 'integrity_mismatch', + ); + assert.equal(wrongExecutor.commands.length, 0); +}); + +test('lost delete response converges through a later absent inspection', async () => { + const client = new MemoryS3Client({ throwAfterPermanentDelete: true }); + const adapter = store(client); + await adapter.put(COMMAND, chunks()); + client.commands.length = 0; + + await assert.rejects( + adapter.retire(retentionCandidate()), + (error) => + error instanceof S3ClusterRemoteWorkerArtifactStoreError && + error.reason === 'unavailable', + ); + client.options.throwAfterPermanentDelete = false; + assert.deepEqual(await adapter.retire(retentionCandidate()), { + disposition: 'already_absent', + byteLength: 0, + truncation: { truncated: 'unknown' }, + }); +}); diff --git a/scripts/ql3-postgres-ha-contract.cjs b/scripts/ql3-postgres-ha-contract.cjs index bd16f86f..8d0f0c23 100644 --- a/scripts/ql3-postgres-ha-contract.cjs +++ b/scripts/ql3-postgres-ha-contract.cjs @@ -24,6 +24,7 @@ const { PostgresClusterScheduleRepository, PostgresRemoteWorkerCompletionRepository, PostgresRemoteWorkerLeaseControlRepository, + PostgresRunAttemptLogRetentionClaimRepository, PostgresToolInvocationArtifactRepository, PostgresWorkerSessionRepository, } = require('../packages/ql3-cluster-postgres/dist/entrypoints/runtime.js'); @@ -213,6 +214,9 @@ const { const { resolveClusterScheduleDecision, } = require('../packages/ql3-runtime-core/dist/scheduler/clusterScheduler.js'); +const { + createRunAttemptLogRetirementRecord, +} = require('../packages/ql3-runtime-core/dist/run/log-retention/runAttemptLogRetention.js'); const { PluginPackageManagementQuotaExceededError, PluginPackageManagementUnavailableError, @@ -9109,6 +9113,240 @@ async function assertAutomationManagementInspectionFailsClosed({ report.failedClosedWithoutSynchronousStandby = true; } +async function persistRunAttemptLogRetentionClaimBeforePromotion({ + primaryPort, + primaryDatabase, + standbyDatabase, +}) { + const fixture = Object.freeze({ + projectId: 'ha-log-retention-project', + runId: 'ha-log-retention-run', + attemptId: 'ha-log-retention-attempt', + logArtifactId: `wlog-${'f'.repeat(30)}`, + ownerId: 'ha-log-retention-primary', + token: 'ha-log-retention-primary-token-0001', + retentionMs: 60_000, + leaseMs: 5_000, + }); + const clock = await primaryDatabase.pool.query( + `SELECT floor(extract(epoch FROM clock_timestamp()) * 1000)::bigint::text + AS "observedAtMs"`, + ); + const observedAtMs = Number(clock.rows[0].observedAtMs); + const finishedAtMs = observedAtMs - 120_000; + await primaryDatabase.pool.query( + `INSERT INTO "ql3"."projects" ( + id, name, slug, status, version, created_at_ms, updated_at_ms + ) VALUES ($1, 'HA Log Retention', 'ha-log-retention', 'active', 1, $2, $2)`, + [fixture.projectId, finishedAtMs], + ); + await primaryDatabase.pool.query( + `INSERT INTO "ql3"."runs" ( + id, project_id, task_id, task_revision, trigger_type, + execution_origin, execution_owner, status, created_at_ms, + queued_at_ms, started_at_ms, finished_at_ms, version, event_sequence + ) VALUES ($1, $2, 'ha-log-retention-task', 'v1', 'manual', 'api', + 'runtime', 'succeeded', $3, $3, $3, $3, 3, 0)`, + [fixture.runId, fixture.projectId, finishedAtMs], + ); + await primaryDatabase.pool.query( + `INSERT INTO "ql3"."run_attempts" ( + id, run_id, attempt, status, executor_type, log_artifact_id, + callback_sequence, created_at_ms, started_at_ms, finished_at_ms, + exit_code + ) VALUES ($1, $2, 1, 'succeeded', 'remote_worker', $3, 0, + $4, $4, $4, 0)`, + [fixture.attemptId, fixture.runId, fixture.logArtifactId, finishedAtMs], + ); + + const runtimeDatabase = await databaseOpener( + 'runtime', + databaseUrl(RUNTIME_USER, RUNTIME_PASSWORD, primaryPort), + 'ql3-ha-log-retention-primary', + )(); + let claim; + try { + const repository = new PostgresRunAttemptLogRetentionClaimRepository( + runtimeDatabase.pool, + () => fixture.token, + ); + const page = await repository.claim({ + ownerId: fixture.ownerId, + retentionMs: fixture.retentionMs, + limit: 1, + leaseMs: fixture.leaseMs, + }); + assert.equal(page.claims.length, 1); + claim = page.claims[0]; + assert.deepEqual(claim.candidate, { + projectId: fixture.projectId, + runId: fixture.runId, + attemptId: fixture.attemptId, + logArtifactId: fixture.logArtifactId, + executorType: 'remote_worker', + finishedAtMs, + }); + assert.equal(claim.ownerId, fixture.ownerId); + assert.equal(claim.token, fixture.token); + assert.equal(claim.version, 1); + assert.equal(claim.failureCount, 0); + assert.equal(claim.expiresAtMs - claim.observedAtMs, fixture.leaseMs); + } finally { + await runtimeDatabase.close(); + } + + const replicated = await waitFor(async () => { + const result = await standbyDatabase.pool.query( + `SELECT claim_owner AS "ownerId", claim_token AS token, + claim_version AS version, claim_expires_at_ms::text AS "expiresAtMs", + (SELECT count(*)::integer + FROM "ql3"."run_attempt_log_artifact_tombstones" + WHERE attempt_id = $1) AS "tombstoneCount" + FROM "ql3"."run_attempt_log_retention_controls" + WHERE attempt_id = $1`, + [fixture.attemptId], + ); + return result.rowCount === 1 ? result.rows[0] : null; + }, 'Run Attempt log retention claim remote apply'); + assert.deepEqual(replicated, { + ownerId: fixture.ownerId, + token: fixture.token, + version: 1, + expiresAtMs: String(claim.expiresAtMs), + tombstoneCount: 0, + }); + + return { + fixture, + claim, + report: { + replicatedBeforePromotion: true, + initialOwnerId: claim.ownerId, + initialClaimVersion: claim.version, + initialClaimExpiresAtMs: claim.expiresAtMs, + initialTombstoneCount: replicated.tombstoneCount, + }, + }; +} + +async function verifyRunAttemptLogRetentionAfterPromotion({ + promotedPort, + promotedDatabase, + evidence, +}) { + const runtimeDatabase = await databaseOpener( + 'runtime', + databaseUrl(RUNTIME_USER, RUNTIME_PASSWORD, promotedPort), + 'ql3-ha-log-retention-promoted', + )(); + try { + const expired = await waitFor(async () => { + const result = await promotedDatabase.pool.query( + `SELECT floor(extract(epoch FROM clock_timestamp()) * 1000)::bigint::text + AS "observedAtMs"`, + ); + return Number(result.rows[0].observedAtMs) >= evidence.claim.expiresAtMs + ? result.rows[0] + : null; + }, 'Run Attempt log retention lease expiry'); + const oldRepository = new PostgresRunAttemptLogRetentionClaimRepository( + runtimeDatabase.pool, + ); + const staleRecord = createRunAttemptLogRetirementRecord({ + ...evidence.claim.candidate, + eligibleAtMs: evidence.claim.eligibleAtMs, + retiredAtMs: Math.max( + Number(expired.observedAtMs), + evidence.claim.eligibleAtMs, + ), + disposition: 'already_absent', + byteLength: 0, + truncation: { truncated: 'unknown' }, + }); + assert.equal( + await oldRepository.settle(evidence.claim, { + status: 'retired', + record: staleRecord, + }), + 'fenced', + ); + + const promotedToken = 'ha-log-retention-promoted-token-0002'; + const promotedOwnerId = 'ha-log-retention-promoted'; + const promotedRepository = + new PostgresRunAttemptLogRetentionClaimRepository( + runtimeDatabase.pool, + () => promotedToken, + ); + const page = await promotedRepository.claim({ + ownerId: promotedOwnerId, + retentionMs: evidence.fixture.retentionMs, + limit: 1, + leaseMs: evidence.fixture.leaseMs, + }); + assert.equal(page.claims.length, 1); + const promotedClaim = page.claims[0]; + assert.deepEqual(promotedClaim.candidate, evidence.claim.candidate); + assert.equal(promotedClaim.ownerId, promotedOwnerId); + assert.equal(promotedClaim.token, promotedToken); + assert.equal(promotedClaim.version, evidence.claim.version + 1); + assert.ok(promotedClaim.observedAtMs >= evidence.claim.expiresAtMs); + + const record = createRunAttemptLogRetirementRecord({ + ...promotedClaim.candidate, + eligibleAtMs: promotedClaim.eligibleAtMs, + retiredAtMs: Math.max( + promotedClaim.observedAtMs, + promotedClaim.eligibleAtMs, + ), + disposition: 'already_absent', + byteLength: 0, + truncation: { truncated: 'unknown' }, + }); + assert.equal( + await promotedRepository.settle(promotedClaim, { + status: 'retired', + record, + }), + 'settled', + ); + assert.deepEqual( + await promotedRepository.inspect({ + projectId: evidence.fixture.projectId, + runId: evidence.fixture.runId, + attemptId: evidence.fixture.attemptId, + logArtifactId: evidence.fixture.logArtifactId, + }), + { status: 'retired', record }, + ); + const durable = await promotedDatabase.pool.query( + `SELECT + (SELECT count(*)::integer + FROM "ql3"."run_attempt_log_retention_controls" + WHERE attempt_id = $1) AS "controlCount", + (SELECT count(*)::integer + FROM "ql3"."run_attempt_log_artifact_tombstones" + WHERE attempt_id = $1 AND record_digest = $2) AS "tombstoneCount"`, + [evidence.fixture.attemptId, record.recordDigest], + ); + assert.deepEqual(durable.rows, [{ controlCount: 0, tombstoneCount: 1 }]); + return { + ...evidence.report, + stalePrimarySettlementFenced: true, + promotedOwnerId, + promotedClaimVersion: promotedClaim.version, + promotedClaimObservedAtMs: promotedClaim.observedAtMs, + controlCountAfterSettlement: durable.rows[0].controlCount, + tombstoneCountAfterSettlement: durable.rows[0].tombstoneCount, + tombstoneDisposition: record.disposition, + recordDigest: record.recordDigest, + survivedPromotion: true, + }; + } finally { + await runtimeDatabase.close(); + } +} + async function main(argv = process.argv.slice(2)) { const reportFile = privateReportPath(argv); const nodeMajor = Number(process.versions.node.split('.')[0]); @@ -9192,6 +9430,8 @@ async function main(argv = process.argv.slice(2)) { let modelInvocationFeaturePromotion; let modelProviderCredentialCatalog; let modelProviderCredentialTestConnection; + let runAttemptLogRetentionEvidence; + let runAttemptLogRetention; const startedAt = performance.now(); const timeline = []; let report; @@ -10290,6 +10530,17 @@ async function main(argv = process.argv.slice(2)) { atMs: Number((performance.now() - startedAt).toFixed(3)), }); + runAttemptLogRetentionEvidence = + await persistRunAttemptLogRetentionClaimBeforePromotion({ + primaryPort, + primaryDatabase, + standbyDatabase, + }); + timeline.push({ + state: 'run_attempt_log_retention_claim_replicated', + atMs: Number((performance.now() - startedAt).toFixed(3)), + }); + await primaryDatabase.pool.query( `INSERT INTO "ql3"."projects" ( id, name, slug, status, version, created_at_ms, updated_at_ms @@ -10914,6 +11165,15 @@ async function main(argv = process.argv.slice(2)) { ? { afterRejoinMarkers, partitionOutcomeUnknownMarkers } : null; }, 'post-rewind synchronous WAL replay'); + runAttemptLogRetention = await verifyRunAttemptLogRetentionAfterPromotion({ + promotedPort: standbyPort, + promotedDatabase, + evidence: runAttemptLogRetentionEvidence, + }); + timeline.push({ + state: 'run_attempt_log_retention_tombstone_survived_promotion', + atMs: Number((performance.now() - startedAt).toFixed(3)), + }); const recoveredAutomationManagerDatabase = await databaseOpener( 'automation-manager', databaseUrl( @@ -11300,7 +11560,7 @@ async function main(argv = process.argv.slice(2)) { FROM "ql3"."worker_credential_deliveries") AS "credentialDeliveries"`, ); assert.deepEqual(sideEffects.rows, [ - { runs: 10, runEvents: 39, credentialDeliveries: 4 }, + { runs: 11, runEvents: 39, credentialDeliveries: 4 }, ]); timeline.push({ state: 'two_fresh_control_replicas_ready', @@ -11414,8 +11674,16 @@ async function main(argv = process.argv.slice(2)) { modelInvocationFeaturePromotion, modelProviderCredentialCatalog, modelProviderCredentialTestConnection, + runAttemptLogRetention, timeline, gates: { + runAttemptLogRetentionLeaseTakeoverAndTombstoneConverge: + runAttemptLogRetention.replicatedBeforePromotion && + runAttemptLogRetention.stalePrimarySettlementFenced && + runAttemptLogRetention.promotedClaimVersion === 2 && + runAttemptLogRetention.controlCountAfterSettlement === 0 && + runAttemptLogRetention.tombstoneCountAfterSettlement === 1 && + runAttemptLogRetention.survivedPromotion, packageAuthoritySplitReadinessBeforeAndAfterPromotion: true, optionalAiFeatureSchemaSurvivesPromotion: true, modelProviderCredentialCatalogSurvivesPromotion: diff --git a/test/back/ql3PackageBoundaryAudit.test.cjs b/test/back/ql3PackageBoundaryAudit.test.cjs index 65bcf5d8..d973c393 100644 --- a/test/back/ql3PackageBoundaryAudit.test.cjs +++ b/test/back/ql3PackageBoundaryAudit.test.cjs @@ -385,10 +385,10 @@ test('current QL3 workspace has exactly eighteen reviewed package boundaries', ( rootSourceFileRoles: clusterControl.rootSourceFileRoles, }, { - sourceFiles: 50, + sourceFiles: 51, rootSourceFiles: 2, rootSourceLines: 195, - nestedSourceFiles: 48, + nestedSourceFiles: 49, rootSourceFileRoles: { 'aiCli.ts': 'binary_entry', 'cli.ts': 'binary_entry',