diff --git a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md index 1ac7f3ad16..6b0e84eb98 100644 --- a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md +++ b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md @@ -3,15 +3,15 @@ - Status: Draft, under maintainer review - Initially proposed by: NoKV Lab - Widened by: LoopX maintainers -- Date: 2026-08-05; revised 2026-09-01 +- Date: 2026-08-05; revised 2026-09-02 - Scope: one provider-neutral LoopX authority contract with built-in file, optional NoKV, and optional PostgreSQL provider profiles, complementing [`host-integration-surface-v0`](../../reference/protocols/host-integration-surface-v0.md) - Source baseline: LoopX `a0c20f1779d273e7aaa4bd3ea166d145d466e6d5` -- Provider API baseline: NoKV `3d75d96965` (0.11.0 line). The Python - `publish_bytes` generation-CAS mapping was exercised once by hand against a - live NoKV stack at that pin (see the example README); the run is evidence for - the mapping only, not part of any merge gate +- Provider API baseline: NoKV `7bb3ffd6512fd57d9c0f193aa6d9c5b935d77f30` + (release 0.11.0, Python API 1, Holt pinned to 0.8.6). The Stage 2A executable + qualification admits only that SDK contract and this checkout's helper. It + remains candidate evidence, not a merge gate or authority promotion - PostgreSQL baseline: the TypeScript Stage 2B candidate implements the store contract and has passed a real PostgreSQL 16 transaction matrix. No shared authority service, runtime caller, authentication boundary, or authority @@ -1428,6 +1428,30 @@ claim that the complete P0 acceptance gate above passes. Historical latency or fault results are informative only; they are not a durability, recovery, HA, or production qualification claim. +The additional TEST ONLY Stage 2A probe in +`examples/nokv-authority-store/` opens three independent SDK helper processes +and checks fresh create, exact generation update, reconciliation after an +applied CAS response is deliberately lost, a one-winner/two-contender CAS, +winner/loser receipt behavior, and fresh-process receipt/history readback. Its +executable fixes argv to one absolute Python executable, the interpreter +isolation flag `-I`, and this checkout's reviewed helper, so `PYTHONPATH` +cannot substitute the `nokv` module; the helper fails closed unless the SDK +reports NoKV 0.11.0 and Python API 1, and the report repeats those two +admission constants rather than server-observed values. It validates read metadata against the current workbench +incarnation and validates publish responses against the requested workbench, +path, operation, revision, and generation. The AuthorityStore accepts even a +successful publish response only after a fresh read proves the exact persisted +transaction under the current workbench incarnation. This closes false success +at the LoopX boundary, but NoKV's current Python API does not atomically bind an +expected workspace incarnation into `publish_bytes`; preventing a write after a +concurrent remove/recreate remains an explicit provider-contract hold. +Only a successful live JSON report is evidence for that single-node Stage 2A +store-conformance run; deterministic tests are sequence tests only. This +LoopX-only candidate changes neither NoKV source nor its workbench/artifact data +model, and it does not prove runtime shadow parity, a multi-Agent canary, +authority-source promotion, HA, restart recovery, capacity, or production +routing. + ## Appendix B: Handoff-Mode Decision Record (2026-08-10) This appendix writes down a direction already agreed during the PR #2787 diff --git a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md index 0bb4d1d851..d0c260bba9 100644 --- a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md +++ b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md @@ -3,14 +3,15 @@ - 状态:Draft,正在接受 maintainer review - 最初提案方:NoKV Lab - 扩展修订方:LoopX maintainer -- 日期:2026-08-05;修订于 2026-09-01 +- 日期:2026-08-05;修订于 2026-09-02 - 范围:一个 provider-neutral 的 LoopX 权威合同,支持内置 file、可选 NoKV 与可选 PostgreSQL provider profile,用来补充 [`host-integration-surface-v0`](../../reference/protocols/host-integration-surface-v0.md) - 源码基线:LoopX `a0c20f1779d273e7aaa4bd3ea166d145d466e6d5` -- Provider API 基线:NoKV `3d75d96965`(0.11.0 线)。Python `publish_bytes` - generation-CAS 映射已在该基线的真实 NoKV stack 上手工跑过一次(见示例 README); - 该次运行只是映射本身的证据,不属于任何合并门槛 +- Provider API 基线:NoKV `7bb3ffd6512fd57d9c0f193aa6d9c5b935d77f30` + (release 0.11.0、Python API 1、Holt 固定为 0.8.6)。Stage 2A 的可执行资格 + 验证只接受这份 SDK 合同与本 checkout 的 helper;它仍是候选证据,不是合并门槛 + 或 authority promotion - PostgreSQL 基线:TypeScript Stage 2B candidate 已实现 store contract,且已通过 真实 PostgreSQL 16 transaction matrix;shared authority service、runtime caller、 authentication boundary 与 authority promotion 均尚未交付 @@ -1145,6 +1146,24 @@ migration/promotion、service recovery 或 HA。 验收门通过。历史 latency 或 fault 结果只具有参考意义,不构成 durability、recovery、 HA 或 production qualification 声明。 +`examples/nokv-authority-store/` 还包含一个 TEST ONLY 的 Stage 2A probe:它会 +打开三个相互独立的 SDK helper 进程,验证 fresh create、精确 generation update、 +一次 CAS 已落盘但响应被刻意丢弃后的回读 reconciliation、两个竞争者恰一胜出的 +CAS、胜负双方的 receipt 行为,以及新进程对 receipt 与完整 history 的回读。该可 +执行入口把 argv 固定为一个绝对 Python executable、解释器隔离标志 `-I`,加本 +checkout 中经过评审的 helper,因此 `PYTHONPATH` 无法替换 `nokv` 模块;helper 只 +接受 NoKV 0.11.0 / Python API 1,report 中重复的是这两个准入常量,而非从服务端读 +回的值。它把 read metadata 与当前 +workbench incarnation 对照,并针对请求的 workbench、path、operation、revision、 +generation 校验 publish 回包。即使 publish 回包报告成功,AuthorityStore 也必须 +重新读取并证明当前 workbench incarnation 下持久化了完全相同的 transaction,才会 +接受成功。这在 LoopX 边界消除了错误成功,但 NoKV 当前 Python API 尚不能把 +expected workspace incarnation 原子绑定到 `publish_bytes`;阻止 concurrent +remove/recreate 后的写入仍是明确的 provider-contract hold。只有成功的 live JSON report 才是该次单节点 Stage 2A store-conformance +运行的证据;确定性测试只证明场景序列。这个纯 LoopX 候选既不修改 NoKV 源码,也 +不改变其 workbench/artifact 数据模型;它不证明 runtime shadow parity、multi-Agent +canary、authority-source promotion、HA、重启恢复、容量或生产路由。 + ## 附录 B:交接模式决策记录(2026-08-10) 本附录把 PR #2787 评审中已同意的方向落成文字,作为实施前置条件的一部分。 diff --git a/examples/nokv-authority-store/README.md b/examples/nokv-authority-store/README.md new file mode 100644 index 0000000000..e8ed9d31e6 --- /dev/null +++ b/examples/nokv-authority-store/README.md @@ -0,0 +1,134 @@ +# Stage 2A NoKV AuthorityStore candidate qualification (TEST ONLY) + +This directory contains an explicit, write-producing single-node conformance +probe for the Stage 2A `NoKVAuthorityStore` candidate. It is **TEST ONLY**. +LoopX does not select NoKV by importing this code, and a successful run does +not connect a runtime shadow, run a multi-Agent canary, or flip an authority +source. + +The integration priority remains: + +1. the native, full NoKV CLI as the primary operator and production surface; +2. the NoKV Python SDK as the secondary programmable surface; and +3. optional sidecars only as adapters around those surfaces, never as the main + authority API. + +This probe intentionally exercises the current Python SDK bridge because it is +the available raw byte-CAS seam. It always starts this checkout's reviewed +`NoKVJsonLinesTransport` and `nokv_jsonl_helper.py`; it has no fake, skip, or +"unverified but successful" CLI path. The helper admits exactly NoKV SDK +`0.11.0` / Python API `1`, and the successful report repeats both values. + +## What it proves + +Against one **already existing** NoKV workbench, the probe starts three +independent helper processes and verifies: + +- the selected tenant/goal path is initially absent and can be created; +- the stored path generation advances from 1 to 2 under exact generation CAS; +- after the generation-2 CAS lands, an injected response loss is reconciled + from the durable authority envelope and operation receipt rather than from + the lower-layer response; +- every successful CAS response is accepted only after a fresh read proves the + exact transaction in the current workbench incarnation; +- two writes released together against generation 2 produce exactly one + generation-3 winner and one typed conflict; +- the losing operation does not acquire a durable receipt; and +- a third, freshly opened transport reads the winning envelope, its complete + three-entry history, the response-lost operation receipt, and the retained + winner receipt. + +If the SDK, helper, workbench, backend, CAS, or independent readback cannot be +proved, the process exits nonzero. The normal test suite uses deterministic +fakes only to test this sequence and does **not** count as live evidence. + +The probe does not prove an atomic expected-incarnation publication fence, +runtime shadow parity, a multi-Agent canary, authority promotion, HA, failover, +restart recovery, capacity, or performance. NoKV generation can restart after +workbench recreation, so the current adapter fails closed through authoritative +post-write readback; preventing the stale-incarnation write itself requires a +future provider primitive. The probe also does not create a workbench. A green +run is Stage 2A single-node storage conformance evidence only. + +## Inputs + +Use a current NoKV Python environment. Keep the client configuration in an +ignored local file; do not commit credentials. Static routing is valid for a +single-node NoKV deployment—etcd is not required by this probe. The following +shape is illustrative: + +```json +{ + "root_id": "00000000000000000000000000000000", + "routing": { + "kind": "static", + "endpoint": "127.0.0.1:7412", + "logical_shard_id": "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb", + "object_namespace_id": "cccccccccccccccccccccccccccccccc", + "placement_generation": 1, + "owner_epoch": 1 + }, + "object_store": { + "kind": "s3", + "bucket": "qualification-bucket", + "region": "us-east-1", + "root": "/loopx-qualification", + "endpoint": "http://127.0.0.1:9000", + "access_key_id": "set-in-your-ignored-local-file", + "secret_access_key": "set-in-your-ignored-local-file", + "virtual_host_style": false, + "skip_signature": false + } +} +``` + +Configuration objects are exact-key contracts. An unknown top-level, routing, +or object-store key fails before an SDK routing config, object-store config, or +client is constructed. In particular, a misspelled explicit credential cannot +silently fall through to NoKV's ambient provider chain. Intentionally omitted +optional S3 credential fields retain the NoKV SDK's normal behavior. + +Pass only the absolute path to the Python executable that resolves the qualified +NoKV SDK. The probe itself fixes the remaining argv to the interpreter isolation +flag `-I` followed by the reviewed helper in this checkout; callers cannot +supply a wrapper argument or an alternate helper path, and `PYTHONPATH`, +`PYTHONHOME`, or user site-packages cannot redirect the `nokv` import away from +that executable's own environment: + +```text +/path/to/nokv-python-environment/bin/python +``` + +Choose a fresh tenant/goal pair for every run. The probe refuses to overwrite +an existing authority envelope and deliberately leaves its three-generation +test envelope behind for inspection. Use a disposable qualification namespace +or remove it later with the native NoKV CLI according to that environment's +retention policy. + +## Run + +From the LoopX repository root: + +```bash +node --no-warnings --experimental-strip-types \ + examples/nokv-authority-store/live-qualification.ts \ + --execute-live \ + --config-json /path/to/ignored/nokv-client.json \ + --python-executable /path/to/nokv-python-environment/bin/python \ + --tenant-id qualification-tenant-20260902 \ + --goal-id qualification-goal-20260902-01 \ + --workbench existing-qualification-workbench +``` + +`--execute-live` is mandatory and is checked before any helper starts. Exit 0 +means every listed live check passed. Any unavailable, failed, ambiguous, +unfenced, pre-existing, or unreadable state exits nonzero with a compact JSON +reason; provider stderr, endpoints, credentials, and raw SDK errors are not +copied into that result. A successful JSON report includes +`"qualification_scope":"stage_2a_single_node_store_conformance"`, +`"nokv_sdk_version":"0.11.0"`, and `"nokv_api_version":1`. The two version +fields are the helper's admission constants: the helper refuses to open a client +for any other SDK version or API version, so a successful report implies them, +but they are not values read back from the NoKV server. The report is Stage +2A/helper-admission evidence only, not runtime-shadow, canary, HA, or +production-readiness evidence. diff --git a/examples/nokv-authority-store/live-qualification.ts b/examples/nokv-authority-store/live-qualification.ts new file mode 100644 index 0000000000..8301f67aac --- /dev/null +++ b/examples/nokv-authority-store/live-qualification.ts @@ -0,0 +1,629 @@ +#!/usr/bin/env -S node --no-warnings --experimental-strip-types + +import { randomUUID } from "node:crypto"; +import { readFile } from "node:fs/promises"; +import { isAbsolute, resolve } from "node:path"; +import { fileURLToPath } from "node:url"; +import { parseArgs } from "node:util"; + +import type { JsonObject } from "../../loopx/control_plane/effect_program.ts"; +import type { + AuthorityStoreCommit, + AuthorityStoreCommitResult, +} from "../../loopx/control_plane/coordination/authority_store.ts"; +import { + canonicalAuthorityObject, +} from "../../loopx/control_plane/coordination/authority_store_codec.ts"; +import { + NoKVAuthorityStore, + NoKVTransportProtocolError, + NoKVTransportUnavailableError, + type NoKVBlobCasRequest, + type NoKVBlobCasResult, + type NoKVBlobReadResult, + type NoKVBlobTransport, + type NoKVStoreIdentityResult, +} from "../../loopx/control_plane/coordination/nokv_authority_store.ts"; +import { + NoKVJsonLinesTransport, +} from "../../loopx/control_plane/coordination/nokv_jsonl_transport.ts"; + +const REPORT_SCHEMA = "loopx_nokv_authority_live_qualification_v0"; +export const QUALIFICATION_SCOPE = "stage_2a_single_node_store_conformance"; +export const QUALIFIED_NOKV_SDK_VERSION = "0.11.0"; +export const QUALIFIED_NOKV_API_VERSION = 1; +const REPOSITORY_HELPER = fileURLToPath( + new URL("../../loopx/control_plane/coordination/nokv_jsonl_helper.py", import.meta.url), +); +const COMPETITION_BARRIER_TIMEOUT_MS = 10_000; + +export class QualificationFailure extends Error { + readonly reasonCode: string; + + constructor(reasonCode: string, message: string) { + super(message); + this.reasonCode = reasonCode; + } +} + +export interface QualificationTransport extends NoKVBlobTransport { + close(): Promise; +} + +export interface QualificationOptions { + python_executable: string; + client_config: JsonObject; + tenant_id: string; + goal_id: string; + workbench: string; + request_timeout_ms?: number; +} + +export interface QualificationReport { + schema_version: typeof REPORT_SCHEMA; + qualification_scope: typeof QUALIFICATION_SCOPE; + ok: true; + checks: readonly { id: string; status: "passed" }[]; + final_generation: number; + final_cursor: string; + durable_test_data_left: true; + authority_source_changed: false; + availability_or_ha_proven: false; + nokv_sdk_version: typeof QUALIFIED_NOKV_SDK_VERSION; + nokv_api_version: typeof QUALIFIED_NOKV_API_VERSION; +} + +export interface QualificationSequenceResult { + checks: readonly { id: string; status: "passed" }[]; + final_generation: number; + final_cursor: string; +} + +export interface QualificationCliArguments { + configJsonPath: string; + pythonExecutable: string; + tenantId: string; + goalId: string; + workbench: string; + requestTimeoutMs?: number; +} + +type QualificationTransportFactory = () => Promise; + +function fail(reasonCode: string, message: string): never { + throw new QualificationFailure(reasonCode, message); +} + +function requiredString(value: unknown, name: string): string { + if (typeof value !== "string" || value.trim() !== value || value.length === 0) { + return fail("invalid_arguments", `${name} must be a non-empty trimmed string`); + } + return value; +} + +/** Interpreter flag: ignore PYTHONPATH, PYTHONHOME, and user site-packages. */ +export const HELPER_INTERPRETER_ISOLATION_FLAG = "-I"; + +/** + * Build the only helper argv accepted by the destructive live probe. + * + * The interpreter runs in isolated mode, so the `nokv` module can come only + * from the chosen executable's own environment. An environment variable cannot + * substitute a stand-in SDK for live evidence. + */ +export function qualificationHelperArgv( + pythonExecutable: string, +): readonly [string, string, string] { + const executable = requiredString(pythonExecutable, "Python executable"); + if (!isAbsolute(executable)) { + return fail("invalid_arguments", "Python executable must be an absolute path"); + } + return [executable, HELPER_INTERPRETER_ISOLATION_FLAG, REPOSITORY_HELPER]; +} + +function positiveSafeInteger(value: unknown, name: string): number { + const parsed = typeof value === "string" && /^[1-9]\d*$/.test(value) + ? Number(value) + : value; + if (!Number.isSafeInteger(parsed) || (parsed as number) < 1) { + return fail("invalid_arguments", `${name} must be a positive safe integer`); + } + return parsed as number; +} + +function expect( + condition: unknown, + reasonCode: string, + message: string, +): asserts condition { + if (!condition) fail(reasonCode, message); +} + +function commit( + expectedProviderRevision: string | null, + runId: string, + operationId: string, + sequence: number, + candidate: string, +): AuthorityStoreCommit { + return { + expected_provider_revision: expectedProviderRevision, + operation_id: operationId, + events: [{ + schema_version: "loopx_nokv_authority_qualification_event_v0", + run_id: runId, + sequence, + candidate, + }], + next_projection: { + schema_version: "loopx_nokv_authority_qualification_head_v0", + run_id: runId, + sequence, + candidate, + }, + receipts: [{ + schema_version: "loopx_nokv_authority_qualification_receipt_v0", + run_id: runId, + operation_id: operationId, + sequence, + candidate, + }], + }; +} + +class OneShotBarrier { + private arrivals = 0; + private readonly released: Promise; + private release!: () => void; + private readonly timer: NodeJS.Timeout; + + constructor() { + this.released = new Promise((resolveBarrier) => { + this.release = resolveBarrier; + }); + this.timer = setTimeout(() => { + this.release(); + }, COMPETITION_BARRIER_TIMEOUT_MS); + this.timer.unref(); + } + + async arrive(): Promise { + this.arrivals += 1; + if (this.arrivals === 2) { + clearTimeout(this.timer); + this.release(); + } + await this.released; + if (this.arrivals !== 2) { + fail("competition_barrier_failed", "both competitors did not reach the NoKV CAS"); + } + } +} + +class BarrierTransport implements NoKVBlobTransport { + readonly inner: NoKVBlobTransport; + readonly barrier: OneShotBarrier; + + constructor(inner: NoKVBlobTransport, barrier: OneShotBarrier) { + this.inner = inner; + this.barrier = barrier; + } + + async storeIdentity(workbench: string): Promise { + return await this.inner.storeIdentity(workbench); + } + + async readBlob(workbench: string, path: string): Promise { + return await this.inner.readBlob(workbench, path); + } + + async casPublishBlob(request: NoKVBlobCasRequest): Promise { + await this.barrier.arrive(); + return await this.inner.casPublishBlob(request); + } +} + +/** + * Turn one proven lower-layer success into an unknown transport outcome. + * + * The wrapper is deliberately above the real NoKV transport: the durable CAS + * still runs, while the TypeScript AuthorityStore must settle the result from + * the persisted operation receipt instead of trusting the lost response. + */ +class LoseOneAppliedResponseTransport implements NoKVBlobTransport { + readonly inner: NoKVBlobTransport; + appliedResponseDropped = false; + + constructor(inner: NoKVBlobTransport) { + this.inner = inner; + } + + async storeIdentity(workbench: string): Promise { + return await this.inner.storeIdentity(workbench); + } + + async readBlob(workbench: string, path: string): Promise { + return await this.inner.readBlob(workbench, path); + } + + async casPublishBlob(request: NoKVBlobCasRequest): Promise { + const result = await this.inner.casPublishBlob(request); + if (!this.appliedResponseDropped && result.status === "applied") { + this.appliedResponseDropped = true; + throw new NoKVTransportUnavailableError( + "qualification injected response loss after an applied NoKV CAS", + ); + } + return result; + } +} + +function applied( + result: AuthorityStoreCommitResult, + reasonCode: string, +): Extract { + expect(result.status === "applied", reasonCode, "authority commit was not applied"); + return result; +} + +async function rawGeneration( + transport: NoKVBlobTransport, + store: NoKVAuthorityStore, + expectedGeneration: number, + reasonCode: string, +): Promise { + const result = await transport.readBlob(store.workbench, store.path); + expect( + result.status === "loaded" && result.generation === expectedGeneration, + reasonCode, + `authority envelope did not read back at generation ${expectedGeneration}`, + ); +} + +/** + * Exercise the destructive sequence against three independent handles. + * + * This lower-level export intentionally does not return a live qualification + * report: tests may supply a deterministic transport factory, which proves the + * sequence but cannot prove a reachable NoKV backend. + */ +export async function exerciseQualificationSequence( + options: QualificationOptions, + openTransport: QualificationTransportFactory, +): Promise { + const tenantId = requiredString(options.tenant_id, "tenant id"); + const goalId = requiredString(options.goal_id, "goal id"); + const workbench = requiredString(options.workbench, "workbench"); + const opened: QualificationTransport[] = []; + const open = async (): Promise => { + const transport = await openTransport(); + opened.push(transport); + return transport; + }; + const checks: { id: string; status: "passed" }[] = []; + const passed = (id: string): void => { + checks.push({ id, status: "passed" }); + }; + const runId = randomUUID(); + const operationIds = { + create: `${runId}:create`, + advance: `${runId}:advance`, + contenderA: `${runId}:contender-a`, + contenderB: `${runId}:contender-b`, + }; + + try { + const firstTransport = await open(); + const secondTransport = await open(); + const first = new NoKVAuthorityStore(firstTransport, { + tenant_id: tenantId, + goal_id: goalId, + workbench, + }); + const responseLossTransport = new LoseOneAppliedResponseTransport(secondTransport); + const second = new NoKVAuthorityStore(responseLossTransport, { + tenant_id: tenantId, + goal_id: goalId, + workbench, + }); + + const [firstIdentity, secondIdentity] = await Promise.all([ + first.storeIdentity(), + second.storeIdentity(), + ]); + expect( + firstIdentity.status === "available" && + secondIdentity.status === "available" && + firstIdentity.store_identity === secondIdentity.store_identity, + "workbench_identity_failed", + "independent transports did not resolve the same existing workbench", + ); + passed("existing_workbench_identity"); + + const [firstInitial, secondInitial] = await Promise.all([ + first.loadAuthority(), + second.loadAuthority(), + ]); + expect( + firstInitial.status === "missing" && secondInitial.status === "missing", + "qualification_target_not_fresh", + "qualification requires a fresh tenant and goal target", + ); + passed("fresh_authority_target"); + + const created = applied( + await first.commitAuthority( + commit(null, runId, operationIds.create, 1, "create"), + ), + "create_failed", + ); + passed("create_applied"); + await rawGeneration(firstTransport, first, 1, "create_generation_failed"); + passed("create_generation_one"); + + const advanced = applied( + await second.commitAuthority( + commit(created.provider_revision, runId, operationIds.advance, 2, "advance"), + ), + "generation_cas_failed", + ); + expect( + responseLossTransport.appliedResponseDropped, + "response_loss_not_injected", + "qualification did not discard an applied NoKV CAS response", + ); + passed("response_lost_success_reconciled"); + passed("generation_cas_applied"); + await rawGeneration(firstTransport, first, 2, "generation_two_readback_failed"); + passed("generation_two_readback"); + + const barrier = new OneShotBarrier(); + const contenderA = new NoKVAuthorityStore( + new BarrierTransport(firstTransport, barrier), + { tenant_id: tenantId, goal_id: goalId, workbench }, + ); + const contenderB = new NoKVAuthorityStore( + new BarrierTransport(secondTransport, barrier), + { tenant_id: tenantId, goal_id: goalId, workbench }, + ); + const race = await Promise.all([ + contenderA.commitAuthority( + commit(advanced.provider_revision, runId, operationIds.contenderA, 3, "a"), + ), + contenderB.commitAuthority( + commit(advanced.provider_revision, runId, operationIds.contenderB, 3, "b"), + ), + ]); + const winnerIndex = race.findIndex((result) => result.status === "applied"); + const loserIndex = race.findIndex((result) => result.status === "conflict"); + expect( + winnerIndex >= 0 && loserIndex >= 0 && winnerIndex !== loserIndex && + race.filter((result) => result.status === "applied").length === 1 && + race.filter((result) => result.status === "conflict").length === 1, + "competition_not_fenced", + "competing generation CAS did not produce exactly one winner", + ); + const winner = race[winnerIndex]!; + const loser = race[loserIndex]!; + expect( + winner.status === "applied" && + loser.status === "conflict" && + loser.conflict_kind === "provider_revision_mismatch" && + loser.current_provider_revision === winner.provider_revision && + loser.current_cursor === "3", + "competition_not_fenced", + "competition conflict was not bound to the winning generation", + ); + passed("competing_generation_cas_one_winner"); + await rawGeneration(secondTransport, second, 3, "competition_generation_failed"); + passed("competition_did_not_double_advance"); + + await Promise.all([firstTransport.close(), secondTransport.close()]); + + const readbackTransport = await open(); + const readback = new NoKVAuthorityStore(readbackTransport, { + tenant_id: tenantId, + goal_id: goalId, + workbench, + }); + const loaded = await readback.loadAuthority(); + expect( + loaded.status === "loaded" && + loaded.provider_revision === winner.provider_revision && + loaded.cursor === "3" && + loaded.head.run_id === runId && + loaded.head.sequence === 3, + "independent_readback_failed", + "a fresh transport did not read the winning authority envelope", + ); + await rawGeneration( + readbackTransport, + readback, + 3, + "independent_readback_failed", + ); + const history = await readback.scanCommitted(null, 10); + expect( + history.status === "page" && + history.transactions.length === 3 && + history.next_cursor === "3" && + history.has_more === false, + "independent_readback_failed", + "a fresh transport did not read the complete committed history", + ); + passed("independent_transport_readback"); + + const winnerOperation = winnerIndex === 0 + ? operationIds.contenderA + : operationIds.contenderB; + const loserOperation = winnerIndex === 0 + ? operationIds.contenderB + : operationIds.contenderA; + const [ambiguousCommitReceipt, winnerReceipt, loserReceipt] = await Promise.all([ + readback.readReceipt(operationIds.advance), + readback.readReceipt(winnerOperation), + readback.readReceipt(loserOperation), + ]); + expect( + ambiguousCommitReceipt.status === "found" && + ambiguousCommitReceipt.cursor === "2" && + ambiguousCommitReceipt.receipts[0]?.operation_id === operationIds.advance, + "ambiguous_commit_receipt_missing", + "the response-lost operation receipt was not retained", + ); + passed("ambiguous_commit_receipt_retained"); + expect( + winnerReceipt.status === "found" && winnerReceipt.cursor === "3", + "winner_receipt_missing", + "the winning operation receipt was not retained", + ); + passed("winner_receipt_retained"); + expect( + loserReceipt.status === "missing", + "loser_receipt_present", + "the losing operation unexpectedly acquired a durable receipt", + ); + passed("loser_receipt_absent"); + + return { + checks, + final_generation: 3, + final_cursor: "3", + }; + } finally { + await Promise.allSettled(opened.map(async (transport) => await transport.close())); + } +} + +/** Run the live probe only through this checkout's reviewed JSONL transport. */ +export async function qualifyNoKVAuthorityStore( + options: QualificationOptions, +): Promise { + const sequence = await exerciseQualificationSequence(options, async () => + await NoKVJsonLinesTransport.open({ + argv: qualificationHelperArgv(options.python_executable), + config: options.client_config, + request_timeout_ms: options.request_timeout_ms, + })); + return { + schema_version: REPORT_SCHEMA, + qualification_scope: QUALIFICATION_SCOPE, + ok: true, + ...sequence, + durable_test_data_left: true, + authority_source_changed: false, + availability_or_ha_proven: false, + nokv_sdk_version: QUALIFIED_NOKV_SDK_VERSION, + nokv_api_version: QUALIFIED_NOKV_API_VERSION, + }; +} + +export function parseQualificationArguments( + argv: readonly string[], +): QualificationCliArguments { + let values: ReturnType["values"]; + try { + values = parseArgs({ + args: [...argv], + strict: true, + allowPositionals: false, + options: { + "execute-live": { type: "boolean", default: false }, + "config-json": { type: "string" }, + "python-executable": { type: "string" }, + "tenant-id": { type: "string" }, + "goal-id": { type: "string" }, + workbench: { type: "string" }, + "request-timeout-ms": { type: "string" }, + }, + }).values; + } catch { + return fail("invalid_arguments", "qualification arguments are invalid"); + } + if (values["execute-live"] !== true) { + return fail( + "live_opt_in_required", + "--execute-live is required because this probe writes durable test data", + ); + } + return { + configJsonPath: requiredString(values["config-json"], "--config-json"), + pythonExecutable: qualificationHelperArgv( + requiredString(values["python-executable"], "--python-executable"), + )[0], + tenantId: requiredString(values["tenant-id"], "--tenant-id"), + goalId: requiredString(values["goal-id"], "--goal-id"), + workbench: requiredString(values.workbench, "--workbench"), + requestTimeoutMs: values["request-timeout-ms"] === undefined + ? undefined + : positiveSafeInteger(values["request-timeout-ms"], "--request-timeout-ms"), + }; +} + +async function readJson(path: string, reasonCode: string): Promise { + let bytes: string; + try { + bytes = await readFile(path, "utf8"); + } catch { + return fail(reasonCode, "qualification JSON input could not be read"); + } + try { + return JSON.parse(bytes); + } catch { + return fail(reasonCode, "qualification JSON input is invalid"); + } +} + +async function loadQualificationOptions( + cli: QualificationCliArguments, +): Promise { + const configValue = await readJson(cli.configJsonPath, "config_json_invalid"); + let clientConfig: JsonObject; + try { + clientConfig = canonicalAuthorityObject(configValue, "NoKV client config"); + } catch { + return fail("config_json_invalid", "NoKV client config must be strict JSON"); + } + return { + python_executable: cli.pythonExecutable, + client_config: clientConfig, + tenant_id: cli.tenantId, + goal_id: cli.goalId, + workbench: cli.workbench, + request_timeout_ms: cli.requestTimeoutMs, + }; +} + +async function main(): Promise { + try { + const cli = parseQualificationArguments(process.argv.slice(2)); + const options = await loadQualificationOptions(cli); + const report = await qualifyNoKVAuthorityStore(options); + process.stdout.write(`${JSON.stringify(report)}\n`); + return 0; + } catch (error) { + let reasonCode = "qualification_failed"; + let reason = "NoKV authority-store qualification failed"; + if (error instanceof QualificationFailure) { + reasonCode = error.reasonCode; + reason = error.message; + } else if (error instanceof NoKVTransportUnavailableError) { + reasonCode = "nokv_backend_unavailable"; + reason = "NoKV backend or SDK helper is unavailable"; + } else if (error instanceof NoKVTransportProtocolError) { + reasonCode = "nokv_transport_protocol_failed"; + reason = "NoKV SDK helper violated the transport protocol"; + } + process.stderr.write(`${JSON.stringify({ + schema_version: REPORT_SCHEMA, + ok: false, + reason_code: reasonCode, + reason, + })}\n`); + return 1; + } +} + +if (process.argv[1] && resolve(process.argv[1]) === resolve(fileURLToPath(import.meta.url))) { + process.exitCode = await main(); +} diff --git a/loopx/control_plane/coordination/authority_store.ts b/loopx/control_plane/coordination/authority_store.ts index 124d4bdf6c..21727d434a 100644 --- a/loopx/control_plane/coordination/authority_store.ts +++ b/loopx/control_plane/coordination/authority_store.ts @@ -60,6 +60,7 @@ export const AUTHORITY_STORE_PROVIDER_PROFILES = { trust_boundary: "loopx_authority_owned_nokv_credentials", qualification_holds: [ "service_grade_contract_adapter", + "atomic_workspace_incarnation_publication_fence", "restart_and_restore_recovery", "capacity_and_receipt_retention", "availability_and_ha", diff --git a/loopx/control_plane/coordination/nokv_authority_store.ts b/loopx/control_plane/coordination/nokv_authority_store.ts new file mode 100644 index 0000000000..5775680caa --- /dev/null +++ b/loopx/control_plane/coordination/nokv_authority_store.ts @@ -0,0 +1,671 @@ +import { createHash, randomUUID } from "node:crypto"; +import { TextDecoder } from "node:util"; + +import type { JsonObject } from "../effect_program.ts"; +import type { + AuthorityStore, + AuthorityStoreCommit, + AuthorityStoreCommittedTransaction, + AuthorityStoreCommitResult, + AuthorityStoreIdentityResult, + AuthorityStoreLoadResult, + AuthorityStoreReadFailure, + AuthorityStoreReceiptResult, + AuthorityStoreScanResult, +} from "./authority_store.ts"; +import { + AuthorityStoreProtocolError, + canonicalAuthorityBytes, + canonicalAuthorityObject, + canonicalAuthorityObjectList, + hasExactAuthorityKeys, + isAuthorityJsonObject, + normalizeAuthorityStoreCommit, + parseAuthorityCursor, + requireAuthorityStoreId, +} from "./authority_store_codec.ts"; + +const NOKV_AUTHORITY_STORE_SCHEMA = "loopx_nokv_authority_store_v0"; +const DEFAULT_MAX_ENVELOPE_BYTES = 16 * 1024 * 1024; +const HEX_128_PATTERN = /^[0-9a-f]{32}$/; + +export class NoKVTransportUnavailableError extends Error {} +export class NoKVTransportProtocolError extends Error {} + +export type NoKVTransportFailure = { + status: "unavailable" | "failed"; + reason_code: string; + reason: string; +}; + +export type NoKVStoreIdentityResult = + | { status: "available"; store_identity: string } + | NoKVTransportFailure; + +export type NoKVBlobReadResult = + | { status: "loaded"; bytes: Uint8Array; generation: number } + | { status: "missing" } + | NoKVTransportFailure; + +export interface NoKVBlobCasRequest { + workbench: string; + path: string; + expected_generation: number | null; + bytes: Uint8Array; + operation_id: string; + artifact_revision_id: string; +} + +export type NoKVBlobCasResult = + | { status: "applied"; generation: number } + | { status: "conflict"; current_generation: number | null } + | { status: "ambiguous"; reason_code: string; reason: string } + | { status: "failed"; reason_code: string; reason: string }; + +/** Raw byte-storage contract implemented by the Python SDK JSON-lines helper. */ +export interface NoKVBlobTransport { + storeIdentity(workbench: string): Promise; + readBlob(workbench: string, path: string): Promise; + casPublishBlob(request: NoKVBlobCasRequest): Promise; +} + +export interface NoKVAuthorityStoreOptions { + tenant_id: string; + goal_id: string; + workbench: string; + max_envelope_bytes?: number; +} + +interface NoKVAuthorityStoreDocument extends JsonObject { + schema_version: typeof NOKV_AUTHORITY_STORE_SCHEMA; + tenant_id: string; + goal_id: string; + store_identity: string; + storage_generation: number; + provider_revision: string; + cursor: string; + head: JsonObject; + committed: AuthorityStoreCommittedTransaction[]; +} + +type EnvelopeReadResult = + | { + status: "loaded"; + identity: string; + generation: number; + document: NoKVAuthorityStoreDocument; + } + | { status: "missing"; identity: string } + | AuthorityStoreReadFailure; + +function cloneTransaction( + value: AuthorityStoreCommittedTransaction, +): AuthorityStoreCommittedTransaction { + return structuredClone(value); +} + +function transactionWithoutRevision(value: AuthorityStoreCommittedTransaction) { + return { + cursor: value.cursor, + operation_id: value.operation_id, + events: value.events, + projection: value.projection, + receipts: value.receipts, + }; +} + +function providerRevision( + tenantId: string, + goalId: string, + storeIdentity: string, + storageGeneration: number, + previousRevision: string | null, + transaction: ReturnType, +): string { + const digest = createHash("sha256") + .update(canonicalAuthorityBytes({ + provider: "nokv", + tenant_id: tenantId, + goal_id: goalId, + store_identity: storeIdentity, + storage_generation: storageGeneration, + previous_provider_revision: previousRevision, + transaction, + })) + .digest("hex") + .slice(0, 24); + return `nokv:${transaction.cursor}:${digest}`; +} + +function physicalAttemptIdentity( + domain: "operation" | "artifact_revision", + attemptNonce: string, + operationId: string, + expectedGeneration: number | null, + payload: Uint8Array, +): string { + return createHash("sha256") + .update(`loopx.nokv.${domain}.v1\0`, "utf8") + .update(attemptNonce, "utf8") + .update("\0", "utf8") + .update(operationId, "utf8") + .update("\0", "utf8") + .update(expectedGeneration === null ? "create" : String(expectedGeneration), "utf8") + .update("\0", "utf8") + .update(createHash("sha256").update(payload).digest()) + .digest("hex") + .slice(0, 32); +} + +function requireGeneration(value: unknown, name: string): number { + if (!Number.isSafeInteger(value) || (value as number) < 1) { + throw new AuthorityStoreProtocolError(`${name} must be a positive safe integer`); + } + return value as number; +} + +function decodeTransaction(value: unknown): AuthorityStoreCommittedTransaction { + if (!isAuthorityJsonObject(value) || !hasExactAuthorityKeys(value, [ + "cursor", "provider_revision", "operation_id", "events", "projection", "receipts", + ])) { + throw new AuthorityStoreProtocolError("committed transaction is invalid"); + } + return { + cursor: requireAuthorityStoreId(value.cursor, "transaction cursor"), + provider_revision: requireAuthorityStoreId( + value.provider_revision, + "transaction provider revision", + ), + operation_id: requireAuthorityStoreId(value.operation_id, "operation id"), + events: canonicalAuthorityObjectList(value.events, "transaction events"), + projection: canonicalAuthorityObject(value.projection, "transaction projection"), + receipts: canonicalAuthorityObjectList(value.receipts, "transaction receipts"), + }; +} + +function decodeDocument( + value: unknown, + tenantId: string, + goalId: string, + storeIdentity: string, + observedGeneration: number, +): NoKVAuthorityStoreDocument { + if (!isAuthorityJsonObject(value) || !hasExactAuthorityKeys(value, [ + "schema_version", "tenant_id", "goal_id", "store_identity", "storage_generation", + "provider_revision", "cursor", "head", "committed", + ]) || value.schema_version !== NOKV_AUTHORITY_STORE_SCHEMA) { + throw new AuthorityStoreProtocolError("NoKV authority store schema mismatch"); + } + if (value.tenant_id !== tenantId) { + throw new AuthorityStoreProtocolError("NoKV authority store tenant mismatch"); + } + if (value.goal_id !== goalId) { + throw new AuthorityStoreProtocolError("NoKV authority store goal mismatch"); + } + if (value.store_identity !== storeIdentity) { + throw new AuthorityStoreProtocolError("NoKV authority store lineage mismatch"); + } + const storageGeneration = requireGeneration( + value.storage_generation, + "NoKV storage generation", + ); + if (storageGeneration !== observedGeneration) { + throw new AuthorityStoreProtocolError( + "NoKV authority store storage generation does not match read metadata", + ); + } + const revision = requireAuthorityStoreId(value.provider_revision, "provider revision"); + const cursor = requireAuthorityStoreId(value.cursor, "provider cursor"); + const head = canonicalAuthorityObject(value.head, "NoKV authority store head"); + if (!Array.isArray(value.committed)) { + throw new AuthorityStoreProtocolError("NoKV authority store history is invalid"); + } + const committed = value.committed.map(decodeTransaction); + if ( + committed.length === 0 || + storageGeneration !== committed.length || + parseAuthorityCursor(cursor) !== BigInt(committed.length) + ) { + throw new AuthorityStoreProtocolError("NoKV authority store generation lineage is invalid"); + } + let previousRevision: string | null = null; + const operationIds = new Set(); + for (const [index, entry] of committed.entries()) { + const generation = index + 1; + if (parseAuthorityCursor(entry.cursor) !== BigInt(generation)) { + throw new AuthorityStoreProtocolError("NoKV authority store cursor lineage is invalid"); + } + if (operationIds.has(entry.operation_id)) { + throw new AuthorityStoreProtocolError( + "NoKV authority store operation identity is duplicated", + ); + } + operationIds.add(entry.operation_id); + const expectedRevision = providerRevision( + tenantId, + goalId, + storeIdentity, + generation, + previousRevision, + transactionWithoutRevision(entry), + ); + if (entry.provider_revision !== expectedRevision) { + throw new AuthorityStoreProtocolError("NoKV authority store revision lineage is invalid"); + } + previousRevision = entry.provider_revision; + } + const last = committed.at(-1)!; + if ( + last.cursor !== cursor || + last.provider_revision !== revision || + !canonicalAuthorityBytes(last.projection).equals(canonicalAuthorityBytes(head)) + ) { + throw new AuthorityStoreProtocolError("NoKV authority store head lineage is invalid"); + } + return { + schema_version: NOKV_AUTHORITY_STORE_SCHEMA, + tenant_id: tenantId, + goal_id: goalId, + store_identity: storeIdentity, + storage_generation: storageGeneration, + provider_revision: revision, + cursor, + head, + committed, + }; +} + +function readFailure(error: unknown): AuthorityStoreReadFailure { + if ( + error instanceof AuthorityStoreProtocolError || + error instanceof NoKVTransportProtocolError || + error instanceof SyntaxError + ) { + return { + status: "failed", + reason_code: "provider_protocol_violation", + reason: error.message, + }; + } + return { + status: "unavailable", + reason_code: "nokv_transport_unavailable", + reason: error instanceof Error ? error.message : "NoKV transport unavailable", + }; +} + +function validStoreIdentity(value: string, workbench: string): boolean { + const prefix = `nokv:${workbench}:`; + return value.startsWith(prefix) && HEX_128_PATTERN.test(value.slice(prefix.length)); +} + +/** Stage 2A candidate. No runtime constructs this provider by default. */ +export class NoKVAuthorityStore implements AuthorityStore { + readonly transport: NoKVBlobTransport; + readonly tenantId: string; + readonly goalId: string; + readonly workbench: string; + readonly maxEnvelopeBytes: number; + readonly path: string; + + constructor(transport: NoKVBlobTransport, options: NoKVAuthorityStoreOptions) { + this.transport = transport; + this.tenantId = requireAuthorityStoreId(options.tenant_id, "tenant id"); + this.goalId = requireAuthorityStoreId(options.goal_id, "goal id"); + this.workbench = requireAuthorityStoreId(options.workbench, "workbench"); + this.maxEnvelopeBytes = options.max_envelope_bytes ?? DEFAULT_MAX_ENVELOPE_BYTES; + if (!Number.isSafeInteger(this.maxEnvelopeBytes) || this.maxEnvelopeBytes < 1) { + throw new AuthorityStoreProtocolError( + "max envelope bytes must be a positive safe integer", + ); + } + const digest = createHash("sha256") + .update(canonicalAuthorityBytes({ + tenant_id: this.tenantId, + goal_id: this.goalId, + })) + .digest("hex") + .slice(0, 32); + this.path = `metadata/loopx-authority/${digest}.json`; + } + + async storeIdentity(): Promise { + try { + const result = await this.transport.storeIdentity(this.workbench); + if (result.status !== "available") return result; + if (!validStoreIdentity(result.store_identity, this.workbench)) { + return { + status: "failed", + reason_code: "store_identity_invalid", + reason: "NoKV store identity does not match the bound workbench incarnation", + }; + } + return result; + } catch (error) { + return readFailure(error); + } + } + + private async readEnvelope(): Promise { + const identityResult = await this.storeIdentity(); + if (identityResult.status !== "available") return identityResult; + let result: NoKVBlobReadResult; + try { + result = await this.transport.readBlob(this.workbench, this.path); + } catch (error) { + return readFailure(error); + } + if (result.status === "missing") { + return { status: "missing", identity: identityResult.store_identity }; + } + if (result.status !== "loaded") return result; + try { + const generation = requireGeneration(result.generation, "NoKV read generation"); + let raw: string; + try { + raw = new TextDecoder("utf-8", { fatal: true }).decode(result.bytes); + } catch (error) { + throw new AuthorityStoreProtocolError( + `NoKV authority store bytes are not UTF-8: ${ + error instanceof Error ? error.message : "invalid UTF-8" + }`, + ); + } + let value: unknown; + try { + value = JSON.parse(raw); + } catch (error) { + throw new AuthorityStoreProtocolError( + `NoKV authority store JSON is invalid: ${ + error instanceof Error ? error.message : "invalid JSON" + }`, + ); + } + return { + status: "loaded", + identity: identityResult.store_identity, + generation, + document: decodeDocument( + value, + this.tenantId, + this.goalId, + identityResult.store_identity, + generation, + ), + }; + } catch (error) { + return readFailure(error); + } + } + + async loadAuthority(): Promise { + const result = await this.readEnvelope(); + if (result.status === "missing") return { status: "missing" }; + if (result.status !== "loaded") return result; + return { + status: "loaded", + head: structuredClone(result.document.head), + provider_revision: result.document.provider_revision, + cursor: result.document.cursor, + }; + } + + private async settleCommitFromReadback( + expectedProviderRevision: string | null, + intended: AuthorityStoreCommittedTransaction, + fallback: Extract, + ): Promise { + const observed = await this.readEnvelope(); + if (observed.status === "failed") return observed; + if (observed.status !== "loaded") return fallback; + const sameOperation = observed.document.committed.find( + (entry) => entry.operation_id === intended.operation_id, + ); + if (sameOperation) { + if ( + canonicalAuthorityBytes(sameOperation).equals(canonicalAuthorityBytes(intended)) + ) { + return { + status: "applied", + provider_revision: sameOperation.provider_revision, + cursor: sameOperation.cursor, + }; + } + return { + status: "conflict", + conflict_kind: "operation_id_exists", + current_provider_revision: observed.document.provider_revision, + current_cursor: observed.document.cursor, + }; + } + if (observed.document.provider_revision !== expectedProviderRevision) { + return { + status: "conflict", + conflict_kind: "provider_revision_mismatch", + current_provider_revision: observed.document.provider_revision, + current_cursor: observed.document.cursor, + }; + } + return fallback; + } + + async commitAuthority(commit: AuthorityStoreCommit): Promise { + let normalized: AuthorityStoreCommit; + try { + normalized = normalizeAuthorityStoreCommit(commit); + } catch (error) { + return { + status: "failed", + reason_code: "invalid_commit_request", + reason: error instanceof Error ? error.message : "invalid commit request", + }; + } + const current = await this.readEnvelope(); + if (current.status !== "loaded" && current.status !== "missing") { + return { + status: "failed", + reason_code: current.reason_code, + reason: current.reason, + }; + } + const currentDocument = current.status === "loaded" ? current.document : null; + if ((currentDocument?.provider_revision ?? null) !== normalized.expected_provider_revision) { + return { + status: "conflict", + conflict_kind: "provider_revision_mismatch", + current_provider_revision: currentDocument?.provider_revision ?? null, + current_cursor: currentDocument?.cursor ?? null, + }; + } + if ( + currentDocument?.committed.some( + (entry) => entry.operation_id === normalized.operation_id, + ) + ) { + return { + status: "conflict", + conflict_kind: "operation_id_exists", + current_provider_revision: currentDocument.provider_revision, + current_cursor: currentDocument.cursor, + }; + } + const cursor = (parseAuthorityCursor(currentDocument?.cursor ?? null) + 1n).toString(); + const generation = (current.status === "loaded" ? current.generation : 0) + 1; + const base = { + cursor, + operation_id: normalized.operation_id, + events: normalized.events, + projection: normalized.next_projection, + receipts: normalized.receipts, + }; + const revision = providerRevision( + this.tenantId, + this.goalId, + current.identity, + generation, + currentDocument?.provider_revision ?? null, + base, + ); + const transaction: AuthorityStoreCommittedTransaction = { + ...base, + provider_revision: revision, + }; + const document: NoKVAuthorityStoreDocument = { + schema_version: NOKV_AUTHORITY_STORE_SCHEMA, + tenant_id: this.tenantId, + goal_id: this.goalId, + store_identity: current.identity, + storage_generation: generation, + provider_revision: revision, + cursor, + head: normalized.next_projection, + committed: [...(currentDocument?.committed ?? []), transaction], + }; + const payload = canonicalAuthorityBytes(document); + if (payload.byteLength > this.maxEnvelopeBytes) { + return { + status: "failed", + reason_code: "authority_envelope_too_large", + reason: `authority envelope exceeds ${this.maxEnvelopeBytes} bytes`, + }; + } + const expectedGeneration = current.status === "loaded" ? current.generation : null; + // NoKV publication identities are terminal after a failed or quarantined + // attempt. Keep the LoopX operation id stable in the authority envelope, + // while giving each physical retry a fresh pair of lower-layer ids. A + // response-lost success is still settled only by reading that envelope. + const attemptNonce = randomUUID(); + let result: NoKVBlobCasResult; + try { + result = await this.transport.casPublishBlob({ + workbench: this.workbench, + path: this.path, + expected_generation: expectedGeneration, + bytes: payload, + operation_id: physicalAttemptIdentity( + "operation", + attemptNonce, + normalized.operation_id, + expectedGeneration, + payload, + ), + artifact_revision_id: physicalAttemptIdentity( + "artifact_revision", + attemptNonce, + normalized.operation_id, + expectedGeneration, + payload, + ), + }); + } catch (error) { + result = { + status: "ambiguous", + reason_code: "nokv_transport_lost", + reason: error instanceof Error ? error.message : "NoKV transport outcome unknown", + }; + } + if (result.status === "applied" && result.generation === generation) { + // Generation is not a workbench-incarnation fence: NoKV may restart it + // after remove/recreate. Never expose success until a fresh read proves + // this exact transaction in the current incarnation. Preventing the + // stale-incarnation write itself still requires an atomic provider + // primitive that accepts the expected incarnation. + return await this.settleCommitFromReadback( + normalized.expected_provider_revision, + transaction, + { + status: "ambiguous", + reason_code: "nokv_applied_readback_unproved", + reason: "NoKV applied response requires current-incarnation readback", + }, + ); + } + if (result.status === "failed") return result; + const fallback: Extract = + result.status === "ambiguous" + ? result + : { + status: "ambiguous", + reason_code: result.status === "conflict" + ? "nokv_cas_conflict_unresolved" + : "nokv_publish_response_invalid", + reason: result.status === "conflict" + ? "NoKV CAS conflict requires authority-envelope readback" + : "NoKV publish generation did not match the attempted envelope", + }; + return await this.settleCommitFromReadback( + normalized.expected_provider_revision, + transaction, + fallback, + ); + } + + async readReceipt(operationId: string): Promise { + let normalized: string; + try { + normalized = requireAuthorityStoreId(operationId, "operation id"); + } catch (error) { + return { + status: "failed", + reason_code: "invalid_operation_id", + reason: error instanceof Error ? error.message : "invalid operation id", + }; + } + const result = await this.readEnvelope(); + if (result.status === "missing") return { status: "missing" }; + if (result.status !== "loaded") return result; + const transaction = result.document.committed.find( + (entry) => entry.operation_id === normalized, + ); + return transaction + ? { + status: "found", + cursor: transaction.cursor, + provider_revision: transaction.provider_revision, + receipts: structuredClone(transaction.receipts), + } + : { status: "missing" }; + } + + async scanCommitted( + afterCursor: string | null, + limit: number, + ): Promise { + let offset: bigint; + try { + offset = parseAuthorityCursor(afterCursor); + if (!Number.isSafeInteger(limit) || limit < 1) { + throw new AuthorityStoreProtocolError("scan limit must be a positive safe integer"); + } + } catch (error) { + return { + status: "failed", + reason_code: "invalid_scan_request", + reason: error instanceof Error ? error.message : "invalid scan request", + }; + } + const result = await this.readEnvelope(); + if (result.status === "missing") { + return { status: "page", transactions: [], next_cursor: afterCursor, has_more: false }; + } + if (result.status !== "loaded") return result; + const headCursor = parseAuthorityCursor(result.document.cursor); + if (offset > headCursor || offset > BigInt(Number.MAX_SAFE_INTEGER)) { + return { + status: "failed", + reason_code: "scan_cursor_out_of_range", + reason: "scan cursor is ahead of the provider head", + }; + } + const start = Number(offset); + const transactions = result.document.committed + .slice(start, start + limit) + .map(cloneTransaction); + return { + status: "page", + transactions, + next_cursor: transactions.at(-1)?.cursor ?? afterCursor, + has_more: start + transactions.length < result.document.committed.length, + }; + } +} diff --git a/loopx/control_plane/coordination/nokv_jsonl_helper.py b/loopx/control_plane/coordination/nokv_jsonl_helper.py new file mode 100644 index 0000000000..bbcf8094d9 --- /dev/null +++ b/loopx/control_plane/coordination/nokv_jsonl_helper.py @@ -0,0 +1,610 @@ +"""Narrow JSON-lines bridge from LoopX to the NoKV Python SDK. + +The bridge deliberately knows only byte storage. Authority transitions, +receipts, cursor ordering, and ambiguous-outcome reconciliation stay in the +TypeScript ``NoKVAuthorityStore``. The first JSON line configures one SDK +client; every later line invokes exactly one of ``store_identity``, +``read_blob``, or ``cas_publish_blob``. +""" + +from __future__ import annotations + +import base64 +import binascii +import json +import re +import sys +from collections.abc import Callable, Mapping +from typing import Any, TextIO + +_HEX_128 = re.compile(r"^[0-9a-f]{32}$") +_CLIENT_AVAILABILITY_ERRORS = (RuntimeError, OSError) +QUALIFIED_NOKV_SDK_VERSION = "0.11.0" +QUALIFIED_NOKV_API_VERSION = 1 + +_CONFIG_KEYS = frozenset( + { + "root_id", + "routing", + "object_store", + "max_attempts", + "connect_timeout_ms", + "read_timeout_ms", + "write_timeout_ms", + "handshake_timeout_ms", + "workbench_root", + } +) +_ETCD_ROUTING_KEYS = frozenset({"kind", "endpoints", "key_prefix", "lease_ttl_seconds"}) +_STATIC_ROUTING_KEYS = frozenset( + { + "kind", + "endpoint", + "logical_shard_id", + "object_namespace_id", + "placement_generation", + "owner_epoch", + } +) +_MEMORY_OBJECT_STORE_KEYS = frozenset({"kind"}) +_S3_OBJECT_STORE_KEYS = frozenset( + { + "kind", + "bucket", + "region", + "root", + "endpoint", + "access_key_id", + "secret_access_key", + "session_token", + "virtual_host_style", + "skip_signature", + } +) + + +class RequestError(ValueError): + """The JSON-lines caller violated the raw storage protocol.""" + + +class ProviderProtocolError(RuntimeError): + """The NoKV SDK returned a shape that violates its reviewed contract.""" + + +class ClientAdmissionUnavailable(RuntimeError): + """The configured SDK could not admit a live route or object provider.""" + + +def _request_id(value: object) -> str | None: + if isinstance(value, str) and value.strip() == value and value: + return value + return None + + +def _response( + request_id: str | None, status: str, **values: object +) -> dict[str, object]: + return {"request_id": request_id, "status": status, **values} + + +def _failure( + request_id: str | None, + status: str, + reason_code: str, + error: object, +) -> dict[str, object]: + return _response( + request_id, + status, + reason_code=reason_code, + reason=str(error) if str(error) else reason_code, + ) + + +def _opaque_failure( + request_id: str | None, + status: str, + reason_code: str, + reason: str, +) -> dict[str, object]: + return _response( + request_id, + status, + reason_code=reason_code, + reason=reason, + ) + + +def _mapping(value: object, name: str) -> Mapping[str, Any]: + if not isinstance(value, Mapping): + raise RequestError(f"{name} must be an object") + return value + + +def _require_exact_keys( + values: Mapping[str, Any], + allowed: frozenset[str], + name: str, +) -> None: + if any(key not in allowed for key in values): + raise RequestError(f"{name} contains unsupported fields") + + +def _required_string(values: Mapping[str, Any], name: str) -> str: + value = values.get(name) + if not isinstance(value, str) or not value or value.strip() != value: + raise RequestError(f"{name} must be a non-empty trimmed string") + return value + + +def _generation(value: object, name: str, *, nullable: bool = False) -> int | None: + if value is None and nullable: + return None + if not isinstance(value, int) or isinstance(value, bool) or value < 1: + suffix = " or null" if nullable else "" + raise RequestError(f"{name} must be a positive integer{suffix}") + return value + + +def _sdk_generation(value: object, name: str) -> int: + try: + generation = _generation(value, name) + except RequestError as error: + raise ProviderProtocolError(str(error)) from error + assert generation is not None + return generation + + +def _decode_bytes(value: object) -> bytes: + if not isinstance(value, str): + raise RequestError("bytes_base64 must be a string") + try: + return base64.b64decode(value, validate=True) + except (binascii.Error, ValueError) as error: + raise RequestError("bytes_base64 is not canonical base64") from error + + +def _identity(client: Any, workbench: str) -> str: + cursor: bytes | None = None + seen_cursors: set[bytes] = set() + identities: set[str] = set() + while True: + page = _mapping( + client.find_workspaces(cursor=cursor, limit=100), + "find_workspaces result", + ) + workspaces = page.get("workspaces") + if not isinstance(workspaces, list): + raise ProviderProtocolError("find_workspaces omitted workspaces") + for item in workspaces: + entry = _mapping(item, "workspace entry") + workspace = _mapping(entry.get("workspace"), "workspace summary") + if workspace.get("workbench") != workbench: + continue + incarnation = workspace.get("workspace_incarnation_id") + if not isinstance(incarnation, str) or not _HEX_128.fullmatch(incarnation): + raise ProviderProtocolError( + "workspace incarnation identity must be 32 lowercase hex" + ) + identities.add(incarnation) + next_cursor = page.get("next_cursor") + if next_cursor is None: + break + if not isinstance(next_cursor, bytes) or not next_cursor: + raise ProviderProtocolError("find_workspaces next_cursor is invalid") + if next_cursor in seen_cursors: + raise ProviderProtocolError("find_workspaces cursor did not advance") + seen_cursors.add(next_cursor) + cursor = next_cursor + if len(identities) != 1: + raise ProviderProtocolError( + "workbench did not resolve to one incarnation identity" + ) + return f"nokv:{workbench}:{next(iter(identities))}" + + +def _store_identity( + client: Any, + request_id: str, + values: Mapping[str, Any], +) -> dict[str, object]: + workbench = _required_string(values, "workbench") + try: + identity = _identity(client, workbench) + except (ProviderProtocolError, TypeError, ValueError) as error: + return _failure( + request_id, + "failed", + "provider_protocol_violation", + error, + ) + except _CLIENT_AVAILABILITY_ERRORS: + return _opaque_failure( + request_id, + "unavailable", + "nokv_identity_unavailable", + "NoKV identity lookup is unavailable", + ) + return _response(request_id, "available", store_identity=identity) + + +def _read_blob( + client: Any, + request_id: str, + values: Mapping[str, Any], +) -> dict[str, object]: + workbench = _required_string(values, "workbench") + path = _required_string(values, "path") + try: + result = _mapping(client.read(workbench, path), "read result") + except FileNotFoundError: + return _response(request_id, "missing") + except _CLIENT_AVAILABILITY_ERRORS: + return _opaque_failure( + request_id, + "unavailable", + "nokv_read_unavailable", + "NoKV blob read is unavailable", + ) + except (TypeError, ValueError) as error: + return _failure(request_id, "failed", "provider_protocol_violation", error) + try: + raw = result.get("bytes") + if not isinstance(raw, (bytes, bytearray, memoryview)): + raise ProviderProtocolError("read result omitted bytes") + metadata = _mapping(result.get("metadata"), "read metadata") + if metadata.get("workbench") != workbench: + raise ProviderProtocolError("read metadata workbench mismatch") + if metadata.get("path") != path: + raise ProviderProtocolError("read metadata path mismatch") + incarnation = metadata.get("workspace_incarnation_id") + if not isinstance(incarnation, str) or not _HEX_128.fullmatch(incarnation): + raise ProviderProtocolError( + "read metadata workspace incarnation identity is invalid" + ) + current_identity = _identity(client, workbench) + if current_identity != f"nokv:{workbench}:{incarnation}": + raise ProviderProtocolError("read metadata workspace incarnation mismatch") + generation = _sdk_generation(metadata.get("generation"), "read generation") + except (ProviderProtocolError, TypeError, ValueError) as error: + return _failure(request_id, "failed", "provider_protocol_violation", error) + except _CLIENT_AVAILABILITY_ERRORS: + return _opaque_failure( + request_id, + "unavailable", + "nokv_read_unavailable", + "NoKV blob read is unavailable", + ) + return _response( + request_id, + "loaded", + bytes_base64=base64.b64encode(bytes(raw)).decode("ascii"), + generation=generation, + ) + + +def _cas_publish_blob( + client: Any, + request_id: str, + values: Mapping[str, Any], +) -> dict[str, object]: + workbench = _required_string(values, "workbench") + path = _required_string(values, "path") + expected_generation = _generation( + values.get("expected_generation"), + "expected_generation", + nullable=True, + ) + payload = _decode_bytes(values.get("bytes_base64")) + operation_id = _required_string(values, "operation_id") + artifact_revision_id = _required_string(values, "artifact_revision_id") + if not _HEX_128.fullmatch(operation_id): + raise RequestError("operation_id must be 32 lowercase hex") + if not _HEX_128.fullmatch(artifact_revision_id): + raise RequestError("artifact_revision_id must be 32 lowercase hex") + try: + raw_result = client.publish_bytes( + workbench, + path, + payload, + content_type="application/json", + expected_generation=expected_generation, + operation_id=operation_id, + artifact_revision_id=artifact_revision_id, + ) + except FileExistsError: + return _response( + request_id, + "conflict", + current_generation=None, + ) + except (RuntimeError, OSError, TypeError, ValueError): + # RuntimeError covers both a rejected generation and a response lost + # after commit. Human error text cannot distinguish them, so only a + # later authority-envelope readback may settle the outcome. + return _opaque_failure( + request_id, + "ambiguous", + "nokv_publish_outcome_unknown", + "NoKV publish outcome is unknown", + ) + try: + result = _mapping(raw_result, "publish result") + except (TypeError, ValueError): + return _opaque_failure( + request_id, + "ambiguous", + "provider_protocol_violation", + "NoKV publish response violated the storage protocol", + ) + try: + generation = _sdk_generation(result.get("generation"), "publish generation") + expected_result_generation = (expected_generation or 0) + 1 + if result.get("workbench") != workbench: + raise ProviderProtocolError("publish result workbench mismatch") + if result.get("path") != path: + raise ProviderProtocolError("publish result path mismatch") + if result.get("operation_id") != operation_id: + raise ProviderProtocolError("publish result operation identity mismatch") + if result.get("artifact_revision_id") != artifact_revision_id: + raise ProviderProtocolError("publish result artifact revision mismatch") + if generation != expected_result_generation: + raise ProviderProtocolError("publish result generation mismatch") + except ProviderProtocolError: + return _opaque_failure( + request_id, + "ambiguous", + "provider_protocol_violation", + "NoKV publish response violated the storage protocol", + ) + return _response(request_id, "applied", generation=generation) + + +def handle_request(client: Any, value: object) -> dict[str, object]: + """Execute one raw storage request without interpreting stored bytes.""" + + request_id = _request_id( + value.get("request_id") if isinstance(value, Mapping) else None + ) + try: + values = _mapping(value, "request") + if request_id is None: + raise RequestError("request_id must be a non-empty trimmed string") + operation = _required_string(values, "operation") + handlers: dict[ + str, + Callable[[Any, str, Mapping[str, Any]], dict[str, object]], + ] = { + "store_identity": _store_identity, + "read_blob": _read_blob, + "cas_publish_blob": _cas_publish_blob, + } + handler = handlers.get(operation) + if handler is None: + raise RequestError(f"unknown operation {operation!r}") + return handler(client, request_id, values) + except RequestError as error: + return _failure(request_id, "failed", "invalid_request", error) + + +def serve(client: Any, incoming: TextIO, outgoing: TextIO) -> None: + """Serve raw storage requests until EOF.""" + + for line in incoming: + try: + request = json.loads(line) + except json.JSONDecodeError as error: + result = _failure(None, "failed", "invalid_json", error) + else: + result = handle_request(client, request) + outgoing.write(json.dumps(result, sort_keys=True, separators=(",", ":")) + "\n") + outgoing.flush() + + +def _string_list(values: Mapping[str, Any], name: str) -> list[str]: + value = values.get(name) + if not isinstance(value, list) or not value: + raise RequestError(f"{name} must be a non-empty array") + result: list[str] = [] + for item in value: + if not isinstance(item, str) or not item: + raise RequestError(f"{name} entries must be non-empty strings") + result.append(item) + return result + + +def build_client(config_value: object) -> Any: + """Construct one eagerly admitted NoKV client from the open handshake.""" + + config = _mapping(config_value, "config") + _require_exact_keys(config, _CONFIG_KEYS, "config") + routing_value = _mapping(config.get("routing"), "routing") + routing_kind = _required_string(routing_value, "kind") + if routing_kind == "etcd": + _require_exact_keys(routing_value, _ETCD_ROUTING_KEYS, "routing") + routing_arguments: tuple[object, ...] = ( + _string_list(routing_value, "endpoints"), + _required_string(routing_value, "key_prefix"), + _generation( + routing_value.get("lease_ttl_seconds", 10), + "lease_ttl_seconds", + ), + ) + elif routing_kind == "static": + _require_exact_keys(routing_value, _STATIC_ROUTING_KEYS, "routing") + routing_arguments = ( + _required_string(routing_value, "endpoint"), + _required_string(routing_value, "logical_shard_id"), + _required_string(routing_value, "object_namespace_id"), + _generation( + routing_value.get("placement_generation"), + "placement_generation", + ), + _generation(routing_value.get("owner_epoch"), "owner_epoch"), + ) + else: + raise RequestError(f"unsupported routing kind {routing_kind!r}") + + object_value = _mapping(config.get("object_store"), "object_store") + object_kind = _required_string(object_value, "kind") + if object_kind == "memory": + _require_exact_keys( + object_value, + _MEMORY_OBJECT_STORE_KEYS, + "object_store", + ) + object_arguments: dict[str, object] | None = None + elif object_kind == "s3": + _require_exact_keys(object_value, _S3_OBJECT_STORE_KEYS, "object_store") + object_arguments = { + "bucket": _required_string(object_value, "bucket"), + "region": object_value.get("region", "us-east-1"), + "root": object_value.get("root", "/"), + "endpoint": object_value.get("endpoint"), + "access_key_id": object_value.get("access_key_id"), + "secret_access_key": object_value.get("secret_access_key"), + "session_token": object_value.get("session_token"), + "virtual_host_style": object_value.get("virtual_host_style", False), + "skip_signature": object_value.get("skip_signature", False), + } + else: + raise RequestError(f"unsupported object store kind {object_kind!r}") + + root_id = _required_string(config, "root_id") + if not _HEX_128.fullmatch(root_id): + raise RequestError("root_id must be 32 lowercase hex") + max_attempts = _generation(config.get("max_attempts", 3), "max_attempts") + connect_timeout_ms = _generation( + config.get("connect_timeout_ms", 5_000), + "connect_timeout_ms", + ) + read_timeout_ms = _generation( + config.get("read_timeout_ms", 30_000), + "read_timeout_ms", + ) + write_timeout_ms = _generation( + config.get("write_timeout_ms", 30_000), + "write_timeout_ms", + ) + handshake_timeout_ms = _generation( + config.get("handshake_timeout_ms", 5_000), + "handshake_timeout_ms", + ) + workbench_root = config.get("workbench_root") + if workbench_root is not None and ( + not isinstance(workbench_root, str) or not workbench_root + ): + raise RequestError("workbench_root must be a non-empty string or null") + + try: + import nokv + except ImportError as error: # pragma: no cover - exercised by live packaging + raise RequestError("the NoKV Python SDK is not installed") from error + if ( + getattr(nokv, "__version__", None) != QUALIFIED_NOKV_SDK_VERSION + or getattr(nokv, "API_VERSION", None) != QUALIFIED_NOKV_API_VERSION + ): + raise RequestError( + "the NoKV Python SDK must be version " + f"{QUALIFIED_NOKV_SDK_VERSION} with API version " + f"{QUALIFIED_NOKV_API_VERSION}" + ) + try: + Client = nokv.Client + ObjectStoreConfig = nokv.ObjectStoreConfig + RoutingConfig = nokv.RoutingConfig + except AttributeError as error: + raise RequestError("the NoKV Python SDK surface is incomplete") from error + + try: + routing = ( + RoutingConfig.etcd(*routing_arguments) + if routing_kind == "etcd" + else RoutingConfig.static(*routing_arguments) + ) + except (TypeError, ValueError) as error: + raise RequestError("NoKV routing configuration is invalid") from error + try: + object_store = ( + ObjectStoreConfig.memory() + if object_arguments is None + else ObjectStoreConfig.s3(**object_arguments) + ) + except (TypeError, ValueError) as error: + raise RequestError("NoKV object-store configuration is invalid") from error + try: + return Client( + root_id=root_id, + routing=routing, + object_store=object_store, + max_attempts=max_attempts, + connect_timeout_ms=connect_timeout_ms, + read_timeout_ms=read_timeout_ms, + write_timeout_ms=write_timeout_ms, + handshake_timeout_ms=handshake_timeout_ms, + workbench_root=workbench_root, + ) + except (RuntimeError, OSError, ValueError) as error: + raise ClientAdmissionUnavailable from error + except TypeError as error: + raise RequestError("NoKV client configuration is invalid") from error + + +def main() -> int: + first = sys.stdin.readline() + try: + value = json.loads(first) + values = _mapping(value, "open request") + request_id = _request_id(values.get("request_id")) + if request_id is None or values.get("operation") != "open": + raise RequestError("first request must be an open handshake") + client = build_client(values.get("config")) + except json.JSONDecodeError as error: + result = _failure(None, "failed", "invalid_json", error) + except RequestError as error: + result = _failure( + _request_id(value.get("request_id")) + if isinstance(value, Mapping) + else None, + "failed", + "invalid_config", + error, + ) + except _CLIENT_AVAILABILITY_ERRORS: + result = _opaque_failure( + request_id, + "unavailable", + "nokv_open_unavailable", + "NoKV client admission is unavailable", + ) + except (TypeError, ValueError): + result = _opaque_failure( + request_id, + "failed", + "invalid_config", + "NoKV helper configuration is invalid", + ) + else: + sys.stdout.write( + json.dumps( + _response( + request_id, + "ready", + nokv_sdk_version=QUALIFIED_NOKV_SDK_VERSION, + nokv_api_version=QUALIFIED_NOKV_API_VERSION, + ), + sort_keys=True, + separators=(",", ":"), + ) + + "\n" + ) + sys.stdout.flush() + serve(client, sys.stdin, sys.stdout) + return 0 + sys.stdout.write(json.dumps(result, sort_keys=True, separators=(",", ":")) + "\n") + sys.stdout.flush() + return 2 + + +if __name__ == "__main__": # pragma: no cover - exercised by subprocess E2E + raise SystemExit(main()) diff --git a/loopx/control_plane/coordination/nokv_jsonl_transport.ts b/loopx/control_plane/coordination/nokv_jsonl_transport.ts new file mode 100644 index 0000000000..2986b4af31 --- /dev/null +++ b/loopx/control_plane/coordination/nokv_jsonl_transport.ts @@ -0,0 +1,358 @@ +import { randomUUID } from "node:crypto"; +import { + spawn, + type ChildProcessWithoutNullStreams, + type SpawnOptionsWithoutStdio, +} from "node:child_process"; +import { once } from "node:events"; + +import type { JsonObject } from "../effect_program.ts"; +import { + NoKVTransportProtocolError, + NoKVTransportUnavailableError, + type NoKVBlobCasRequest, + type NoKVBlobCasResult, + type NoKVBlobReadResult, + type NoKVBlobTransport, + type NoKVStoreIdentityResult, + type NoKVTransportFailure, +} from "./nokv_authority_store.ts"; +import { + isAuthorityJsonObject, +} from "./authority_store_codec.ts"; + +const DEFAULT_REQUEST_TIMEOUT_MS = 30_000; +const DEFAULT_MAX_RESPONSE_BYTES = 32 * 1024 * 1024; +const TRANSPORT_CONSTRUCTION_TOKEN = Symbol("NoKVJsonLinesTransport.open"); + +export type NoKVJsonLinesProcessFactory = ( + command: string, + args: readonly string[], + options: SpawnOptionsWithoutStdio, +) => ChildProcessWithoutNullStreams; + +export interface NoKVJsonLinesTransportOptions { + /** Explicit command plus arguments, for example `[python, helper.py]`. */ + argv: readonly string[]; + /** Sent only over stdin in the helper's open handshake. */ + config: JsonObject; + cwd?: string; + env?: NodeJS.ProcessEnv; + request_timeout_ms?: number; + max_response_bytes?: number; + process_factory?: NoKVJsonLinesProcessFactory; +} + +interface PendingResponse { + resolve(value: JsonObject): void; + reject(error: Error): void; + timer: NodeJS.Timeout; +} + +function positiveSafeInteger(value: unknown, name: string): number { + if (!Number.isSafeInteger(value) || (value as number) < 1) { + throw new NoKVTransportProtocolError(`${name} must be a positive safe integer`); + } + return value as number; +} + +function requiredString(value: unknown, name: string): string { + if (typeof value !== "string" || value.trim() !== value || value.length === 0) { + throw new NoKVTransportProtocolError(`${name} must be a non-empty trimmed string`); + } + return value; +} + +function responseFailure(value: JsonObject): NoKVTransportFailure { + const status = value.status; + if (status !== "unavailable" && status !== "failed") { + throw new NoKVTransportProtocolError("helper response status is invalid"); + } + return { + status, + reason_code: requiredString(value.reason_code, "helper reason code"), + reason: requiredString(value.reason, "helper reason"), + }; +} + +function canonicalBase64(value: unknown): Uint8Array { + if (typeof value !== "string") { + throw new NoKVTransportProtocolError("helper bytes_base64 must be a string"); + } + const bytes = Buffer.from(value, "base64"); + if (bytes.toString("base64") !== value) { + throw new NoKVTransportProtocolError("helper bytes_base64 is not canonical base64"); + } + return bytes; +} + +/** + * Reusable process connection to the NoKV Python SDK helper. + * + * Construction is intentionally asynchronous and requires an explicit argv; + * no runtime path selects NoKV merely by importing this module. + */ +export class NoKVJsonLinesTransport implements NoKVBlobTransport { + private readonly child: ChildProcessWithoutNullStreams; + private readonly requestTimeoutMs: number; + private readonly maxResponseBytes: number; + private readonly pending = new Map(); + private stdoutBuffer = Buffer.alloc(0); + private terminalError: Error | null = null; + private closing = false; + + constructor( + options: NoKVJsonLinesTransportOptions, + constructionToken: typeof TRANSPORT_CONSTRUCTION_TOKEN, + ) { + if (constructionToken !== TRANSPORT_CONSTRUCTION_TOKEN) { + throw new NoKVTransportProtocolError( + "NoKV JSON-lines transport must be created with open()", + ); + } + if (!Array.isArray(options.argv) || options.argv.length === 0) { + throw new NoKVTransportProtocolError("NoKV helper argv must not be empty"); + } + for (const value of options.argv) { + requiredString(value, "NoKV helper argv entry"); + } + this.requestTimeoutMs = options.request_timeout_ms ?? DEFAULT_REQUEST_TIMEOUT_MS; + this.maxResponseBytes = options.max_response_bytes ?? DEFAULT_MAX_RESPONSE_BYTES; + positiveSafeInteger(this.requestTimeoutMs, "NoKV helper request timeout"); + positiveSafeInteger(this.maxResponseBytes, "NoKV helper max response bytes"); + const factory = options.process_factory ?? ((command, args, spawnOptions) => + spawn(command, args, { ...spawnOptions, stdio: ["pipe", "pipe", "pipe"] })); + const [command, ...args] = options.argv; + try { + this.child = factory(command!, args, { + cwd: options.cwd, + env: options.env, + }); + } catch { + throw new NoKVTransportUnavailableError("NoKV helper failed to start"); + } + this.child.stdout.on("data", (chunk: Buffer | string) => { + this.onStdout(Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk)); + }); + // Drain stderr so a noisy SDK cannot block, but never promote arbitrary + // provider output into LoopX errors where endpoints or credentials could + // escape the provider boundary. + this.child.stderr.on("data", () => {}); + this.child.stdin.on("error", () => { + if (!this.closing) this.failUnavailable("NoKV helper stdin failed"); + }); + this.child.on("error", () => { + this.failUnavailable("NoKV helper failed to start"); + }); + this.child.on("exit", (code, signal) => { + if (this.closing && this.pending.size === 0) return; + this.failUnavailable( + `NoKV helper disconnected (code=${String(code)}, signal=${String(signal)})`, + false, + ); + }); + } + + static async open(options: NoKVJsonLinesTransportOptions): Promise { + const transport = new NoKVJsonLinesTransport( + options, + TRANSPORT_CONSTRUCTION_TOKEN, + ); + let response: JsonObject; + try { + response = await transport.exchange("open", { config: options.config }); + if (response.status === "unavailable") { + throw new NoKVTransportUnavailableError( + requiredString(response.reason, "helper open reason"), + ); + } + if (response.status !== "ready") { + throw new NoKVTransportProtocolError( + response.status === "failed" && typeof response.reason === "string" + ? response.reason + : "NoKV helper did not acknowledge its open handshake", + ); + } + return transport; + } catch (error) { + await transport.close(); + throw error; + } + } + + private fail(error: Error, terminate: boolean): void { + if (this.terminalError === null) this.terminalError = error; + for (const pending of this.pending.values()) { + clearTimeout(pending.timer); + pending.reject(this.terminalError); + } + this.pending.clear(); + if (terminate && this.child.exitCode === null && this.child.signalCode === null) { + this.child.kill(); + } + } + + private failUnavailable(message: string, terminate = true): void { + this.fail(new NoKVTransportUnavailableError(message), terminate); + } + + private failProtocol(message: string): void { + this.fail(new NoKVTransportProtocolError(message), true); + } + + private onStdout(chunk: Buffer): void { + if (this.terminalError !== null) return; + this.stdoutBuffer = Buffer.concat([this.stdoutBuffer, chunk]); + while (true) { + const newline = this.stdoutBuffer.indexOf(0x0a); + if (newline < 0) break; + if (newline > this.maxResponseBytes) { + this.failProtocol("NoKV helper response exceeded max_response_bytes"); + return; + } + const line = this.stdoutBuffer.subarray(0, newline); + this.stdoutBuffer = this.stdoutBuffer.subarray(newline + 1); + if (line.byteLength === 0) { + this.failProtocol("NoKV helper emitted an empty response line"); + return; + } + let value: unknown; + try { + value = JSON.parse(line.toString("utf8")); + } catch (error) { + this.failProtocol( + `NoKV helper emitted invalid JSON: ${ + error instanceof Error ? error.message : "invalid JSON" + }`, + ); + return; + } + if (!isAuthorityJsonObject(value)) { + this.failProtocol("NoKV helper response must be an object"); + return; + } + const requestId = value.request_id; + if (typeof requestId !== "string") { + this.failProtocol("NoKV helper response omitted request_id"); + return; + } + const pending = this.pending.get(requestId); + if (!pending) { + this.failProtocol("NoKV helper responded with an unknown request_id"); + return; + } + this.pending.delete(requestId); + clearTimeout(pending.timer); + pending.resolve(value); + } + if (this.stdoutBuffer.byteLength > this.maxResponseBytes) { + this.failProtocol("NoKV helper response exceeded max_response_bytes"); + } + } + + private async exchange(operation: string, values: JsonObject): Promise { + if (this.terminalError) throw this.terminalError; + if (this.closing) { + throw new NoKVTransportUnavailableError("NoKV helper transport is closed"); + } + const requestId = randomUUID(); + const response = new Promise((resolve, reject) => { + const timer = setTimeout(() => { + this.pending.delete(requestId); + const error = new NoKVTransportUnavailableError( + `NoKV helper request ${operation} timed out`, + ); + reject(error); + this.fail(error, true); + }, this.requestTimeoutMs); + this.pending.set(requestId, { resolve, reject, timer }); + }); + const line = JSON.stringify({ request_id: requestId, operation, ...values }) + "\n"; + try { + this.child.stdin.write(line); + } catch { + this.failUnavailable("NoKV helper request write failed"); + } + return await response; + } + + async storeIdentity(workbench: string): Promise { + const response = await this.exchange("store_identity", { workbench }); + if (response.status === "available") { + return { + status: "available", + store_identity: requiredString( + response.store_identity, + "helper store identity", + ), + }; + } + return responseFailure(response); + } + + async readBlob(workbench: string, path: string): Promise { + const response = await this.exchange("read_blob", { workbench, path }); + if (response.status === "missing") return { status: "missing" }; + if (response.status === "loaded") { + return { + status: "loaded", + bytes: canonicalBase64(response.bytes_base64), + generation: positiveSafeInteger(response.generation, "helper read generation"), + }; + } + return responseFailure(response); + } + + async casPublishBlob(request: NoKVBlobCasRequest): Promise { + const response = await this.exchange("cas_publish_blob", { + workbench: request.workbench, + path: request.path, + expected_generation: request.expected_generation, + bytes_base64: Buffer.from(request.bytes).toString("base64"), + operation_id: request.operation_id, + artifact_revision_id: request.artifact_revision_id, + }); + if (response.status === "applied") { + return { + status: "applied", + generation: positiveSafeInteger(response.generation, "helper publish generation"), + }; + } + if (response.status === "conflict") { + const current = response.current_generation; + return { + status: "conflict", + current_generation: current === null + ? null + : positiveSafeInteger(current, "helper conflict generation"), + }; + } + if (response.status === "ambiguous") { + return { + status: "ambiguous", + reason_code: requiredString(response.reason_code, "helper reason code"), + reason: requiredString(response.reason, "helper reason"), + }; + } + const failure = responseFailure(response); + return { + status: "failed", + reason_code: failure.reason_code, + reason: failure.reason, + }; + } + + async close(): Promise { + if (this.closing) return; + this.closing = true; + if (this.child.exitCode !== null || this.child.signalCode !== null) return; + this.child.stdin.end(); + await Promise.race([ + once(this.child, "exit"), + new Promise((resolve) => setTimeout(resolve, 1_000)), + ]); + if (this.child.exitCode === null && this.child.signalCode === null) { + this.child.kill(); + } + } +} diff --git a/tests/control_plane_ts/nokv_authority_store.test.ts b/tests/control_plane_ts/nokv_authority_store.test.ts new file mode 100644 index 0000000000..0c08c66c13 --- /dev/null +++ b/tests/control_plane_ts/nokv_authority_store.test.ts @@ -0,0 +1,333 @@ +import assert from "node:assert/strict"; +import test from "node:test"; + +import { + NoKVAuthorityStore, + type NoKVBlobCasRequest, + type NoKVBlobCasResult, + type NoKVBlobReadResult, + type NoKVBlobTransport, + type NoKVStoreIdentityResult, +} from "../../loopx/control_plane/coordination/nokv_authority_store.ts"; +import { + authorityStoreCommitFixture as commit, + registerAuthorityStoreConformance, +} from "./authority_store_conformance.ts"; + +type PublishFault = + | "ambiguous_before" + | "terminal_ambiguous_before" + | "ambiguous_after" + | "ambiguous_after_then_read_unavailable" + | "failed" + | null; + +interface FakeNoKVBackend { + identity: string; + blob: { bytes: Uint8Array; generation: number } | null; + identityUnavailable: boolean; + rotateIdentityBeforePublish: string | null; + readUnavailable: number; + publishFault: PublishFault; + casRequests: NoKVBlobCasRequest[]; + terminalPhysicalIds: Set; +} + +function fakeBackend(): FakeNoKVBackend { + return { + identity: `nokv:authority-workbench:${"a".repeat(32)}`, + blob: null, + identityUnavailable: false, + rotateIdentityBeforePublish: null, + readUnavailable: 0, + publishFault: null, + casRequests: [], + terminalPhysicalIds: new Set(), + }; +} + +class FakeNoKVTransport implements NoKVBlobTransport { + readonly backend: FakeNoKVBackend; + + constructor(backend: FakeNoKVBackend) { + this.backend = backend; + } + + async storeIdentity(_workbench: string): Promise { + if (this.backend.identityUnavailable) { + return { + status: "unavailable", + reason_code: "injected_identity_unavailable", + reason: "identity lookup unavailable", + }; + } + return { status: "available", store_identity: this.backend.identity }; + } + + async readBlob(_workbench: string, _path: string): Promise { + if (this.backend.readUnavailable > 0) { + this.backend.readUnavailable -= 1; + return { + status: "unavailable", + reason_code: "injected_read_unavailable", + reason: "blob read unavailable", + }; + } + return this.backend.blob + ? { + status: "loaded", + bytes: this.backend.blob.bytes.slice(), + generation: this.backend.blob.generation, + } + : { status: "missing" }; + } + + async casPublishBlob(request: NoKVBlobCasRequest): Promise { + this.backend.casRequests.push({ ...request, bytes: request.bytes.slice() }); + if (this.backend.rotateIdentityBeforePublish !== null) { + this.backend.identity = this.backend.rotateIdentityBeforePublish; + this.backend.rotateIdentityBeforePublish = null; + } + if ( + this.backend.terminalPhysicalIds.has(request.operation_id) || + this.backend.terminalPhysicalIds.has(request.artifact_revision_id) + ) { + return { + status: "ambiguous", + reason_code: "injected_terminal_identity_spent", + reason: "physical publication identity is terminal", + }; + } + const current = this.backend.blob?.generation ?? null; + if (current !== request.expected_generation) { + return { status: "conflict", current_generation: current }; + } + if (this.backend.publishFault === "failed") { + this.backend.publishFault = null; + return { + status: "failed", + reason_code: "injected_publish_rejected", + reason: "publish rejected before SDK call", + }; + } + if (this.backend.publishFault === "ambiguous_before") { + this.backend.publishFault = null; + return { + status: "ambiguous", + reason_code: "injected_lost_response", + reason: "publish outcome unknown", + }; + } + if (this.backend.publishFault === "terminal_ambiguous_before") { + this.backend.publishFault = null; + this.backend.terminalPhysicalIds.add(request.operation_id); + this.backend.terminalPhysicalIds.add(request.artifact_revision_id); + return { + status: "ambiguous", + reason_code: "injected_terminal_identity_spent", + reason: "physical publication identity failed terminally", + }; + } + const generation = (request.expected_generation ?? 0) + 1; + this.backend.blob = { bytes: request.bytes.slice(), generation }; + if ( + this.backend.publishFault === "ambiguous_after" || + this.backend.publishFault === "ambiguous_after_then_read_unavailable" + ) { + if (this.backend.publishFault === "ambiguous_after_then_read_unavailable") { + this.backend.readUnavailable += 1; + } + this.backend.publishFault = null; + return { + status: "ambiguous", + reason_code: "injected_lost_response", + reason: "publish response was lost", + }; + } + return { status: "applied", generation }; + } +} + +function store(backend: FakeNoKVBackend, tenantId = "tenant-a", goalId = "goal-a") { + return new NoKVAuthorityStore(new FakeNoKVTransport(backend), { + tenant_id: tenantId, + goal_id: goalId, + workbench: "authority-workbench", + }); +} + +registerAuthorityStoreConformance("NoKV single-envelope provider", async () => { + const backend = fakeBackend(); + return { store: store(backend), contender: store(backend) }; +}); + +test("NoKV provider uses a deterministic CLI-readable metadata path", () => { + const backend = fakeBackend(); + const first = store(backend); + const same = store(backend); + const otherTenant = store(backend, "tenant-b"); + + assert.equal(first.path, same.path); + assert.match(first.path, /^metadata\/loopx-authority\/[0-9a-f]{32}\.json$/); + assert.notEqual(first.path, otherTenant.path); +}); + +test("NoKV provider keeps proven missing distinct from identity and read unavailability", async () => { + const backend = fakeBackend(); + const provider = store(backend); + assert.deepEqual(await provider.loadAuthority(), { status: "missing" }); + + backend.readUnavailable = 1; + const readUnavailable = await provider.loadAuthority(); + assert.equal(readUnavailable.status, "unavailable"); + if (readUnavailable.status === "unavailable") { + assert.equal(readUnavailable.reason_code, "injected_read_unavailable"); + } + + backend.identityUnavailable = true; + const identityUnavailable = await provider.loadAuthority(); + assert.equal(identityUnavailable.status, "unavailable"); + assert.equal((await provider.storeIdentity()).status, "unavailable"); +}); + +test("NoKV provider reconciles a lost success from the embedded operation receipt", async () => { + const backend = fakeBackend(); + const provider = store(backend); + backend.publishFault = "ambiguous_after"; + + const applied = await provider.commitAuthority(commit(null, "operation-a", 1, 7)); + assert.equal(applied.status, "applied"); + const receipt = await provider.readReceipt("operation-a"); + assert.equal(receipt.status, "found"); + if (receipt.status === "found") assert.equal(receipt.receipts[0]?.lease_epoch, 7); +}); + +test("NoKV provider leaves an outcome ambiguous until readback becomes available", async () => { + const backend = fakeBackend(); + const provider = store(backend); + backend.publishFault = "ambiguous_after_then_read_unavailable"; + + const unknown = await provider.commitAuthority(commit(null, "operation-a", 1, 8)); + assert.equal(unknown.status, "ambiguous"); + const receipt = await provider.readReceipt("operation-a"); + assert.equal(receipt.status, "found"); + if (receipt.status === "found") assert.equal(receipt.receipts[0]?.lease_epoch, 8); +}); + +test("NoKV provider does not invent a receipt when an ambiguous publish did not land", async () => { + const backend = fakeBackend(); + const provider = store(backend); + backend.publishFault = "ambiguous_before"; + + const unknown = await provider.commitAuthority(commit(null, "operation-a", 1, 9)); + assert.equal(unknown.status, "ambiguous"); + assert.deepEqual(await provider.readReceipt("operation-a"), { status: "missing" }); + assert.deepEqual(await provider.loadAuthority(), { status: "missing" }); +}); + +test("NoKV provider retries one logical commit with fresh physical identities", async () => { + const backend = fakeBackend(); + const provider = store(backend); + const request = commit(null, "operation-a", 1, 1); + backend.publishFault = "terminal_ambiguous_before"; + + assert.equal((await provider.commitAuthority(request)).status, "ambiguous"); + assert.equal((await provider.commitAuthority(request)).status, "applied"); + assert.equal(backend.casRequests.length, 2); + assert.notEqual( + backend.casRequests[0]?.operation_id, + backend.casRequests[1]?.operation_id, + ); + assert.notEqual( + backend.casRequests[0]?.artifact_revision_id, + backend.casRequests[1]?.artifact_revision_id, + ); + assert.match(backend.casRequests[0]!.operation_id, /^[0-9a-f]{32}$/); + assert.match(backend.casRequests[0]!.artifact_revision_id, /^[0-9a-f]{32}$/); + assert.notEqual( + backend.casRequests[0]?.operation_id, + backend.casRequests[0]?.artifact_revision_id, + ); +}); + +test("NoKV provider fences restored bytes with a different workspace incarnation", async () => { + const backend = fakeBackend(); + const original = store(backend); + const applied = await original.commitAuthority(commit(null, "operation-a", 1, 1)); + assert.equal(applied.status, "applied"); + const callsBeforeRestore = backend.casRequests.length; + + backend.identity = `nokv:authority-workbench:${"b".repeat(32)}`; + const restored = store(backend); + const loaded = await restored.loadAuthority(); + assert.equal(loaded.status, "failed"); + if (loaded.status === "failed") assert.match(loaded.reason, /lineage mismatch/); + + const rejected = await restored.commitAuthority( + commit( + applied.status === "applied" ? applied.provider_revision : null, + "operation-b", + 2, + 2, + ), + ); + assert.equal(rejected.status, "failed"); + assert.equal(backend.casRequests.length, callsBeforeRestore); +}); + +test("NoKV provider does not report applied across a workbench-incarnation race", async () => { + const backend = fakeBackend(); + const provider = store(backend); + backend.rotateIdentityBeforePublish = `nokv:authority-workbench:${"b".repeat(32)}`; + + const result = await provider.commitAuthority(commit(null, "operation-a", 1, 1)); + + assert.equal(result.status, "failed"); + if (result.status === "failed") { + assert.equal(result.reason_code, "provider_protocol_violation"); + assert.match(result.reason, /lineage mismatch/); + } + assert.equal((await provider.loadAuthority()).status, "failed"); +}); + +test("NoKV provider fails closed when persisted generation and bytes diverge", async () => { + const backend = fakeBackend(); + const provider = store(backend); + assert.equal( + (await provider.commitAuthority(commit(null, "operation-a", 1, 1))).status, + "applied", + ); + backend.blob!.generation += 1; + + const loaded = await provider.loadAuthority(); + assert.equal(loaded.status, "failed"); + if (loaded.status === "failed") assert.match(loaded.reason, /storage generation/); +}); + +test("NoKV provider treats non-UTF8 persisted bytes as a protocol failure", async () => { + const backend = fakeBackend(); + backend.blob = { bytes: Uint8Array.of(0xff), generation: 1 }; + + const loaded = await store(backend).loadAuthority(); + assert.equal(loaded.status, "failed"); + if (loaded.status === "failed") { + assert.equal(loaded.reason_code, "provider_protocol_violation"); + } +}); + +test("NoKV provider enforces its candidate envelope capacity before CAS", async () => { + const backend = fakeBackend(); + const provider = new NoKVAuthorityStore(new FakeNoKVTransport(backend), { + tenant_id: "tenant-a", + goal_id: "goal-a", + workbench: "authority-workbench", + max_envelope_bytes: 64, + }); + + const result = await provider.commitAuthority(commit(null, "operation-a", 1, 1)); + assert.equal(result.status, "failed"); + if (result.status === "failed") { + assert.equal(result.reason_code, "authority_envelope_too_large"); + } + assert.equal(backend.casRequests.length, 0); +}); diff --git a/tests/control_plane_ts/nokv_jsonl_transport.test.ts b/tests/control_plane_ts/nokv_jsonl_transport.test.ts new file mode 100644 index 0000000000..54641a31b3 --- /dev/null +++ b/tests/control_plane_ts/nokv_jsonl_transport.test.ts @@ -0,0 +1,203 @@ +import assert from "node:assert/strict"; +import { fileURLToPath } from "node:url"; +import test from "node:test"; + +import { + NoKVAuthorityStore, + NoKVTransportProtocolError, + NoKVTransportUnavailableError, +} from "../../loopx/control_plane/coordination/nokv_authority_store.ts"; +import { NoKVJsonLinesTransport } from "../../loopx/control_plane/coordination/nokv_jsonl_transport.ts"; +import { registerAuthorityStoreConformance } from "./authority_store_conformance.ts"; + +const PYTHON = process.env.LOOPX_TEST_PYTHON ?? "python3"; +const FAULT_HELPER = fileURLToPath( + new URL("../fixtures/nokv_jsonl_fake_helper.py", import.meta.url), +); +const SDK_HELPER = fileURLToPath( + new URL("../../loopx/control_plane/coordination/nokv_jsonl_helper.py", import.meta.url), +); +const FAKE_SDK_ROOT = fileURLToPath( + new URL("../fixtures/nokv_fake_sdk", import.meta.url), +); + +async function openSdkHelper() { + return await NoKVJsonLinesTransport.open({ + argv: [PYTHON, SDK_HELPER], + config: { + root_id: "0".repeat(32), + routing: { + kind: "etcd", + endpoints: ["http://127.0.0.1:2379"], + key_prefix: "/nokv/control", + lease_ttl_seconds: 10, + }, + object_store: { kind: "memory" }, + }, + env: { + ...process.env, + PYTHONPATH: process.env.PYTHONPATH + ? `${FAKE_SDK_ROOT}:${process.env.PYTHONPATH}` + : FAKE_SDK_ROOT, + }, + request_timeout_ms: 2_000, + }); +} + +async function openFaultHelper(mode: string, maxResponseBytes?: number) { + return await NoKVJsonLinesTransport.open({ + argv: [PYTHON, FAULT_HELPER, mode], + config: {}, + request_timeout_ms: 2_000, + max_response_bytes: maxResponseBytes, + }); +} + +test("JSON-lines transport cannot bypass its open handshake", () => { + let processStarted = false; + const DirectTransport = NoKVJsonLinesTransport as unknown as new ( + options: { + argv: readonly string[]; + config: Record; + process_factory: () => never; + }, + constructionToken: symbol, + ) => NoKVJsonLinesTransport; + + assert.throws( + () => + new DirectTransport( + { + argv: ["injected-helper"], + config: {}, + process_factory: () => { + processStarted = true; + throw new Error("constructor reached process creation"); + }, + }, + Symbol("caller-token"), + ), + /must be created with open\(\)/, + ); + assert.equal(processStarted, false); +}); + +registerAuthorityStoreConformance("NoKV JSON-lines process", async (t) => { + const transport = await openSdkHelper(); + t.after(async () => await transport.close()); + return { + store: new NoKVAuthorityStore(transport, { + tenant_id: "tenant-a", + goal_id: "goal-a", + workbench: "authority-workbench", + }), + contender: new NoKVAuthorityStore(transport, { + tenant_id: "tenant-a", + goal_id: "goal-a", + workbench: "authority-workbench", + }), + }; +}); + +test("JSON-lines transport starts once and reuses the helper process", async (t) => { + const transport = await openSdkHelper(); + t.after(async () => await transport.close()); + + assert.deepEqual(await transport.storeIdentity("authority-workbench"), { + status: "available", + store_identity: `nokv:authority-workbench:${"a".repeat(32)}`, + }); + assert.deepEqual( + await transport.readBlob("authority-workbench", "metadata/head.json"), + { status: "missing" }, + ); + assert.deepEqual( + await transport.casPublishBlob({ + workbench: "authority-workbench", + path: "metadata/head.json", + expected_generation: null, + bytes: Buffer.from("payload", "utf8"), + operation_id: "a".repeat(32), + artifact_revision_id: "b".repeat(32), + }), + { status: "applied", generation: 1 }, + ); + const loaded = await transport.readBlob("authority-workbench", "metadata/head.json"); + assert.equal(loaded.status, "loaded"); + if (loaded.status === "loaded") { + assert.equal(Buffer.from(loaded.bytes).toString("utf8"), "payload"); + assert.equal(loaded.generation, 1); + } +}); + +test("JSON-lines helper disconnect is typed unavailable", async (t) => { + const transport = await openFaultHelper("disconnect"); + t.after(async () => await transport.close()); + + await assert.rejects( + transport.readBlob("authority-workbench", "metadata/head.json"), + (error: unknown) => { + assert.ok(error instanceof NoKVTransportUnavailableError); + assert.doesNotMatch(error.message, /provider-private-diagnostic/); + assert.match(error.message, /code=17/); + return true; + }, + ); +}); + +test("JSON-lines synchronous start failure is typed and sanitized", async () => { + await assert.rejects( + NoKVJsonLinesTransport.open({ + argv: ["injected-helper"], + config: {}, + process_factory: () => { + throw new Error("private helper path and provider detail"); + }, + }), + (error: unknown) => { + assert.ok(error instanceof NoKVTransportUnavailableError); + assert.equal(error.message, "NoKV helper failed to start"); + return true; + }, + ); +}); + +test("JSON-lines invalid response is a protocol failure", async (t) => { + const transport = await openFaultHelper("invalid"); + t.after(async () => await transport.close()); + + await assert.rejects( + transport.readBlob("authority-workbench", "metadata/head.json"), + NoKVTransportProtocolError, + ); +}); + +test("JSON-lines response limit fails closed as a protocol error", async (t) => { + const transport = await openFaultHelper("oversized", 256); + t.after(async () => await transport.close()); + + await assert.rejects( + transport.readBlob("authority-workbench", "metadata/head.json"), + (error: unknown) => { + assert.ok(error instanceof NoKVTransportProtocolError); + assert.match(error.message, /max_response_bytes/); + return true; + }, + ); +}); + +test("NoKV AuthorityStore preserves helper protocol failure as failed, not missing", async (t) => { + const transport = await openFaultHelper("invalid"); + t.after(async () => await transport.close()); + const store = new NoKVAuthorityStore(transport, { + tenant_id: "tenant-a", + goal_id: "goal-a", + workbench: "authority-workbench", + }); + + const loaded = await store.loadAuthority(); + assert.equal(loaded.status, "failed"); + if (loaded.status === "failed") { + assert.equal(loaded.reason_code, "provider_protocol_violation"); + } +}); diff --git a/tests/control_plane_ts/nokv_stage2a_qualification_harness.test.ts b/tests/control_plane_ts/nokv_stage2a_qualification_harness.test.ts new file mode 100644 index 0000000000..5b290895cd --- /dev/null +++ b/tests/control_plane_ts/nokv_stage2a_qualification_harness.test.ts @@ -0,0 +1,359 @@ +import assert from "node:assert/strict"; +import { spawnSync } from "node:child_process"; +import { existsSync, mkdirSync, mkdtempSync, rmSync, writeFileSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { isAbsolute, join } from "node:path"; +import { fileURLToPath } from "node:url"; +import test from "node:test"; + +import { + exerciseQualificationSequence, + parseQualificationArguments, + qualificationHelperArgv, + QUALIFICATION_SCOPE, + QUALIFIED_NOKV_API_VERSION, + QUALIFIED_NOKV_SDK_VERSION, + QualificationFailure, + type QualificationTransport, +} from "../../examples/nokv-authority-store/live-qualification.ts"; +import { + NoKVTransportProtocolError, + NoKVTransportUnavailableError, + type NoKVBlobCasRequest, + type NoKVBlobCasResult, + type NoKVBlobReadResult, + type NoKVStoreIdentityResult, +} from "../../loopx/control_plane/coordination/nokv_authority_store.ts"; +import { + NoKVJsonLinesTransport, +} from "../../loopx/control_plane/coordination/nokv_jsonl_transport.ts"; + +const REPOSITORY_HELPER = fileURLToPath( + new URL("../../loopx/control_plane/coordination/nokv_jsonl_helper.py", import.meta.url), +); +const PYTHON = process.env.LOOPX_TEST_PYTHON ?? "python3"; + +/** Minimal module that satisfies helper admission and records that it was imported. */ +const STAND_IN_SDK_SOURCE = `import os + +__version__ = "0.11.0" +API_VERSION = 1 + +_marker = os.environ.get("LOOPX_TEST_STAND_IN_MARKER") +if _marker: + with open(_marker, "w", encoding="utf-8") as handle: + handle.write("stand-in nokv imported") + + +class RoutingConfig: + @staticmethod + def etcd(*values): + return ("etcd", values) + + @staticmethod + def static(*values): + return ("static", values) + + +class ObjectStoreConfig: + @staticmethod + def memory(): + return ("memory",) + + @staticmethod + def s3(**values): + return ("s3", values) + + +class Client: + def __init__(self, **values): + self.values = values +`; + +function absolutePythonExecutable(): string | null { + const probe = spawnSync(PYTHON, ["-c", "import sys; print(sys.executable)"], { + encoding: "utf8", + }); + if (probe.status !== 0) return null; + const executable = probe.stdout.trim(); + return isAbsolute(executable) ? executable : null; +} + +interface Backend { + blob: { bytes: Uint8Array; generation: number } | null; + ignoreCas: boolean; + publishCalls: number; + pretendAppliedWithoutWriteOnCall: number | null; +} + +class FakeQualificationTransport implements QualificationTransport { + readonly backend: Backend; + closed = false; + + constructor(backend: Backend) { + this.backend = backend; + } + + async storeIdentity(workbench: string): Promise { + return { + status: "available", + store_identity: `nokv:${workbench}:${"a".repeat(32)}`, + }; + } + + async readBlob(_workbench: string, _path: string): Promise { + return this.backend.blob + ? { + status: "loaded", + bytes: this.backend.blob.bytes.slice(), + generation: this.backend.blob.generation, + } + : { status: "missing" }; + } + + async casPublishBlob(request: NoKVBlobCasRequest): Promise { + this.backend.publishCalls += 1; + const current = this.backend.blob?.generation ?? null; + if (!this.backend.ignoreCas && current !== request.expected_generation) { + return { status: "conflict", current_generation: current }; + } + const generation = (current ?? 0) + 1; + if (this.backend.pretendAppliedWithoutWriteOnCall === this.backend.publishCalls) { + return { status: "applied", generation }; + } + this.backend.blob = { bytes: request.bytes.slice(), generation }; + return { status: "applied", generation }; + } + + async close(): Promise { + this.closed = true; + } +} + +const BASE_OPTIONS = { + python_executable: "/usr/bin/python3", + client_config: { + root_id: "0".repeat(32), + routing: { kind: "etcd" }, + object_store: { kind: "memory" }, + }, + tenant_id: "qualification-tenant", + goal_id: "qualification-goal", + workbench: "authority-workbench", +} as const; + +test("Stage 2A qualification harness requires explicit write opt-in", () => { + assert.throws( + () => parseQualificationArguments([ + "--config-json", "/tmp/client.json", + "--python-executable", "/usr/bin/python3", + "--tenant-id", "qualification-tenant", + "--goal-id", "qualification-goal", + "--workbench", "authority-workbench", + ]), + (error: unknown) => { + assert.ok(error instanceof QualificationFailure); + assert.equal(error.reasonCode, "live_opt_in_required"); + return true; + }, + ); +}); + +test("Stage 2A qualification harness fixes the executable, isolation flag, and repository helper", () => { + const executable = "/opt/loopx-qualification/bin/python"; + + assert.deepEqual(qualificationHelperArgv(executable), [executable, "-I", REPOSITORY_HELPER]); + const parsed = parseQualificationArguments([ + "--execute-live", + "--config-json", "/tmp/client.json", + "--python-executable", executable, + "--tenant-id", "qualification-tenant", + "--goal-id", "qualification-goal", + "--workbench", "authority-workbench", + ]); + assert.equal(parsed.pythonExecutable, executable); +}); + +test("Stage 2A qualification harness rejects a relative Python executable", () => { + assert.throws( + () => qualificationHelperArgv("python3"), + (error: unknown) => { + assert.ok(error instanceof QualificationFailure); + assert.equal(error.reasonCode, "invalid_arguments"); + return true; + }, + ); +}); + +test("Stage 2A qualification harness names the exact NoKV SDK contract", () => { + assert.equal(QUALIFICATION_SCOPE, "stage_2a_single_node_store_conformance"); + assert.equal(QUALIFIED_NOKV_SDK_VERSION, "0.11.0"); + assert.equal(QUALIFIED_NOKV_API_VERSION, 1); +}); + +test("qualification proves create, ambiguous reconciliation, contention, and fresh readback", async () => { + const backend: Backend = { + blob: null, + ignoreCas: false, + publishCalls: 0, + pretendAppliedWithoutWriteOnCall: null, + }; + const opened: FakeQualificationTransport[] = []; + const report = await exerciseQualificationSequence(BASE_OPTIONS, async () => { + const transport = new FakeQualificationTransport(backend); + opened.push(transport); + return transport; + }); + + assert.equal(report.final_generation, 3); + assert.equal(report.final_cursor, "3"); + assert.deepEqual(report.checks.map((check) => check.id), [ + "existing_workbench_identity", + "fresh_authority_target", + "create_applied", + "create_generation_one", + "response_lost_success_reconciled", + "generation_cas_applied", + "generation_two_readback", + "competing_generation_cas_one_winner", + "competition_did_not_double_advance", + "independent_transport_readback", + "ambiguous_commit_receipt_retained", + "winner_receipt_retained", + "loser_receipt_absent", + ]); + assert.deepEqual(report.checks.map((check) => check.status), + Array(report.checks.length).fill("passed")); + assert.equal(opened.length, 3); + assert.ok(opened.every((transport) => transport.closed)); +}); + +test("qualification rejects a backend that does not enforce generation CAS", async () => { + const backend: Backend = { + blob: null, + ignoreCas: true, + publishCalls: 0, + pretendAppliedWithoutWriteOnCall: null, + }; + await assert.rejects( + exerciseQualificationSequence(BASE_OPTIONS, async () => + new FakeQualificationTransport(backend)), + (error: unknown) => { + assert.ok(error instanceof QualificationFailure); + assert.equal(error.reasonCode, "competition_not_fenced"); + return true; + }, + ); +}); + +test("qualification rejects an independent transport that cannot read the envelope", async () => { + const shared: Backend = { + blob: null, + ignoreCas: false, + publishCalls: 0, + pretendAppliedWithoutWriteOnCall: null, + }; + let opened = 0; + await assert.rejects( + exerciseQualificationSequence(BASE_OPTIONS, async () => { + opened += 1; + return new FakeQualificationTransport( + opened === 3 + ? { + blob: null, + ignoreCas: false, + publishCalls: 0, + pretendAppliedWithoutWriteOnCall: null, + } + : shared, + ); + }), + (error: unknown) => { + assert.ok(error instanceof QualificationFailure); + assert.equal(error.reasonCode, "independent_readback_failed"); + return true; + }, + ); +}); + +test("qualification requires durable readback after the injected response loss", async () => { + const backend: Backend = { + blob: null, + ignoreCas: false, + publishCalls: 0, + pretendAppliedWithoutWriteOnCall: 2, + }; + await assert.rejects( + exerciseQualificationSequence(BASE_OPTIONS, async () => + new FakeQualificationTransport(backend)), + (error: unknown) => { + assert.ok(error instanceof QualificationFailure); + assert.equal(error.reasonCode, "generation_cas_failed"); + return true; + }, + ); +}); + +test("qualification helper runs Python isolated so PYTHONPATH cannot substitute the nokv module", async (t) => { + const executable = absolutePythonExecutable(); + if (executable === null) { + t.skip("no Python interpreter is available for the helper"); + return; + } + const root = mkdtempSync(join(tmpdir(), "loopx-nokv-stand-in-")); + try { + const marker = join(root, "stand-in-imported"); + mkdirSync(join(root, "nokv")); + writeFileSync(join(root, "nokv", "__init__.py"), STAND_IN_SDK_SOURCE); + const config = { + root_id: "0".repeat(32), + routing: { + kind: "etcd", + endpoints: ["http://127.0.0.1:1"], + key_prefix: "/loopx-stand-in", + lease_ttl_seconds: 1, + }, + object_store: { kind: "memory" }, + }; + const env = { + ...process.env, + PYTHONPATH: root, + LOOPX_TEST_STAND_IN_MARKER: marker, + }; + + // Without isolation the PYTHONPATH stand-in is imported and admitted. This + // is the vector the qualification argv closes, so prove it exists first. + const unguarded = await NoKVJsonLinesTransport.open({ + argv: [executable, REPOSITORY_HELPER], + config, + env, + request_timeout_ms: 30_000, + }); + await unguarded.close(); + assert.ok(existsSync(marker), "stand-in module must be importable without -I"); + rmSync(marker); + + // The qualification argv must never consult the stand-in: the marker stays + // absent and the open handshake fails closed, either because the isolated + // interpreter has no nokv module or because the real SDK cannot reach the + // closed endpoint. + await assert.rejects( + NoKVJsonLinesTransport.open({ + argv: qualificationHelperArgv(executable), + config, + env, + request_timeout_ms: 30_000, + }), + (error: unknown) => + error instanceof NoKVTransportUnavailableError || + error instanceof NoKVTransportProtocolError, + ); + assert.equal( + existsSync(marker), + false, + "isolated interpreter must not import the PYTHONPATH stand-in", + ); + } finally { + rmSync(root, { recursive: true, force: true }); + } +}); diff --git a/tests/fixtures/nokv_fake_sdk/nokv/__init__.py b/tests/fixtures/nokv_fake_sdk/nokv/__init__.py new file mode 100644 index 0000000000..4ad7970ef2 --- /dev/null +++ b/tests/fixtures/nokv_fake_sdk/nokv/__init__.py @@ -0,0 +1,80 @@ +from __future__ import annotations + +from typing import Any + +__version__ = "0.11.0" +API_VERSION = 1 + + +class RoutingConfig: + @staticmethod + def etcd(endpoints: list[str], key_prefix: str, lease_ttl_seconds: int) -> object: + return ("etcd", endpoints, key_prefix, lease_ttl_seconds) + + @staticmethod + def static(*values: Any) -> object: + return ("static", values) + + +class ObjectStoreConfig: + @staticmethod + def memory() -> object: + return ("memory",) + + @staticmethod + def s3(**values: Any) -> object: + return ("s3", values) + + +class Client: + def __init__(self, **_values: Any) -> None: + self._bytes: bytes | None = None + self._generation: int | None = None + + def find_workspaces(self, **_values: Any) -> dict[str, Any]: + return { + "workspaces": [ + { + "workspace": { + "workbench": "authority-workbench", + "workspace_incarnation_id": "a" * 32, + } + } + ], + "next_cursor": None, + } + + def read(self, workbench: str, path: str) -> dict[str, Any]: + if self._bytes is None or self._generation is None: + raise FileNotFoundError("missing") + return { + "bytes": self._bytes, + "metadata": { + "workbench": workbench, + "path": path, + "workspace_incarnation_id": "a" * 32, + "generation": self._generation, + }, + } + + def publish_bytes( + self, + workbench: str, + path: str, + payload: bytes, + **values: Any, + ) -> dict[str, Any]: + expected = values["expected_generation"] + if expected is None and self._generation is not None: + raise FileExistsError("already exists") + if expected is not None and expected != self._generation: + raise RuntimeError("generation conflict") + self._generation = (self._generation or 0) + 1 + self._bytes = payload + return { + "operation_id": values["operation_id"], + "artifact_revision_id": values["artifact_revision_id"], + "workbench": workbench, + "path": path, + "generation": self._generation, + } diff --git a/tests/fixtures/nokv_jsonl_fake_helper.py b/tests/fixtures/nokv_jsonl_fake_helper.py new file mode 100644 index 0000000000..d5744293e5 --- /dev/null +++ b/tests/fixtures/nokv_jsonl_fake_helper.py @@ -0,0 +1,82 @@ +from __future__ import annotations + +import base64 +import json +import sys + + +mode = sys.argv[1] if len(sys.argv) > 1 else "normal" +blob: bytes | None = None +generation: int | None = None + + +def emit(value: object) -> None: + sys.stdout.write(json.dumps(value, separators=(",", ":")) + "\n") + sys.stdout.flush() + + +first = json.loads(sys.stdin.readline()) +emit({"request_id": first["request_id"], "status": "ready"}) + +for line in sys.stdin: + request = json.loads(line) + request_id = request["request_id"] + if mode == "disconnect": + sys.stderr.write("provider-private-diagnostic=must-not-escape\n") + sys.stderr.flush() + raise SystemExit(17) + if mode == "invalid": + emit({"request_id": request_id, "status": "loaded", "generation": True}) + continue + if mode == "oversized": + emit({"request_id": request_id, "status": "missing", "padding": "x" * 4_096}) + continue + operation = request["operation"] + if operation == "store_identity": + emit( + { + "request_id": request_id, + "status": "available", + "store_identity": f"nokv:{request['workbench']}:{'a' * 32}", + } + ) + elif operation == "read_blob": + if blob is None: + emit({"request_id": request_id, "status": "missing"}) + else: + emit( + { + "request_id": request_id, + "status": "loaded", + "bytes_base64": base64.b64encode(blob).decode("ascii"), + "generation": generation, + } + ) + elif operation == "cas_publish_blob": + if request["expected_generation"] != generation: + emit( + { + "request_id": request_id, + "status": "conflict", + "current_generation": generation, + } + ) + else: + blob = base64.b64decode(request["bytes_base64"], validate=True) + generation = (generation or 0) + 1 + emit( + { + "request_id": request_id, + "status": "applied", + "generation": generation, + } + ) + else: + emit( + { + "request_id": request_id, + "status": "failed", + "reason_code": "unknown_operation", + "reason": "unknown operation", + } + ) diff --git a/tests/test_nokv_jsonl_helper.py b/tests/test_nokv_jsonl_helper.py new file mode 100644 index 0000000000..b67bf1d467 --- /dev/null +++ b/tests/test_nokv_jsonl_helper.py @@ -0,0 +1,590 @@ +from __future__ import annotations + +import base64 +import io +import json +import sys +import types +from typing import Any + +import pytest + +from loopx.control_plane.coordination.nokv_jsonl_helper import ( + ClientAdmissionUnavailable, + RequestError, + build_client, + handle_request, + main, + serve, +) + + +class FakeClient: + def __init__(self) -> None: + self.find_pages: list[dict[str, Any]] = [] + self.read_result: dict[str, Any] | BaseException = FileNotFoundError("missing") + self.publish_result: dict[str, Any] | BaseException = publish_result() + self.publish_calls: list[tuple[tuple[Any, ...], dict[str, Any]]] = [] + + def find_workspaces(self, **kwargs: Any) -> dict[str, Any]: + assert kwargs["limit"] == 100 + return self.find_pages.pop(0) + + def read(self, *args: Any) -> dict[str, Any]: + if isinstance(self.read_result, BaseException): + raise self.read_result + return self.read_result + + def publish_bytes(self, *args: Any, **kwargs: Any) -> dict[str, Any]: + self.publish_calls.append((args, kwargs)) + if isinstance(self.publish_result, BaseException): + raise self.publish_result + return self.publish_result + + +def request(operation: str, **values: Any) -> dict[str, Any]: + return {"request_id": "request-a", "operation": operation, **values} + + +def identity_page() -> dict[str, Any]: + return { + "workspaces": [ + { + "workspace": { + "workbench": "authority-workbench", + "workspace_incarnation_id": "c" * 32, + } + } + ], + "next_cursor": None, + } + + +def read_result() -> dict[str, Any]: + return { + "bytes": b"canonical bytes", + "metadata": { + "workbench": "authority-workbench", + "path": "metadata/head.json", + "workspace_incarnation_id": "c" * 32, + "generation": 7, + }, + } + + +def publish_result(*, generation: int = 1) -> dict[str, Any]: + return { + "operation_id": "a" * 32, + "artifact_revision_id": "b" * 32, + "workbench": "authority-workbench", + "path": "metadata/head.json", + "generation": generation, + } + + +def test_store_identity_follows_all_pages_and_binds_the_workspace_incarnation() -> None: + client = FakeClient() + client.find_pages = [ + { + "workspaces": [{"workspace": {"workbench": "other"}}], + "next_cursor": b"page-two", + }, + { + "workspaces": [ + { + "workspace": { + "workbench": "authority-workbench", + "workspace_incarnation_id": "a" * 32, + } + } + ], + "next_cursor": None, + }, + ] + + result = handle_request( + client, + request("store_identity", workbench="authority-workbench"), + ) + assert result == { + "request_id": "request-a", + "status": "available", + "store_identity": f"nokv:authority-workbench:{'a' * 32}", + } + + +def test_store_identity_never_turns_an_outage_or_missing_workspace_into_identity() -> ( + None +): + unavailable = FakeClient() + unavailable.find_pages = [] + + def fail(**_kwargs: Any) -> dict[str, Any]: + raise RuntimeError("route unavailable") + + unavailable.find_workspaces = fail # type: ignore[method-assign] + result = handle_request( + unavailable, + request("store_identity", workbench="authority-workbench"), + ) + assert result["status"] == "unavailable" + assert result["reason_code"] == "nokv_identity_unavailable" + assert result["reason"] == "NoKV identity lookup is unavailable" + assert "route unavailable" not in result["reason"] + + absent = FakeClient() + absent.find_pages = [{"workspaces": [], "next_cursor": None}] + result = handle_request( + absent, + request("store_identity", workbench="authority-workbench"), + ) + assert result["status"] == "failed" + assert result["reason_code"] == "provider_protocol_violation" + + +def test_read_blob_preserves_missing_unavailable_and_generation() -> None: + client = FakeClient() + assert handle_request( + client, + request( + "read_blob", workbench="authority-workbench", path="metadata/head.json" + ), + ) == {"request_id": "request-a", "status": "missing"} + + client.read_result = RuntimeError("server unavailable") + unavailable = handle_request( + client, + request( + "read_blob", workbench="authority-workbench", path="metadata/head.json" + ), + ) + assert unavailable["status"] == "unavailable" + assert unavailable["reason_code"] == "nokv_read_unavailable" + assert unavailable["reason"] == "NoKV blob read is unavailable" + assert "server unavailable" not in unavailable["reason"] + + client.find_pages = [identity_page()] + client.read_result = read_result() + loaded = handle_request( + client, + request( + "read_blob", workbench="authority-workbench", path="metadata/head.json" + ), + ) + assert loaded == { + "request_id": "request-a", + "status": "loaded", + "bytes_base64": base64.b64encode(b"canonical bytes").decode("ascii"), + "generation": 7, + } + + +def test_cas_publish_blob_forwards_exact_generation_bytes_and_identities() -> None: + client = FakeClient() + client.publish_result = publish_result(generation=5) + payload = b'{"head":true}' + result = handle_request( + client, + request( + "cas_publish_blob", + workbench="authority-workbench", + path="metadata/head.json", + expected_generation=4, + bytes_base64=base64.b64encode(payload).decode("ascii"), + operation_id="a" * 32, + artifact_revision_id="b" * 32, + ), + ) + + assert result == { + "request_id": "request-a", + "status": "applied", + "generation": 5, + } + args, kwargs = client.publish_calls[0] + assert args == ("authority-workbench", "metadata/head.json", payload) + assert kwargs == { + "content_type": "application/json", + "expected_generation": 4, + "operation_id": "a" * 32, + "artifact_revision_id": "b" * 32, + } + + +@pytest.mark.parametrize( + ("field", "wrong_value"), + [ + ("workbench", "other-workbench"), + ("path", "metadata/other.json"), + ("workspace_incarnation_id", "d" * 32), + ], +) +def test_read_blob_rejects_sdk_metadata_bound_to_another_object_or_incarnation( + field: str, + wrong_value: object, +) -> None: + client = FakeClient() + client.find_pages = [identity_page()] + client.read_result = read_result() + client.read_result["metadata"][field] = wrong_value + + result = handle_request( + client, + request( + "read_blob", workbench="authority-workbench", path="metadata/head.json" + ), + ) + + assert result["status"] == "failed" + assert result["reason_code"] == "provider_protocol_violation" + + +@pytest.mark.parametrize( + ("field", "wrong_value"), + [ + ("workbench", "other-workbench"), + ("path", "metadata/other.json"), + ("operation_id", "d" * 32), + ("artifact_revision_id", "e" * 32), + ("generation", 6), + ], +) +def test_publish_never_reports_applied_for_an_sdk_result_bound_to_another_write( + field: str, + wrong_value: object, +) -> None: + client = FakeClient() + client.publish_result = publish_result(generation=5) + client.publish_result[field] = wrong_value + + result = handle_request( + client, + request( + "cas_publish_blob", + workbench="authority-workbench", + path="metadata/head.json", + expected_generation=4, + bytes_base64=base64.b64encode(b"{}").decode("ascii"), + operation_id="a" * 32, + artifact_revision_id="b" * 32, + ), + ) + + assert result["status"] == "ambiguous" + assert result["reason_code"] == "provider_protocol_violation" + + +def test_cas_publish_blob_maps_only_proven_collision_to_conflict() -> None: + client = FakeClient() + client.publish_result = FileExistsError("already exists") + conflict = handle_request( + client, + request( + "cas_publish_blob", + workbench="authority-workbench", + path="metadata/head.json", + expected_generation=None, + bytes_base64=base64.b64encode(b"{}").decode("ascii"), + operation_id="a" * 32, + artifact_revision_id="b" * 32, + ), + ) + assert conflict["status"] == "conflict" + assert conflict["current_generation"] is None + + client.publish_result = RuntimeError("generation conflict or lost response") + ambiguous = handle_request( + client, + request( + "cas_publish_blob", + workbench="authority-workbench", + path="metadata/head.json", + expected_generation=1, + bytes_base64=base64.b64encode(b"{}").decode("ascii"), + operation_id="a" * 32, + artifact_revision_id="b" * 32, + ), + ) + assert ambiguous["status"] == "ambiguous" + assert ambiguous["reason_code"] == "nokv_publish_outcome_unknown" + assert ambiguous["reason"] == "NoKV publish outcome is unknown" + assert "lost response" not in ambiguous["reason"] + + client.publish_result = ValueError("post-call conversion exposed an endpoint") + malformed = handle_request( + client, + request( + "cas_publish_blob", + workbench="authority-workbench", + path="metadata/head.json", + expected_generation=1, + bytes_base64=base64.b64encode(b"{}").decode("ascii"), + operation_id="a" * 32, + artifact_revision_id="b" * 32, + ), + ) + assert malformed["status"] == "ambiguous" + assert "endpoint" not in malformed["reason"] + + +def test_invalid_publish_request_fails_before_calling_the_sdk() -> None: + client = FakeClient() + invalid = handle_request( + client, + request( + "cas_publish_blob", + workbench="authority-workbench", + path="metadata/head.json", + expected_generation=True, + bytes_base64="not base64", + operation_id="short", + artifact_revision_id="b" * 32, + ), + ) + assert invalid["status"] == "failed" + assert invalid["reason_code"] == "invalid_request" + assert client.publish_calls == [] + + +def test_json_lines_server_emits_one_typed_response_per_request() -> None: + client = FakeClient() + incoming = io.StringIO( + json.dumps( + request( + "read_blob", + workbench="authority-workbench", + path="metadata/head.json", + ) + ) + + "\n" + + "not-json\n" + ) + outgoing = io.StringIO() + + serve(client, incoming, outgoing) + + rows = [json.loads(line) for line in outgoing.getvalue().splitlines()] + assert rows[0] == {"request_id": "request-a", "status": "missing"} + assert rows[1]["request_id"] is None + assert rows[1]["status"] == "failed" + assert rows[1]["reason_code"] == "invalid_json" + + +def test_static_route_requires_positive_generation_and_epoch_before_sdk_call( + monkeypatch: pytest.MonkeyPatch, +) -> None: + static_calls: list[tuple[Any, ...]] = [] + + class RoutingConfig: + @staticmethod + def static(*args: Any) -> object: + static_calls.append(args) + return object() + + module = types.SimpleNamespace( + __version__="0.11.0", + API_VERSION=1, + Client=lambda **_kwargs: object(), + ObjectStoreConfig=types.SimpleNamespace(memory=lambda: object()), + RoutingConfig=RoutingConfig, + ) + monkeypatch.setitem(sys.modules, "nokv", module) + base = { + "root_id": "a" * 32, + "routing": { + "kind": "static", + "endpoint": "127.0.0.1:7000", + "logical_shard_id": "b" * 32, + "object_namespace_id": "c" * 32, + "placement_generation": 1, + "owner_epoch": 1, + }, + "object_store": {"kind": "memory"}, + } + for field, value in [ + ("placement_generation", None), + ("placement_generation", True), + ("owner_epoch", 0), + ]: + invalid = json.loads(json.dumps(base)) + invalid["routing"][field] = value + with pytest.raises(RequestError): + build_client(invalid) + assert static_calls == [] + + +def test_client_constructor_value_error_is_typed_as_admission_unavailable( + monkeypatch: pytest.MonkeyPatch, +) -> None: + class RoutingConfig: + @staticmethod + def etcd(*_args: Any) -> object: + return object() + + def unavailable_client(**_kwargs: Any) -> object: + raise ValueError("provider endpoint and credential detail") + + module = types.SimpleNamespace( + __version__="0.11.0", + API_VERSION=1, + Client=unavailable_client, + ObjectStoreConfig=types.SimpleNamespace(memory=lambda: object()), + RoutingConfig=RoutingConfig, + ) + monkeypatch.setitem(sys.modules, "nokv", module) + + with pytest.raises(ClientAdmissionUnavailable) as raised: + build_client( + { + "root_id": "a" * 32, + "routing": { + "kind": "etcd", + "endpoints": ["http://unused.invalid"], + "key_prefix": "/nokv/control", + "lease_ttl_seconds": 10, + }, + "object_store": {"kind": "memory"}, + } + ) + assert "endpoint" not in str(raised.value) + + +@pytest.mark.parametrize("unknown_location", ["top", "routing", "object_store"]) +def test_unknown_config_keys_fail_before_any_sdk_object_is_constructed( + monkeypatch: pytest.MonkeyPatch, + unknown_location: str, +) -> None: + construction_calls: list[str] = [] + + class RoutingConfig: + @staticmethod + def etcd(*_args: Any) -> object: + construction_calls.append("routing") + return object() + + class ObjectStoreConfig: + @staticmethod + def memory() -> object: + construction_calls.append("object_store") + return object() + + def client(**_kwargs: Any) -> object: + construction_calls.append("client") + return object() + + module = types.SimpleNamespace( + __version__="0.11.0", + API_VERSION=1, + Client=client, + ObjectStoreConfig=ObjectStoreConfig, + RoutingConfig=RoutingConfig, + ) + monkeypatch.setitem(sys.modules, "nokv", module) + config = { + "root_id": "a" * 32, + "routing": { + "kind": "etcd", + "endpoints": ["http://unused.invalid"], + "key_prefix": "/nokv/control", + "lease_ttl_seconds": 10, + }, + "object_store": {"kind": "memory"}, + } + secret_marker = "must-not-appear" + if unknown_location == "top": + config["routing_typo"] = secret_marker + elif unknown_location == "routing": + config["routing"]["endpoints_typo"] = secret_marker + else: + config["object_store"]["secret_access_key_typo"] = secret_marker + + with pytest.raises(RequestError) as raised: + build_client(config) + + assert construction_calls == [] + assert secret_marker not in str(raised.value) + + +@pytest.mark.parametrize( + ("sdk_version", "api_version"), + [("incompatible-version", 1), ("0.11.0", 999)], +) +def test_sdk_version_or_api_mismatch_fails_before_provider_construction( + monkeypatch: pytest.MonkeyPatch, + sdk_version: str, + api_version: int, +) -> None: + construction_calls: list[str] = [] + module = types.SimpleNamespace( + __version__=sdk_version, + API_VERSION=api_version, + Client=lambda **_kwargs: construction_calls.append("client"), + ObjectStoreConfig=types.SimpleNamespace( + memory=lambda: construction_calls.append("object_store") + ), + RoutingConfig=types.SimpleNamespace( + etcd=lambda *_args: construction_calls.append("routing") + ), + ) + monkeypatch.setitem(sys.modules, "nokv", module) + + with pytest.raises(RequestError) as raised: + build_client( + { + "root_id": "a" * 32, + "routing": { + "kind": "etcd", + "endpoints": ["http://unused.invalid"], + "key_prefix": "/nokv/control", + "lease_ttl_seconds": 10, + }, + "object_store": {"kind": "memory"}, + } + ) + + assert construction_calls == [] + assert "incompatible-version" not in str(raised.value) + assert "999" not in str(raised.value) + + +def test_open_handshake_reports_the_qualified_sdk_contract( + monkeypatch: pytest.MonkeyPatch, +) -> None: + module = types.SimpleNamespace( + __version__="0.11.0", + API_VERSION=1, + Client=lambda **_kwargs: object(), + ObjectStoreConfig=types.SimpleNamespace(memory=lambda: object()), + RoutingConfig=types.SimpleNamespace(etcd=lambda *_args: object()), + ) + monkeypatch.setitem(sys.modules, "nokv", module) + incoming = io.StringIO( + json.dumps( + { + "request_id": "open-a", + "operation": "open", + "config": { + "root_id": "a" * 32, + "routing": { + "kind": "etcd", + "endpoints": ["http://unused.invalid"], + "key_prefix": "/nokv/control", + "lease_ttl_seconds": 10, + }, + "object_store": {"kind": "memory"}, + }, + } + ) + + "\n" + ) + outgoing = io.StringIO() + monkeypatch.setattr(sys, "stdin", incoming) + monkeypatch.setattr(sys, "stdout", outgoing) + + assert main() == 0 + assert json.loads(outgoing.getvalue().splitlines()[0]) == { + "request_id": "open-a", + "status": "ready", + "nokv_api_version": 1, + "nokv_sdk_version": "0.11.0", + } diff --git a/tsconfig.control-plane.json b/tsconfig.control-plane.json index 4bbfd66289..36a799cca7 100644 --- a/tsconfig.control-plane.json +++ b/tsconfig.control-plane.json @@ -21,6 +21,8 @@ "loopx/control_plane/coordination/authority_store.ts", "loopx/control_plane/coordination/authority_store_codec.ts", "loopx/control_plane/coordination/file_authority_store.ts", + "loopx/control_plane/coordination/nokv_authority_store.ts", + "loopx/control_plane/coordination/nokv_jsonl_transport.ts", "loopx/control_plane/coordination/postgresql_authority_store.ts", "loopx/control_plane/agents/delivery_workspace.ts", "loopx/control_plane/goals/vision_checkpoint.ts", @@ -47,12 +49,16 @@ "loopx/control_plane/work_items/task_lease_acquire.ts", "loopx/control_plane/work_items/task_lease_lifecycle.ts", "loopx/control_plane/work_items/task_lease_acquire_cli.ts", + "examples/nokv-authority-store/live-qualification.ts", "tests/control_plane_ts/effect_program.test.ts", "tests/control_plane_ts/effect_runtime_errors.test.ts", "tests/control_plane_ts/interaction_contract.test.ts", "tests/control_plane_ts/runtime_decode.test.ts", "tests/control_plane_ts/authority_store.test.ts", "tests/control_plane_ts/authority_store_conformance.ts", + "tests/control_plane_ts/nokv_authority_store.test.ts", + "tests/control_plane_ts/nokv_jsonl_transport.test.ts", + "tests/control_plane_ts/nokv_stage2a_qualification_harness.test.ts", "tests/control_plane_ts/postgresql_authority_store.integration.test.ts", "tests/control_plane_ts/delivery_continuity.test.ts", "tests/control_plane_ts/delivery_workspace.test.ts",