diff --git a/.github/workflows/python-integration.yml b/.github/workflows/python-integration.yml index 83c130b..1158e03 100644 --- a/.github/workflows/python-integration.yml +++ b/.github/workflows/python-integration.yml @@ -4,10 +4,12 @@ on: push: paths: - "python/**" + - "spec/SDK_DAEMON_TARGET_0_4_33.json" - ".github/workflows/python-integration.yml" pull_request: paths: - "python/**" + - "spec/SDK_DAEMON_TARGET_0_4_33.json" - ".github/workflows/python-integration.yml" workflow_dispatch: inputs: @@ -24,6 +26,9 @@ jobs: - name: Checkout cccc-sdk uses: actions/checkout@v4 + - name: Check SDK hardening contract + run: python3 scripts/check_sdk_hardening.py + - name: Checkout cccc (daemon) uses: actions/checkout@v4 with: @@ -51,3 +56,80 @@ jobs: trap cleanup EXIT cccc daemon start python python/examples/compat_check.py + + native-daemon-0433: + runs-on: ubuntu-latest + + steps: + - uses: actions/checkout@v4 + + - uses: actions/setup-python@v5 + with: + python-version: "3.12" + + - uses: dtolnay/rust-toolchain@1.88.0 + + - name: Install exact native daemon and Python SDK + run: | + cargo install cccc --version '=0.4.33' --locked + python -m pip install -e python + + - name: Verify stable Web Model completion replay + run: | + set -euo pipefail + export CCCC_HOME="$RUNNER_TEMP/cccc-sdk-python-0433" + cccc daemon start + cleanup() { + cccc daemon stop >/dev/null 2>&1 || true + } + trap cleanup EXIT + + python - <<'PY' + from cccc_sdk import CCCCClient + + client = CCCCClient() + created = client.group_create(title="python-sdk-native-0433") + group_id = str(created.get("group_id") or "") + if not group_id: + raise RuntimeError("group_create did not return group_id") + + try: + client.group_start(group_id=group_id) + client.actor_add( + group_id=group_id, + actor_id="web-ci", + runtime="web_model", + runner="headless", + ) + client.send( + group_id=group_id, + text="native runtime completion probe", + by="user", + to=["web-ci"], + ) + acquired = client.web_model_runtime_wait_next_turn( + group_id=group_id, + actor_id="web-ci", + ) + turn = acquired.get("turn") or {} + turn_id = str(turn.get("turn_id") or "") + event_ids = list(turn.get("event_ids") or []) + if not turn_id or not event_ids: + raise RuntimeError(f"wait_next_turn returned no work: {acquired!r}") + delivery_id = f"python-ci:{turn_id}" + options = { + "group_id": group_id, + "actor_id": "web-ci", + "turn_id": turn_id, + "delivery_id": delivery_id, + "event_ids": event_ids, + } + first = client.web_model_runtime_complete_turn(**options) + replay = client.web_model_runtime_complete_turn(**options) + if first.get("delivery_id") != delivery_id or replay.get("delivery_id") != delivery_id: + raise RuntimeError("completion did not preserve delivery_id") + if (first.get("read_event") or {}).get("id") != (replay.get("read_event") or {}).get("id"): + raise RuntimeError("completion replay created a second receipt") + finally: + client.group_delete(group_id=group_id) + PY diff --git a/.github/workflows/rust-ci.yml b/.github/workflows/rust-ci.yml index 1a7c85f..77bf69a 100644 --- a/.github/workflows/rust-ci.yml +++ b/.github/workflows/rust-ci.yml @@ -4,10 +4,12 @@ on: push: paths: - "rust/**" + - "spec/SDK_DAEMON_TARGET_0_4_33.json" - ".github/workflows/rust-ci.yml" pull_request: paths: - "rust/**" + - "spec/SDK_DAEMON_TARGET_0_4_33.json" - ".github/workflows/rust-ci.yml" jobs: @@ -15,6 +17,8 @@ jobs: runs-on: ubuntu-latest steps: - uses: actions/checkout@v4 + - name: Check SDK hardening contract + run: python3 scripts/check_sdk_hardening.py - uses: dtolnay/rust-toolchain@stable with: components: clippy,rustfmt @@ -23,14 +27,54 @@ jobs: run: cargo fmt --check - name: Lint working-directory: rust - run: cargo clippy --all-targets --all-features -- -D warnings + run: cargo clippy --locked --all-targets --all-features -- -D warnings - name: Test working-directory: rust - run: cargo test --all-targets + run: cargo test --locked --all-targets - name: Package working-directory: rust run: cargo package --locked + msrv: + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v4 + - uses: dtolnay/rust-toolchain@master + with: + toolchain: "1.74.0" + - name: Check declared MSRV + working-directory: rust + run: cargo check --locked + + windows: + runs-on: windows-latest + steps: + - uses: actions/checkout@v4 + - uses: dtolnay/rust-toolchain@stable + - name: Test cross-platform cursor replacement + working-directory: rust + run: cargo test --locked --all-targets + + native-daemon-0433: + runs-on: ubuntu-latest + env: + CCCC_RUN_LIVE_RELIABILITY: "1" + steps: + - uses: actions/checkout@v4 + - uses: dtolnay/rust-toolchain@1.88.0 + - name: Install exact native daemon + run: cargo install cccc --version '=0.4.33' --locked + - name: Run live reliability contract + run: | + set -euo pipefail + export CCCC_HOME="$RUNNER_TEMP/cccc-sdk-rust-0433" + cccc daemon start + cleanup() { + cccc daemon stop >/dev/null 2>&1 || true + } + trap cleanup EXIT + cargo test --locked --manifest-path rust/Cargo.toml --test live_reliability -- --nocapture + integration: runs-on: ubuntu-latest needs: test diff --git a/.github/workflows/spec-drift.yml b/.github/workflows/spec-drift.yml index 7e3f1f5..3e7ac96 100644 --- a/.github/workflows/spec-drift.yml +++ b/.github/workflows/spec-drift.yml @@ -5,11 +5,13 @@ on: paths: - "spec/**" - "scripts/check_specs_against_cccc.sh" + - "scripts/check_sdk_hardening.py" - ".github/workflows/spec-drift.yml" pull_request: paths: - "spec/**" - "scripts/check_specs_against_cccc.sh" + - "scripts/check_sdk_hardening.py" - ".github/workflows/spec-drift.yml" schedule: - cron: "17 3 * * *" @@ -36,3 +38,6 @@ jobs: - name: Compare mirrored standards run: bash scripts/check_specs_against_cccc.sh cccc-core + + - name: Check SDK hardening contract + run: python3 scripts/check_sdk_hardening.py diff --git a/.github/workflows/ts-ci.yml b/.github/workflows/ts-ci.yml index dd1058c..357c42b 100644 --- a/.github/workflows/ts-ci.yml +++ b/.github/workflows/ts-ci.yml @@ -4,10 +4,12 @@ on: push: paths: - "ts/**" + - "spec/SDK_DAEMON_TARGET_0_4_33.json" - ".github/workflows/ts-ci.yml" pull_request: paths: - "ts/**" + - "spec/SDK_DAEMON_TARGET_0_4_33.json" - ".github/workflows/ts-ci.yml" jobs: @@ -17,6 +19,9 @@ jobs: steps: - uses: actions/checkout@v4 + - name: Check SDK hardening contract + run: python3 scripts/check_sdk_hardening.py + - uses: actions/setup-node@v4 with: node-version: "20" @@ -171,3 +176,98 @@ jobs: } console.log('integration smoke passed', { groupId, eventId: event.id }); JS + + native-daemon-0433: + runs-on: ubuntu-latest + needs: test + + steps: + - uses: actions/checkout@v4 + + - uses: actions/setup-node@v4 + with: + node-version: "20" + cache: npm + cache-dependency-path: ts/package-lock.json + + - uses: dtolnay/rust-toolchain@1.88.0 + + - name: Build TypeScript SDK + working-directory: ts + run: | + npm ci + npm run build + + - name: Install exact native daemon + run: cargo install cccc --version '=0.4.33' --locked + + - name: Verify stream gating and stable completion replay + run: | + set -euo pipefail + export CCCC_HOME="$RUNNER_TEMP/cccc-sdk-ts-0433" + cccc daemon start + cleanup() { + cccc daemon stop >/dev/null 2>&1 || true + } + trap cleanup EXIT + + node --input-type=module <<'JS' + import { CCCCClient, IncompatibleDaemonError } from './ts/dist/index.js'; + + const client = await CCCCClient.create({ ccccHome: process.env.CCCC_HOME }); + const created = await client.groupCreate({ title: 'ts-sdk-native-0433' }); + const groupId = String(created.group_id || ''); + if (!groupId) throw new Error('groupCreate did not return group_id'); + + try { + await client.groupStart(groupId); + let streamProbeRejected = false; + try { + await client.assertCompatible({ requireOps: ['events_stream'] }); + } catch (error) { + if (error instanceof IncompatibleDaemonError) streamProbeRejected = true; + else throw error; + } + if (!streamProbeRejected) { + throw new Error('native 0.4.33 unexpectedly accepted events_stream probe'); + } + + await client.actorAdd({ + groupId, + actorId: 'web-ci', + runtime: 'web_model', + runner: 'headless', + }); + await client.send({ + groupId, + text: 'native runtime completion probe', + by: 'user', + to: ['web-ci'], + }); + const acquired = await client.webModelRuntimeWaitNextTurn({ + groupId, + actorId: 'web-ci', + }); + const turn = acquired.turn; + if (!turn || typeof turn !== 'object') { + throw new Error(`waitNextTurn returned no work: ${JSON.stringify(acquired)}`); + } + const turnId = String(turn.turn_id || ''); + const eventIds = Array.isArray(turn.event_ids) ? turn.event_ids.map(String) : []; + if (!turnId || eventIds.length === 0) { + throw new Error(`invalid native turn: ${JSON.stringify(turn)}`); + } + const deliveryId = `ts-ci:${turnId}`; + const options = { groupId, actorId: 'web-ci', turnId, deliveryId, eventIds }; + const first = await client.webModelRuntimeCompleteTurn(options); + const replay = await client.webModelRuntimeCompleteTurn(options); + if (first.delivery_id !== deliveryId || replay.delivery_id !== deliveryId) { + throw new Error('completion did not preserve delivery_id'); + } + if (!first.read_event || !replay.read_event || first.read_event.id !== replay.read_event.id) { + throw new Error('completion replay created a second receipt'); + } + } finally { + await client.groupDelete(groupId); + } + JS diff --git a/CHANGELOG.md b/CHANGELOG.md index ba4e85b..6a5c9d3 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -17,6 +17,10 @@ CCCC line and exposes the IPC surface available on that line. after an exchange has begun. - Scheduled CI drift detection for all three mirrored CCCC standards, automatic Python integration coverage, and a live Rust-SDK/current-daemon smoke job. +- Python and TypeScript Web Model wait/complete helpers with a required stable + `delivery_id`, plus the exact native 0.4.33 replay fixture. +- Rust 0.0.2 identity-bound reliable messaging, durable inbox checkpoints, and + explicit reconciliation for ambiguous send/reply/read outcomes. ### Changed @@ -52,6 +56,10 @@ CCCC line and exposes the IPC surface available on that line. - TypeScript's `INVALID_REQUEST` constant now matches the daemon's `invalid_request` code and includes current request-size and Remote Access administrator-token errors. +- TypeScript connection, response, and stream-handshake phases now share one + deadline; handshake cancellation closes the socket, buffered stream data is + byte-capped before decoding, and transport cleanup removes only SDK-owned + listeners. - The daemon IPC mirror now matches current CCCC core, including Remote Access administrator-token state and enforcement fields. @@ -63,6 +71,18 @@ CCCC line and exposes the IPC surface available on that line. maps and side-effect-free compatibility probing. TypeScript also compiles exported option fixtures so documented contract values cannot silently drift out of the published declaration surface. +- Added exact native 0.4.33 completion/reliability jobs, Rust 1.74 MSRV and + Windows cursor-replacement jobs, and a static SDK hardening contract gate. + +## Rust crate [0.0.2] — Unreleased + +### Added + +- Least-privilege `IdentityBoundClient` adapters for idempotent send, reply, + inbox polling, and mark-read operations. +- Nullable native cursor decoding, authoritative remote-ahead read-state + reconciliation, ledger-order fallback for notifications, and atomic + same-directory `FileCursorStore` replacement. ## Rust crate [0.0.1] — 2026-08-03 diff --git a/README.ja.md b/README.ja.md index fbdac51..21c7860 100644 --- a/README.ja.md +++ b/README.ja.md @@ -29,6 +29,7 @@ SDK と CCCC Web が同じ `CCCC_HOME` を参照していれば、書き込み 主な用途: - リアルタイム更新が必要な Web/IDE プラグイン(`events_stream`) - Working Group を監視して自動応答する bot/service +- identity-bound write と永続 inbox cursor を必要とする高信頼 Rust worker - group / actors / shared context / capability ポリシー / Group Space をプログラムから管理する社内ツール - `tracked_send`、Context Ops v3 task/agent state、capability discovery、ローカル memory API を使う workflow 連携 @@ -90,7 +91,7 @@ python python/examples/auto_ack_attention.py --group g_xxx --actor user ```toml [dependencies] -cccc-sdk = "0.0.1" +cccc-sdk = "0.0.2" ``` Rust クライアントは `CCCC_HOME` の Unix Socket/TCP daemon を自動検出し、 @@ -102,7 +103,7 @@ Rust クライアントは `CCCC_HOME` の Unix Socket/TCP daemon を自動検 ## バージョニングと互換性 SDK リリースは daemon のバージョン文字列ではなく contract に追従します: -- Python と TypeScript は現在の SDK リリースラインに追従し、Rust crate は `0.0.1` から開始します。 +- Python と TypeScript は現在の SDK リリースラインに追従し、Rust crate は当面 `0.0.x` 系列です。 - 実行時互換性は `assert_compatible(...)` で必要な capability/op を指定して確認します。 互換性は “契約/能力” で保証し、バージョン文字列の厳密一致には依存しません: diff --git a/README.md b/README.md index 333fb20..db7cb8b 100644 --- a/README.md +++ b/README.md @@ -30,6 +30,7 @@ If SDK clients and CCCC Web use the same `CCCC_HOME`, all writes are shared imme Typical use cases: - Reactive UI / IDE plugins that need real-time updates (`events_stream`) - Bots/services that watch groups and respond automatically +- Reliable Rust workers that need identity-bound writes and durable inbox checkpoints - Internal tools that create/manage groups, actors, shared context, capability policy, and Group Space programmatically - Workflow integrations that use `tracked_send`, Context Ops v3 task/agent state updates, capability discovery, and first-class local memory @@ -91,7 +92,7 @@ python python/examples/auto_ack_attention.py --group g_xxx --actor user ```toml [dependencies] -cccc-sdk = "0.0.1" +cccc-sdk = "0.0.2" ``` ```rust @@ -115,7 +116,7 @@ fn main() -> Result<(), Box> { SDK releases follow daemon contracts, not strict daemon version strings: - Python and TypeScript package versions track the current SDK release line; the - Rust crate starts at `0.0.1` while its public API settles. + Rust crate is on the `0.0.x` line while its public API settles. - Use `assert_compatible(...)` with required capabilities/ops for runtime gating. Compatibility is enforced by **contracts**, not by strict version string matching: @@ -132,6 +133,8 @@ This repo keeps a mirror under `spec/`: ```bash ./scripts/sync_specs_from_cccc.sh ../cccc +# Reproduce a tagged mirror when auditing an older compatibility fixture: +./scripts/sync_specs_from_cccc.sh ../cccc v0.4.33 ``` The sync command intentionally replaces only the three mirrored standards. diff --git a/README.zh-CN.md b/README.zh-CN.md index 20ecb7b..f0f85a7 100644 --- a/README.zh-CN.md +++ b/README.zh-CN.md @@ -29,6 +29,7 @@ CCCC SDK 是一套用于 CCCC 平台的**客户端 SDK**。 典型场景: - 需要实时更新的 Web/IDE 插件(`events_stream`) - 监听工作组并自动响应的 bot/service +- 需要身份绑定写入与持久 inbox 游标的可靠 Rust worker - 以编程方式创建/管理 group、actors、共享 context、capability 策略与 Group Space 的内部工具 - 使用 `tracked_send`、Context Ops v3 任务/agent state、capability discovery、本地 memory API 的工作流集成 @@ -90,7 +91,7 @@ python python/examples/auto_ack_attention.py --group g_xxx --actor user ```toml [dependencies] -cccc-sdk = "0.0.1" +cccc-sdk = "0.0.2" ``` Rust 客户端会自动发现 `CCCC_HOME` 下的 Unix Socket/TCP daemon,并提供通用 @@ -101,7 +102,7 @@ Rust 客户端会自动发现 `CCCC_HOME` 下的 Unix Socket/TCP daemon,并提 ## 版本策略与兼容性 SDK 发布跟随 daemon 合约,而不是硬匹配 daemon 版本号: -- Python 和 TypeScript 包版本跟随当前 SDK 发布线;Rust crate 从 `0.0.1` 起步。 +- Python 和 TypeScript 包版本跟随当前 SDK 发布线;Rust crate 暂处于 `0.0.x` 版本线。 - 运行时兼容请用 `assert_compatible(...)` 指定所需 capability/op。 我们保证兼容性的手段是“契约/能力”,而不是字符串版本号硬匹配: diff --git a/RELEASING.md b/RELEASING.md index dcdb4cf..37daed6 100644 --- a/RELEASING.md +++ b/RELEASING.md @@ -9,7 +9,7 @@ This repo is a monorepo with three deliverables: - SDK version tracks the supported CCCC line: next release `0.4.34`. - RC sequence is SDK-owned (`0.4.34rcN` for Python, `0.4.34-rc.N` for npm). -- The Rust crate begins at `0.0.1` while its public API settles. +- Rust 0.0.1 is published; the next source release is 0.0.2 while its public API settles. - Compatibility is enforced by contracts/capabilities/op-probing, not by matching RC numbers. ## 0) Sync specs (recommended) @@ -17,6 +17,9 @@ This repo is a monorepo with three deliverables: ```bash ./scripts/sync_specs_from_cccc.sh ../cccc ./scripts/check_specs_against_cccc.sh ../cccc +# Reproducible audit against a committed core revision: +./scripts/check_specs_against_cccc.sh ../cccc +python3 scripts/check_sdk_hardening.py ``` ## 1) Python release (PyPI/TestPyPI) @@ -106,11 +109,15 @@ npm publish --access public ```bash cd rust cargo fmt --check -cargo clippy --all-targets --all-features -- -D warnings -cargo test --all-targets +cargo clippy --locked --all-targets --all-features -- -D warnings +cargo test --locked --all-targets +cargo +1.74.0 check --locked cargo package --locked ``` +The Rust CI matrix also runs the full suite on Windows and the opt-in reliable +messaging test against the exact native CCCC 0.4.33 daemon. + ### Publish ```bash diff --git a/python/README.md b/python/README.md index e35afd6..3c2d007 100644 --- a/python/README.md +++ b/python/README.md @@ -222,8 +222,28 @@ recovered = c.web_model_runtime_recover_turn( actor_id="web-model", event_ids=["e_xxx"], ) + +# Complete a runtime-owned turn. Reuse the exact delivery_id if the caller +# must retry after an unknown outcome. +turn = c.web_model_runtime_wait_next_turn( + group_id="g_xxx", + actor_id="web-model", +) +payload = turn["turn"] +delivery_id = f"worker:{payload['turn_id']}" +c.web_model_runtime_complete_turn( + group_id="g_xxx", + actor_id="web-model", + turn_id=payload["turn_id"], + delivery_id=delivery_id, + event_ids=payload["event_ids"], +) ``` +`delivery_id` is required and is the completion replay key. A retry for the +same acquired turn must use the same value; generating a new value can create a +second completion receipt. + `term_resize()` sends the standard `term_resize` operation. For current Rust daemon builds that still expose `terminal_resize`, the SDK falls back only after receiving a structured `unknown_op`; transport failures are never diff --git a/python/src/cccc_sdk/client_0430_ops.py b/python/src/cccc_sdk/client_0430_ops.py index 7b69ae1..1cab2a7 100644 --- a/python/src/cccc_sdk/client_0430_ops.py +++ b/python/src/cccc_sdk/client_0430_ops.py @@ -3,11 +3,13 @@ from .client_0430_admin_ops import CCCC0430AdminOpsMixin from .client_0430_assistant_ops import CCCC0430AssistantOpsMixin from .client_0430_memory_ops import CCCC0430MemoryOpsMixin +from .client_0430_runtime_ops import CCCC0430RuntimeOpsMixin class CCCC0430OpsMixin( CCCC0430AdminOpsMixin, CCCC0430AssistantOpsMixin, CCCC0430MemoryOpsMixin, + CCCC0430RuntimeOpsMixin, ): pass diff --git a/python/src/cccc_sdk/client_0430_runtime_ops.py b/python/src/cccc_sdk/client_0430_runtime_ops.py new file mode 100644 index 0000000..6c67ee3 --- /dev/null +++ b/python/src/cccc_sdk/client_0430_runtime_ops.py @@ -0,0 +1,63 @@ +from __future__ import annotations + +from typing import Any, Dict, List, Optional + +from .client_0430_shared import _compact + + +class CCCC0430RuntimeOpsMixin: + """Web Model runtime operations retained from the CCCC 0.4.33 contract.""" + + def web_model_runtime_wait_next_turn( + self, + *, + group_id: str, + actor_id: str, + by: Optional[str] = None, + limit: int = 20, + kind_filter: str = "all", + ) -> Dict[str, Any]: + aid = str(actor_id) + return self.call( + "web_model_runtime_wait_next_turn", + { + "group_id": str(group_id), + "actor_id": aid, + "by": str(by) if by is not None else aid, + "limit": min(max(int(limit), 1), 20), + "kind_filter": str(kind_filter), + }, + ) + + def web_model_runtime_complete_turn( + self, + *, + group_id: str, + actor_id: str, + turn_id: str, + delivery_id: str, + event_ids: Optional[List[str]] = None, + latest_event_id: str = "", + status: str = "done", + summary: str = "", + by: Optional[str] = None, + ) -> Dict[str, Any]: + aid = str(actor_id) + return self.call( + "web_model_runtime_complete_turn", + _compact( + { + "group_id": str(group_id), + "actor_id": aid, + "by": str(by) if by is not None else aid, + "turn_id": str(turn_id), + "delivery_id": str(delivery_id), + "event_ids": [str(event_id) for event_id in event_ids] + if event_ids is not None + else None, + "latest_event_id": latest_event_id or None, + "status": str(status), + "summary": summary or None, + } + ), + ) diff --git a/python/tests/test_client_0430_contract.py b/python/tests/test_client_0430_contract.py index 1ce97e1..4b1b7c8 100644 --- a/python/tests/test_client_0430_contract.py +++ b/python/tests/test_client_0430_contract.py @@ -1,6 +1,8 @@ from __future__ import annotations +import json import unittest +from pathlib import Path from unittest.mock import patch from cccc_sdk.client import CCCCClient @@ -8,6 +10,13 @@ from cccc_sdk.transport import DaemonEndpoint +TARGET_FIXTURE = json.loads( + (Path(__file__).resolve().parents[2] / "spec" / "SDK_DAEMON_TARGET_0_4_33.json").read_text( + encoding="utf-8" + ) +) + + class TestClient0433Contract(unittest.TestCase): def _client(self) -> CCCCClient: return CCCCClient(endpoint=DaemonEndpoint(transport="tcp", host="127.0.0.1", port=9000)) @@ -348,6 +357,32 @@ def fake_call_daemon(*, endpoint, request, timeout_s): # type: ignore[no-untype self.assertEqual(captured[1]["args"]["text"], "Include omissions") self.assertEqual(captured[1]["args"]["source_text"], "Current meeting notes") + def test_web_model_completion_requires_and_reuses_delivery_id(self) -> None: + captured: list[dict] = [] + contract = TARGET_FIXTURE["operations"]["web_model_runtime_complete_turn"] + args = contract["request"]["args"] + + def fake_call_daemon(*, endpoint, request, timeout_s): # type: ignore[no-untyped-def] + captured.append(request) + return contract["completion_response"] + + with patch("cccc_sdk.client.call_daemon", side_effect=fake_call_daemon): + client = self._client() + for _ in range(2): + client.web_model_runtime_complete_turn( + group_id=args["group_id"], + actor_id=args["actor_id"], + turn_id=args["turn_id"], + delivery_id=args["delivery_id"], + event_ids=args["event_ids"], + status=args["status"], + ) + + self.assertEqual(captured[0], captured[1]) + self.assertEqual(captured[0]["op"], "web_model_runtime_complete_turn") + self.assertTrue(set(contract["required_args"]).issubset(captured[0]["args"])) + self.assertEqual(captured[0]["args"]["delivery_id"], args["delivery_id"]) + if __name__ == "__main__": unittest.main() diff --git a/rust/Cargo.lock b/rust/Cargo.lock index 73bf675..32111ff 100644 --- a/rust/Cargo.lock +++ b/rust/Cargo.lock @@ -2,27 +2,86 @@ # It is not intended for manual editing. version = 3 +[[package]] +name = "bitflags" +version = "2.13.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b588b76d00fde79687d7646a9b5bdf3cc0f655e0bbd080335a95d7e96f3587da" + [[package]] name = "cccc-sdk" -version = "0.0.1" +version = "0.0.2" dependencies = [ "serde", "serde_json", + "tempfile", "thiserror", ] +[[package]] +name = "cfg-if" +version = "1.0.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9330f8b2ff13f34540b44e946ef35111825727b38d33286ef986142615121801" + +[[package]] +name = "errno" +version = "0.3.14" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "39cab71617ae0d63f51a36d69f866391735b51691dbda63cf6f96d042b63efeb" +dependencies = [ + "libc", + "windows-sys 0.61.2", +] + +[[package]] +name = "fastrand" +version = "2.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "da7c62ceae207dd37ea5b845da6a0696c799f85e97da1ab5b7910be3c1c80223" + +[[package]] +name = "getrandom" +version = "0.3.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "899def5c37c4fd7b2664648c28120ecec138e4d395b459e5ca34f9cce2dd77fd" +dependencies = [ + "cfg-if", + "libc", + "r-efi", + "wasip2", +] + [[package]] name = "itoa" version = "1.0.18" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8f42a60cbdf9a97f5d2305f08a87dc4e09308d1276d28c869c684d7777685682" +[[package]] +name = "libc" +version = "0.2.189" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3eaf3ede3fee6db1a4c2ee091bf8a8b4dccdc6d17f656fb07896ee72867612f2" + +[[package]] +name = "linux-raw-sys" +version = "0.12.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "32a66949e030da00e8c7d4434b251670a91556f4144941d37452769c25d58a53" + [[package]] name = "memchr" version = "2.8.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "cf8baf1c55e62ffcace7a9f06f4bd9cd3f0c4beb022d3b367256b91b87513d98" +[[package]] +name = "once_cell" +version = "1.21.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9f7c3e4beb33f85d45ae3e3a1792185706c8e16d043238c593331cc7cd313b50" + [[package]] name = "proc-macro2" version = "1.0.107" @@ -41,6 +100,25 @@ dependencies = [ "proc-macro2", ] +[[package]] +name = "r-efi" +version = "5.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "69cdb34c158ceb288df11e18b4bd39de994f6657d83847bdffdbd7f346754b0f" + +[[package]] +name = "rustix" +version = "1.1.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b6fe4565b9518b83ef4f91bb47ce29620ca828bd32cb7e408f0062e9930ba190" +dependencies = [ + "bitflags", + "errno", + "libc", + "linux-raw-sys", + "windows-sys 0.61.2", +] + [[package]] name = "serde" version = "1.0.229" @@ -95,6 +173,19 @@ dependencies = [ "unicode-ident", ] +[[package]] +name = "tempfile" +version = "3.20.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e8a64e3985349f2441a1a9ef0b853f869006c3855f2cda6862a94d26ebb9d6a1" +dependencies = [ + "fastrand", + "getrandom", + "once_cell", + "rustix", + "windows-sys 0.59.0", +] + [[package]] name = "thiserror" version = "2.0.19" @@ -121,6 +212,109 @@ version = "1.0.24" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e6e4313cd5fcd3dad5cafa179702e2b244f760991f45397d14d4ebf38247da75" +[[package]] +name = "wasip2" +version = "1.0.4+wasi-0.2.12" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b67efb37e106e55ce722a510d6b5f9c17f083e5fc79afc2badeb12cc313d9487" +dependencies = [ + "wit-bindgen", +] + +[[package]] +name = "windows-link" +version = "0.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f0805222e57f7521d6a62e36fa9163bc891acd422f971defe97d64e70d0a4fe5" + +[[package]] +name = "windows-sys" +version = "0.59.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1e38bc4d79ed67fd075bcc251a1c39b32a1776bbe92e5bef1f0bf1f8c531853b" +dependencies = [ + "windows-targets", +] + +[[package]] +name = "windows-sys" +version = "0.61.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ae137229bcbd6cdf0f7b80a31df61766145077ddf49416a728b02cb3921ff3fc" +dependencies = [ + "windows-link", +] + +[[package]] +name = "windows-targets" +version = "0.52.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9b724f72796e036ab90c1021d4780d4d3d648aca59e491e6b98e725b84e99973" +dependencies = [ + "windows_aarch64_gnullvm", + "windows_aarch64_msvc", + "windows_i686_gnu", + "windows_i686_gnullvm", + "windows_i686_msvc", + "windows_x86_64_gnu", + "windows_x86_64_gnullvm", + "windows_x86_64_msvc", +] + +[[package]] +name = "windows_aarch64_gnullvm" +version = "0.52.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "32a4622180e7a0ec044bb555404c800bc9fd9ec262ec147edd5989ccd0c02cd3" + +[[package]] +name = "windows_aarch64_msvc" +version = "0.52.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "09ec2a7bb152e2252b53fa7803150007879548bc709c039df7627cabbd05d469" + +[[package]] +name = "windows_i686_gnu" +version = "0.52.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8e9b5ad5ab802e97eb8e295ac6720e509ee4c243f69d781394014ebfe8bbfa0b" + +[[package]] +name = "windows_i686_gnullvm" +version = "0.52.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0eee52d38c090b3caa76c563b86c3a4bd71ef1a819287c19d586d7334ae8ed66" + +[[package]] +name = "windows_i686_msvc" +version = "0.52.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "240948bc05c5e7c6dabba28bf89d89ffce3e303022809e73deaefe4f6ec56c66" + +[[package]] +name = "windows_x86_64_gnu" +version = "0.52.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "147a5c80aabfbf0c7d901cb5895d1de30ef2907eb21fbbab29ca94c5b08b1a78" + +[[package]] +name = "windows_x86_64_gnullvm" +version = "0.52.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "24d5b23dc417412679681396f2b49f3de8c1473deb516bd34410872eff51ed0d" + +[[package]] +name = "windows_x86_64_msvc" +version = "0.52.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "589f6da84c646204747d1270a2a5661ea66ed1cced2631d546fdfb155959f9ec" + +[[package]] +name = "wit-bindgen" +version = "0.57.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1ebf944e87a7c253233ad6766e082e3cd714b5d03812acc24c318f549614536e" + [[package]] name = "zmij" version = "1.0.23" diff --git a/rust/Cargo.toml b/rust/Cargo.toml index 3fbcbd6..18e2efe 100644 --- a/rust/Cargo.toml +++ b/rust/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "cccc-sdk" -version = "0.0.1" +version = "0.0.2" edition = "2021" rust-version = "1.74" description = "Official Rust client SDK for CCCC daemon IPC v1" @@ -16,4 +16,5 @@ exclude = ["tests/**"] [dependencies] serde = { version = "1", features = ["derive"] } serde_json = "1" +tempfile = "=3.20.0" thiserror = "2" diff --git a/rust/README.md b/rust/README.md index 12c31b4..0563ad0 100644 --- a/rust/README.md +++ b/rust/README.md @@ -6,7 +6,7 @@ Official blocking Rust client for CCCC Daemon IPC v1. ```toml [dependencies] -cccc-sdk = "0.0.1" +cccc-sdk = "0.0.2" ``` ## Quick start @@ -76,10 +76,25 @@ a connection-establishment failure. Once request exchange begins, failures are reported as `Error::OutcomeUnknown` and are never replayed automatically. Clients created with `new(endpoint)` keep that explicit endpoint. +## Reliable messaging + +Rust 0.0.2 adds an `IdentityBoundClient` that exposes only idempotent chat and +inbox operations. A caller supplies a `WorkloadIdentityHook`; the adapter binds +the verified principal and evidence to every request instead of treating the +wire-level `by` field as authentication. + +`FileCursorStore` persists the last fully processed event through a unique +same-directory temporary file and atomic replacement. `PersistentInbox` +reconciles a saved cursor before polling, accepts native `event_id: null` for a +fresh actor, maps native `duplicate: true` to `replayed`, and verifies a +remote-ahead cursor through `message_read_status` or immutable ledger order. +The adapter intentionally has no generic `call`, shutdown, configuration, or +credential methods. + `assert_compatible` probes requested operation names and rejects an advertised capability whose actual operation returns `unknown_op`. Streaming upgrade operations such as `events_stream` and `term_attach` are not -exposed as iterators in 0.0.1. `assert_compatible` deliberately skips unsafe +exposed as iterators in 0.0.2. `assert_compatible` deliberately skips unsafe duplex probes; a reusable stream API will be added only with stable ownership, close, and backpressure semantics. diff --git a/rust/src/error.rs b/rust/src/error.rs index 5ebd812..ec41d13 100644 --- a/rust/src/error.rs +++ b/rust/src/error.rs @@ -60,6 +60,9 @@ pub enum Error { #[error("incompatible CCCC daemon: {0}")] Incompatible(String), + + #[error("write reconciliation requires operator action: {0}")] + ReconciliationRequired(String), } pub type Result = std::result::Result; diff --git a/rust/src/identity.rs b/rust/src/identity.rs new file mode 100644 index 0000000..ab5a0f1 --- /dev/null +++ b/rust/src/identity.rs @@ -0,0 +1,99 @@ +use serde_json::{Map, Value}; + +use crate::{Error, Result}; + +/// Principal established by a workload identity provider. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct AuthenticatedPrincipal { + pub subject: String, + pub issuer: String, + pub evidence_id: Option, +} + +/// Verifiable carrier produced for one request by an identity hook. +#[derive(Clone, Debug, PartialEq)] +pub struct WorkloadIdentityEvidence { + pub carrier_key: String, + pub carrier: Value, +} + +impl WorkloadIdentityEvidence { + pub fn new(carrier_key: impl Into, carrier: Value) -> Result { + let evidence = Self { + carrier_key: carrier_key.into(), + carrier, + }; + if evidence.carrier_key.trim().is_empty() + || evidence.carrier_key == "by" + || evidence.carrier.is_null() + { + return Err(Error::Incompatible( + "workload identity evidence needs a non-by carrier key and non-null value".into(), + )); + } + Ok(evidence) + } +} + +impl AuthenticatedPrincipal { + pub fn new(subject: impl Into, issuer: impl Into) -> Result { + let principal = Self { + subject: subject.into(), + issuer: issuer.into(), + evidence_id: None, + }; + principal.validate()?; + Ok(principal) + } + + fn validate(&self) -> Result<()> { + if self.subject.trim().is_empty() || self.issuer.trim().is_empty() { + return Err(Error::Incompatible( + "authenticated principal subject and issuer must be non-empty".into(), + )); + } + Ok(()) + } +} + +/// Hook for an external workload identity implementation. +/// +/// `evidence` returns a signature, token, nonce, or other carrier for `args`. +/// The receiving daemon or gateway must verify that carrier; the SDK never +/// treats a caller-provided `by` value as authentication. +pub trait WorkloadIdentityHook { + fn principal(&self) -> Result; + + /// Sign or otherwise bind the operation and canonical args to evidence. + fn evidence( + &self, + operation: &str, + args: &Map, + ) -> Result; +} + +pub(crate) fn bind_identity( + hook: &H, + operation: &str, + args: &mut Map, +) -> Result { + let principal = hook.principal()?; + principal.validate()?; + if let Some(claimed) = args.get("by") { + if claimed.as_str() != Some(&principal.subject) { + return Err(Error::Incompatible( + "request by does not match authenticated principal".into(), + )); + } + } + args.insert("by".into(), Value::String(principal.subject.clone())); + let evidence = hook.evidence(operation, args)?; + if args.contains_key(&evidence.carrier_key) { + return Err(Error::Incompatible(format!( + "workload identity carrier would overwrite request field {}", + evidence.carrier_key + ))); + } + args.insert(evidence.carrier_key, evidence.carrier); + Ok(principal) +} diff --git a/rust/src/lib.rs b/rust/src/lib.rs index cda6a19..ff47a12 100644 --- a/rust/src/lib.rs +++ b/rust/src/lib.rs @@ -7,11 +7,14 @@ mod client; mod endpoint; mod error; +mod identity; mod protocol; +mod reliable; pub use client::{CCCCClient, CompatibilityRequirements}; pub use endpoint::{discover_endpoint, DaemonEndpoint}; pub use error::{DaemonError, Error, Result}; +pub use identity::{AuthenticatedPrincipal, WorkloadIdentityEvidence, WorkloadIdentityHook}; pub use protocol::{ DaemonRequest, DaemonResponse, PingResult, TerminalHistoryOptions, TerminalHistoryResult, TerminalResizeResult, TerminalSinceHistory, TerminalSinceOptions, TerminalSinceResult, @@ -19,3 +22,7 @@ pub use protocol::{ WebModelDeliveryPreference, WebModelDeliveryPreferencesResult, WebModelRecoveredTurn, WebModelRecoveredTurnDelivery, WebModelRuntimeRecoverTurnResult, }; +pub use reliable::{ + CursorStore, Event, FileCursorStore, IdentityBoundClient, InboxCursor, InboxPage, + MarkReadReconciliation, MarkReadResult, MessageWriteResult, PersistentInbox, +}; diff --git a/rust/src/reliable.rs b/rust/src/reliable.rs new file mode 100644 index 0000000..51564ca --- /dev/null +++ b/rust/src/reliable.rs @@ -0,0 +1,898 @@ +use std::collections::BTreeMap; +use std::fs; +use std::io::{self, Write}; +use std::path::{Path, PathBuf}; + +use serde::{Deserialize, Deserializer, Serialize}; +use serde_json::{Map, Value}; + +use crate::identity::{bind_identity, AuthenticatedPrincipal, WorkloadIdentityHook}; +use crate::{CCCCClient, Error, Result}; + +#[derive(Clone, Debug, Deserialize, Serialize)] +pub struct Event { + #[serde(alias = "event_id")] + pub id: String, + pub ts: String, + pub kind: String, + #[serde(default)] + pub by: String, + #[serde(default)] + pub data: Value, + #[serde(flatten)] + pub extra: Map, +} + +#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] +pub struct InboxCursor { + #[serde(default, deserialize_with = "deserialize_nullable_string")] + pub event_id: String, + #[serde(default, deserialize_with = "deserialize_nullable_string")] + pub ts: String, +} + +#[derive(Clone, Debug, Deserialize)] +pub struct InboxPage { + #[serde(default)] + pub messages: Vec, + pub cursor: InboxCursor, +} + +#[derive(Clone, Debug, Deserialize)] +pub struct MarkReadResult { + pub cursor: InboxCursor, + #[serde(default)] + pub event: Option, + #[serde(default)] + pub replayed: bool, +} + +#[derive(Clone, Debug, Deserialize)] +pub struct MessageWriteResult { + pub event: Event, + #[serde(default, alias = "duplicate")] + pub replayed: bool, +} + +#[derive(Debug, Deserialize)] +struct MessageReadStatus { + #[serde(default)] + read_status: BTreeMap, +} + +#[derive(Debug, Deserialize)] +struct LedgerWindow { + #[serde(default)] + events: Vec, + #[serde(default)] + has_more_after: bool, +} + +/// Result of querying the daemon cursor before reconciling `inbox_mark_read`. +#[derive(Clone, Debug)] +pub struct MarkReadReconciliation { + pub cursor: InboxCursor, + pub wrote: bool, + pub result: Option, +} + +/// Durable storage for the last fully processed inbox event. +pub trait CursorStore { + fn load(&self) -> Result>; + fn save(&self, cursor: &InboxCursor) -> Result<()>; +} + +/// JSON cursor file written by a same-directory temporary file + atomic replace. +#[derive(Clone, Debug)] +pub struct FileCursorStore { + path: PathBuf, +} + +impl FileCursorStore { + pub fn new(path: impl Into) -> Self { + Self { path: path.into() } + } + + pub fn path(&self) -> &Path { + &self.path + } +} + +impl CursorStore for FileCursorStore { + fn load(&self) -> Result> { + match fs::read(&self.path) { + Ok(bytes) => Ok(Some(serde_json::from_slice(&bytes)?)), + Err(error) if error.kind() == io::ErrorKind::NotFound => Ok(None), + Err(error) => Err(Error::Io(error)), + } + } + + fn save(&self, cursor: &InboxCursor) -> Result<()> { + validate_cursor(cursor)?; + if let Some(parent) = self + .path + .parent() + .filter(|path| !path.as_os_str().is_empty()) + { + fs::create_dir_all(parent).map_err(local_write_error)?; + } + let parent = self + .path + .parent() + .filter(|path| !path.as_os_str().is_empty()) + .unwrap_or_else(|| Path::new(".")); + let bytes = serde_json::to_vec(cursor)?; + let mut temp = tempfile::NamedTempFile::new_in(parent).map_err(local_write_error)?; + temp.write_all(&bytes).map_err(local_write_error)?; + temp.as_file().sync_all().map_err(local_write_error)?; + temp.persist(&self.path) + .map_err(|error| local_write_error(error.error))?; + #[cfg(unix)] + fs::File::open(parent) + .and_then(|directory| directory.sync_all()) + .map_err(local_write_error)?; + Ok(()) + } +} + +fn local_write_error(source: io::Error) -> Error { + Error::Io(source) +} + +/// Identity-bound API exposing only the P0 chat/inbox operations. +/// +/// There is deliberately no generic `call`, daemon shutdown, configuration, +/// or credential operation on this adapter. +pub struct IdentityBoundClient { + client: CCCCClient, + identity: H, + principal: AuthenticatedPrincipal, +} + +impl IdentityBoundClient { + pub fn new(client: CCCCClient, identity: H) -> Result { + let principal = identity.principal()?; + if principal.subject.trim().is_empty() || principal.issuer.trim().is_empty() { + return Err(Error::Incompatible( + "authenticated principal subject and issuer must be non-empty".into(), + )); + } + Ok(Self { + client, + identity, + principal, + }) + } + + pub fn principal(&self) -> &AuthenticatedPrincipal { + &self.principal + } + + fn call(&self, operation: &str, mut args: Map) -> Result> { + let current = bind_identity(&self.identity, operation, &mut args)?; + if current != self.principal { + return Err(Error::Incompatible( + "workload identity principal changed during the session".into(), + )); + } + self.client.call(operation, args) + } + + fn call_typed Deserialize<'de>>( + &self, + operation: &str, + args: Map, + ) -> Result { + Ok(serde_json::from_value(Value::Object( + self.call(operation, args)?, + ))?) + } + + /// Send with a daemon-stable caller key. Reuse the same key to reconcile + /// an `Unknown` write outcome; never generate a new key for that retry. + pub fn send_idempotent( + &self, + group_id: &str, + text: &str, + idempotency_key: &str, + ) -> Result { + validate_key(idempotency_key)?; + self.call_typed( + "send", + object([ + ("group_id", Value::String(group_id.into())), + ("text", Value::String(text.into())), + ("client_id", Value::String(idempotency_key.into())), + ]), + ) + } + + pub fn reconcile_send( + &self, + group_id: &str, + text: &str, + idempotency_key: &str, + ) -> Result { + self.send_idempotent(group_id, text, idempotency_key) + } + + pub fn reply_idempotent( + &self, + group_id: &str, + reply_to: &str, + text: &str, + idempotency_key: &str, + ) -> Result { + validate_key(idempotency_key)?; + self.call_typed( + "reply", + object([ + ("group_id", Value::String(group_id.into())), + ("reply_to", Value::String(reply_to.into())), + ("text", Value::String(text.into())), + ("client_id", Value::String(idempotency_key.into())), + ]), + ) + } + + pub fn reconcile_reply( + &self, + group_id: &str, + reply_to: &str, + text: &str, + idempotency_key: &str, + ) -> Result { + self.reply_idempotent(group_id, reply_to, text, idempotency_key) + } + + pub fn inbox_list(&self, group_id: &str, actor_id: &str, limit: u32) -> Result { + self.call_typed( + "inbox_list", + object([ + ("group_id", Value::String(group_id.into())), + ("actor_id", Value::String(actor_id.into())), + ("limit", Value::from(limit)), + ]), + ) + } + + pub fn mark_read_idempotent( + &self, + group_id: &str, + actor_id: &str, + cursor: &InboxCursor, + idempotency_key: &str, + ) -> Result { + validate_key(idempotency_key)?; + self.call_typed( + "inbox_mark_read", + object([ + ("group_id", Value::String(group_id.into())), + ("actor_id", Value::String(actor_id.into())), + ("event_id", Value::String(cursor.event_id.clone())), + ("idempotency_key", Value::String(idempotency_key.into())), + ]), + ) + } + + /// Query the daemon cursor and issue `inbox_mark_read` only when the target + /// is not already covered. This closes the write-then-disconnect window + /// even for daemon builds that do not replay `idempotency_key` themselves. + pub fn reconcile_mark_read( + &self, + group_id: &str, + actor_id: &str, + cursor: &InboxCursor, + idempotency_key: &str, + ) -> Result { + validate_cursor(cursor)?; + validate_key(idempotency_key)?; + let page = self.inbox_list(group_id, actor_id, 1)?; + match self.remote_cursor_relation(group_id, actor_id, &page.cursor, cursor)? { + CursorRelation::Covers => Ok(MarkReadReconciliation { + cursor: page.cursor, + wrote: false, + result: None, + }), + CursorRelation::Behind => { + let result = + self.mark_read_idempotent(group_id, actor_id, cursor, idempotency_key)?; + Ok(MarkReadReconciliation { + cursor: result.cursor.clone(), + wrote: true, + result: Some(result), + }) + } + } + } + + fn remote_cursor_relation( + &self, + group_id: &str, + actor_id: &str, + remote: &InboxCursor, + target: &InboxCursor, + ) -> Result { + if remote.event_id == target.event_id { + return Ok(CursorRelation::Covers); + } + if remote.event_id.is_empty() { + return Ok(CursorRelation::Behind); + } + + let status: MessageReadStatus = self.call_typed( + "message_read_status", + object([ + ("group_id", Value::String(group_id.into())), + ("event_id", Value::String(target.event_id.clone())), + ]), + )?; + if let Some(read) = status.read_status.get(actor_id) { + return Ok(if *read { + CursorRelation::Covers + } else { + CursorRelation::Behind + }); + } + + self.remote_cursor_relation_from_ledger(group_id, remote, target) + } + + fn remote_cursor_relation_from_ledger( + &self, + group_id: &str, + remote: &InboxCursor, + target: &InboxCursor, + ) -> Result { + // `message_read_status` is authoritative for chat messages. Native + // daemon 0.4.33 omits system notifications from that result, so prove + // their order against the immutable ledger without comparing IDs. + self.ledger_window(group_id, &target.event_id, 0)?; + let mut center = remote.event_id.clone(); + for _ in 0..10_000 { + let page = self.ledger_window(group_id, ¢er, 200)?; + if page.events.iter().any(|event| event.id == target.event_id) { + return Ok(CursorRelation::Behind); + } + if !page.has_more_after { + return Ok(CursorRelation::Covers); + } + let next = page + .events + .last() + .map(|event| event.id.clone()) + .filter(|event_id| !event_id.is_empty() && event_id != ¢er) + .ok_or_else(|| { + Error::ReconciliationRequired( + "ledger cursor scan did not advance while reconciling inbox state".into(), + ) + })?; + center = next; + } + Err(Error::ReconciliationRequired( + "ledger cursor scan exceeded 2,000,000 events".into(), + )) + } + + fn ledger_window(&self, group_id: &str, center: &str, after: u32) -> Result { + self.call_typed( + "ledger_window", + object([ + ("group_id", Value::String(group_id.into())), + ("center", Value::String(center.into())), + ("kind", Value::String("all".into())), + ("before", Value::from(0)), + ("after", Value::from(after)), + ]), + ) + } + + pub fn persistent_inbox( + &self, + group_id: impl Into, + actor_id: impl Into, + store: S, + ) -> PersistentInbox<'_, H, S> { + PersistentInbox { + client: self, + group_id: group_id.into(), + actor_id: actor_id.into(), + store, + reconciled: false, + } + } +} + +/// At-least-once inbox consumer with a durable, caller-owned checkpoint. +pub struct PersistentInbox<'a, H, S> { + client: &'a IdentityBoundClient, + group_id: String, + actor_id: String, + store: S, + reconciled: bool, +} + +impl PersistentInbox<'_, H, S> { + pub fn poll(&mut self, limit: u32) -> Result { + if !self.reconciled { + if let Some(cursor) = self.store.load()? { + self.client.reconcile_mark_read( + &self.group_id, + &self.actor_id, + &cursor, + &cursor_key(&self.group_id, &self.actor_id, &cursor), + )?; + } + self.reconciled = true; + } + self.client + .inbox_list(&self.group_id, &self.actor_id, limit) + } + + /// Commit only after the event's side effect has completed. The local + /// checkpoint is stored first, so a crash cannot make a processed event + /// disappear merely because daemon cursor state moved ahead. + pub fn commit(&mut self, event: &Event) -> Result { + let incoming = InboxCursor { + event_id: event.id.clone(), + ts: event.ts.clone(), + }; + validate_cursor(&incoming)?; + let (cursor, needs_mark) = match self.store.load()? { + Some(current) => match local_cursor_relation(¤t, &incoming)? { + CursorRelation::Covers => (current, false), + CursorRelation::Behind => { + self.store.save(&incoming)?; + (incoming, true) + } + }, + None => { + self.store.save(&incoming)?; + (incoming, true) + } + }; + if !needs_mark { + return Ok(MarkReadResult { + cursor, + event: None, + replayed: true, + }); + } + let result = self.client.mark_read_idempotent( + &self.group_id, + &self.actor_id, + &cursor, + &cursor_key(&self.group_id, &self.actor_id, &cursor), + ); + if result.is_err() { + self.reconciled = false; + } + result + } +} + +fn cursor_key(group_id: &str, actor_id: &str, cursor: &InboxCursor) -> String { + format!("cursor:{group_id}:{actor_id}:{}", cursor.event_id) +} + +enum CursorRelation { + Covers, + Behind, +} + +fn local_cursor_relation(remote: &InboxCursor, target: &InboxCursor) -> Result { + if remote.event_id == target.event_id || remote.ts > target.ts { + return Ok(CursorRelation::Covers); + } + if remote.ts < target.ts || (remote.ts.is_empty() && remote.event_id.is_empty()) { + return Ok(CursorRelation::Behind); + } + Err(Error::ReconciliationRequired(format!( + "daemon and local cursors have the same timestamp but different event IDs: remote={}, local={}", + remote.event_id, target.event_id + ))) +} + +fn deserialize_nullable_string<'de, D>(deserializer: D) -> std::result::Result +where + D: Deserializer<'de>, +{ + Ok(Option::::deserialize(deserializer)?.unwrap_or_default()) +} + +fn validate_key(key: &str) -> Result<()> { + if key.trim().is_empty() || key.len() > 256 { + return Err(Error::Incompatible( + "idempotency key must contain 1..=256 bytes".into(), + )); + } + Ok(()) +} + +fn validate_cursor(cursor: &InboxCursor) -> Result<()> { + if cursor.event_id.trim().is_empty() || cursor.ts.trim().is_empty() { + return Err(Error::Incompatible( + "persisted inbox cursor needs non-empty event_id and ts".into(), + )); + } + Ok(()) +} + +fn object(entries: [(&str, Value); N]) -> Map { + entries + .into_iter() + .map(|(key, value)| (key.to_owned(), value)) + .collect() +} + +#[cfg(test)] +mod tests { + use super::*; + use std::io::{BufRead, BufReader, Write}; + use std::net::TcpListener; + use std::sync::{Arc, Mutex}; + use std::thread; + use std::time::{SystemTime, UNIX_EPOCH}; + + #[derive(Clone)] + struct SignedIdentity; + + impl WorkloadIdentityHook for SignedIdentity { + fn principal(&self) -> Result { + AuthenticatedPrincipal::new("workload:aquant", "test-spiffe") + } + + fn evidence( + &self, + operation: &str, + _args: &Map, + ) -> Result { + crate::WorkloadIdentityEvidence::new( + "workload_identity", + serde_json::json!({"scheme": "test", "signature": format!("sig:{operation}")}), + ) + } + } + + fn server( + responses: Vec<&'static str>, + ) -> ( + crate::DaemonEndpoint, + Arc>>, + thread::JoinHandle<()>, + ) { + let listener = TcpListener::bind("127.0.0.1:0").expect("bind"); + let address = listener.local_addr().expect("address"); + let requests = Arc::new(Mutex::new(Vec::new())); + let captured = Arc::clone(&requests); + let handle = thread::spawn(move || { + for response in responses { + let (mut stream, _) = listener.accept().expect("accept"); + let mut request = String::new(); + BufReader::new(&mut stream) + .read_line(&mut request) + .expect("request"); + captured + .lock() + .expect("lock") + .push(serde_json::from_str(&request).expect("JSON request")); + stream.write_all(response.as_bytes()).expect("response"); + } + }); + ( + crate::DaemonEndpoint::Tcp { + host: "127.0.0.1".into(), + port: address.port(), + }, + requests, + handle, + ) + } + + fn fixture_response(operation: &str, name: &str) -> &'static str { + match (operation, name) { + ("inbox_list", "fresh_response") => concat!( + r#"{"v":1,"ok":true,"result":{"messages":[],"cursor":{"event_id":null,"ts":""}}}"#, + "\n" + ), + ("inbox_list", "remote_ahead_response") => concat!( + r#"{"v":1,"ok":true,"result":{"messages":[],"cursor":{"event_id":"e8","ts":""}}}"#, + "\n" + ), + ("message_read_status", "covered_response") => concat!( + r#"{"v":1,"ok":true,"result":{"event_id":"e7","read_status":{"a1":true}}}"#, + "\n" + ), + ("send", "initial_response") => concat!( + r#"{"v":1,"ok":true,"result":{"event":{"id":"e1","ts":"2026-08-05T01:00:00Z","kind":"chat.message","by":"workload:aquant","data":{"text":"hello","client_id":"send:stable-1"}},"delivery":{"accepted":true,"state":"queued"}}}"#, + "\n" + ), + ("send", "duplicate_response") => concat!( + r#"{"v":1,"ok":true,"result":{"event":{"id":"e1","ts":"2026-08-05T01:00:00Z","kind":"chat.message","by":"workload:aquant","data":{"text":"hello","client_id":"send:stable-1"}},"delivery":{"accepted":true,"state":"duplicate"},"duplicate":true}}"#, + "\n" + ), + _ => panic!("unknown fixture response: {operation}.{name}"), + } + } + + fn event_json(duplicate: bool) -> &'static str { + fixture_response( + "send", + if duplicate { + "duplicate_response" + } else { + "initial_response" + }, + ) + } + + #[test] + fn identity_is_bound_and_send_reconciliation_reuses_the_key() { + let first = event_json(false); + let replay = event_json(true); + let (endpoint, requests, handle) = server(vec![first, replay]); + let client = IdentityBoundClient::new(CCCCClient::new(endpoint), SignedIdentity) + .expect("bound client"); + + client + .send_idempotent("g1", "hello", "send:stable-1") + .expect("send"); + let result = client + .reconcile_send("g1", "hello", "send:stable-1") + .expect("reconcile"); + assert!(result.replayed); + handle.join().expect("server"); + + let requests = requests.lock().expect("lock"); + for request in requests.iter() { + assert_eq!(request["args"]["by"], "workload:aquant"); + assert_eq!(request["args"]["client_id"], "send:stable-1"); + assert_eq!(request["args"]["workload_identity"]["scheme"], "test"); + } + } + + #[test] + fn persisted_cursor_queries_before_replaying_an_unknown_mark() { + let truncated = ""; + let remote_behind = fixture_response("inbox_list", "fresh_response"); + let marked = "{\"v\":1,\"ok\":true,\"result\":{\"cursor\":{\"event_id\":\"e7\",\"ts\":\"2026-08-05T01:00:00Z\"},\"event\":null}}\n"; + let page = "{\"v\":1,\"ok\":true,\"result\":{\"messages\":[],\"cursor\":{\"event_id\":\"e7\",\"ts\":\"\"}}}\n"; + let (first_endpoint, first_requests, first_handle) = server(vec![truncated]); + let client = IdentityBoundClient::new(CCCCClient::new(first_endpoint), SignedIdentity) + .expect("bound client"); + let path = temp_path(); + let store = FileCursorStore::new(&path); + let mut first = client.persistent_inbox("g1", "a1", store.clone()); + let event = Event { + id: "e7".into(), + ts: "2026-08-05T01:00:00Z".into(), + kind: "chat.message".into(), + by: "peer".into(), + data: Value::Null, + extra: Map::new(), + }; + let error = first.commit(&event).expect_err("disconnect after write"); + assert!(matches!(error, Error::OutcomeUnknown { .. })); + assert_eq!( + store.load().expect("load"), + Some(InboxCursor { + event_id: "e7".into(), + ts: "2026-08-05T01:00:00Z".into(), + }) + ); + first_handle.join().expect("first server"); + + // A new process freshly discovers a different endpoint and reuses only + // the durable cursor and the stable workload identity. + let (second_endpoint, second_requests, second_handle) = + server(vec![remote_behind, marked, page]); + let restarted_client = + IdentityBoundClient::new(CCCCClient::new(second_endpoint), SignedIdentity) + .expect("restarted client"); + let mut restarted = restarted_client.persistent_inbox("g1", "a1", store); + restarted.poll(50).expect("query, reconcile, then poll"); + second_handle.join().expect("second server"); + let first_requests = first_requests.lock().expect("first lock"); + let second_requests = second_requests.lock().expect("second lock"); + assert_eq!(first_requests[0]["op"], "inbox_mark_read"); + assert_eq!(second_requests[0]["op"], "inbox_list"); + assert_eq!(second_requests[1]["op"], "inbox_mark_read"); + assert_eq!(second_requests[2]["op"], "inbox_list"); + assert_eq!( + first_requests[0]["args"]["idempotency_key"], + second_requests[1]["args"]["idempotency_key"] + ); + fs::remove_file(path).expect("cleanup"); + } + + #[test] + fn reconciliation_does_not_repeat_an_already_applied_mark() { + let covered = "{\"v\":1,\"ok\":true,\"result\":{\"messages\":[],\"cursor\":{\"event_id\":\"e7\",\"ts\":\"\"}}}\n"; + let page = covered; + let (endpoint, requests, handle) = server(vec![covered, page]); + let client = IdentityBoundClient::new(CCCCClient::new(endpoint), SignedIdentity) + .expect("bound client"); + let path = temp_path(); + let store = FileCursorStore::new(&path); + store + .save(&InboxCursor { + event_id: "e7".into(), + ts: "2026-08-05T01:00:00Z".into(), + }) + .expect("seed cursor"); + client + .persistent_inbox("g1", "a1", store) + .poll(50) + .expect("reconcile"); + handle.join().expect("server"); + assert!(requests + .lock() + .expect("lock") + .iter() + .all(|request| request["op"] == "inbox_list")); + fs::remove_file(path).expect("cleanup"); + } + + #[test] + fn ambiguous_out_of_order_cursor_fails_closed() { + let remote = InboxCursor { + event_id: "remote".into(), + ts: "2026-08-05T01:00:00Z".into(), + }; + let local = InboxCursor { + event_id: "local".into(), + ts: remote.ts.clone(), + }; + assert!(matches!( + local_cursor_relation(&remote, &local), + Err(Error::ReconciliationRequired(_)) + )); + } + + #[test] + fn fresh_actor_poll_accepts_native_nullable_cursor() { + let fresh = fixture_response("inbox_list", "fresh_response"); + let (endpoint, requests, handle) = server(vec![fresh]); + let client = IdentityBoundClient::new(CCCCClient::new(endpoint), SignedIdentity) + .expect("bound client"); + let path = temp_path(); + let page = client + .persistent_inbox("g1", "a1", FileCursorStore::new(&path)) + .poll(50) + .expect("fresh poll"); + handle.join().expect("server"); + assert_eq!(page.cursor.event_id, ""); + assert_eq!(page.cursor.ts, ""); + assert_eq!(requests.lock().expect("lock")[0]["op"], "inbox_list"); + assert!(!path.exists()); + } + + #[test] + fn remote_ahead_cursor_uses_read_status_without_repeating_mark() { + let remote_ahead = fixture_response("inbox_list", "remote_ahead_response"); + let read = fixture_response("message_read_status", "covered_response"); + let page = remote_ahead; + let (endpoint, requests, handle) = server(vec![remote_ahead, read, page]); + let client = IdentityBoundClient::new(CCCCClient::new(endpoint), SignedIdentity) + .expect("bound client"); + let path = temp_path(); + let store = FileCursorStore::new(&path); + store + .save(&InboxCursor { + event_id: "e7".into(), + ts: "2026-08-05T01:00:00Z".into(), + }) + .expect("seed cursor"); + client + .persistent_inbox("g1", "a1", store) + .poll(50) + .expect("remote cursor covers local target"); + handle.join().expect("server"); + let operations = requests + .lock() + .expect("lock") + .iter() + .map(|request| request["op"].as_str().unwrap_or_default().to_owned()) + .collect::>(); + assert_eq!( + operations, + ["inbox_list", "message_read_status", "inbox_list"] + ); + fs::remove_file(path).expect("cleanup"); + } + + #[test] + fn notification_cursor_falls_back_to_provable_ledger_order() { + let remote_ahead = "{\"v\":1,\"ok\":true,\"result\":{\"messages\":[],\"cursor\":{\"event_id\":\"e8\",\"ts\":\"\"}}}\n"; + let no_chat_status = + "{\"v\":1,\"ok\":true,\"result\":{\"event_id\":\"e7\",\"read_status\":{}}}\n"; + let target = "{\"v\":1,\"ok\":true,\"result\":{\"center_id\":\"e7\",\"center_index\":0,\"events\":[{\"id\":\"e7\",\"ts\":\"2026-08-05T01:00:00Z\",\"kind\":\"system.notify\",\"by\":\"peer\",\"data\":{}}],\"has_more_before\":true,\"has_more_after\":true,\"count\":1}}\n"; + let remote = "{\"v\":1,\"ok\":true,\"result\":{\"center_id\":\"e8\",\"center_index\":0,\"events\":[{\"id\":\"e8\",\"ts\":\"2026-08-05T02:00:00Z\",\"kind\":\"chat.message\",\"by\":\"peer\",\"data\":{}}],\"has_more_before\":true,\"has_more_after\":false,\"count\":1}}\n"; + let page = remote_ahead; + let (endpoint, requests, handle) = + server(vec![remote_ahead, no_chat_status, target, remote, page]); + let client = IdentityBoundClient::new(CCCCClient::new(endpoint), SignedIdentity) + .expect("bound client"); + let path = temp_path(); + let store = FileCursorStore::new(&path); + store + .save(&InboxCursor { + event_id: "e7".into(), + ts: "2026-08-05T01:00:00Z".into(), + }) + .expect("seed cursor"); + client + .persistent_inbox("g1", "a1", store) + .poll(50) + .expect("notification cursor order"); + handle.join().expect("server"); + let requests = requests.lock().expect("lock"); + assert_eq!(requests[2]["op"], "ledger_window"); + assert_eq!(requests[2]["args"]["center"], "e7"); + assert_eq!(requests[3]["op"], "ledger_window"); + assert_eq!(requests[3]["args"]["center"], "e8"); + assert!(requests + .iter() + .all(|request| request["op"] != "inbox_mark_read")); + fs::remove_file(path).expect("cleanup"); + } + + #[test] + fn file_cursor_store_replaces_an_existing_checkpoint() { + let path = temp_path(); + let store = FileCursorStore::new(&path); + let first = InboxCursor { + event_id: "e1".into(), + ts: "2026-08-05T01:00:00Z".into(), + }; + let second = InboxCursor { + event_id: "e2".into(), + ts: "2026-08-05T02:00:00Z".into(), + }; + store.save(&first).expect("first save"); + store.save(&second).expect("atomic replace"); + assert_eq!(store.load().expect("load"), Some(second)); + fs::remove_file(path).expect("cleanup"); + } + + #[test] + fn duplicate_or_older_commit_never_regresses_the_local_cursor() { + let (endpoint, requests, handle) = server(vec![]); + let client = IdentityBoundClient::new(CCCCClient::new(endpoint), SignedIdentity) + .expect("bound client"); + let path = temp_path(); + let store = FileCursorStore::new(&path); + let latest = InboxCursor { + event_id: "e8".into(), + ts: "2026-08-05T02:00:00Z".into(), + }; + store.save(&latest).expect("seed"); + let older = Event { + id: "e7".into(), + ts: "2026-08-05T01:00:00Z".into(), + kind: "chat.message".into(), + by: "peer".into(), + data: Value::Null, + extra: Map::new(), + }; + client + .persistent_inbox("g1", "a1", store.clone()) + .commit(&older) + .expect("commit is monotonic"); + handle.join().expect("server"); + assert_eq!(store.load().expect("load"), Some(latest)); + assert!(requests.lock().expect("lock").is_empty()); + fs::remove_file(path).expect("cleanup"); + } + + fn temp_path() -> PathBuf { + let nonce = SystemTime::now() + .duration_since(UNIX_EPOCH) + .expect("clock") + .as_nanos(); + std::env::temp_dir().join(format!( + "cccc-sdk-cursor-{}-{nonce}.json", + std::process::id() + )) + } +} diff --git a/rust/tests/live_reliability.rs b/rust/tests/live_reliability.rs new file mode 100644 index 0000000..a257d30 --- /dev/null +++ b/rust/tests/live_reliability.rs @@ -0,0 +1,162 @@ +use cccc_sdk::{ + AuthenticatedPrincipal, CCCCClient, CursorStore, IdentityBoundClient, InboxCursor, Result, + WorkloadIdentityEvidence, WorkloadIdentityHook, +}; +use serde_json::{json, Map, Value}; + +struct LiveTestIdentity(&'static str); + +#[derive(Clone, Copy)] +struct EmptyCursorStore; + +impl CursorStore for EmptyCursorStore { + fn load(&self) -> Result> { + Ok(None) + } + + fn save(&self, _cursor: &InboxCursor) -> Result<()> { + Ok(()) + } +} + +impl WorkloadIdentityHook for LiveTestIdentity { + fn principal(&self) -> Result { + AuthenticatedPrincipal::new(self.0, "cccc-sdk-live-test") + } + + fn evidence( + &self, + operation: &str, + _args: &Map, + ) -> Result { + WorkloadIdentityEvidence::new( + "workload_identity", + json!({"test_only": true, "operation": operation}), + ) + } +} + +fn object(entries: impl IntoIterator) -> Map { + entries + .into_iter() + .map(|(key, value)| (key.to_owned(), value)) + .collect() +} + +#[test] +fn daemon_0433_replays_real_writes_and_reconciles_mark_read() { + if std::env::var_os("CCCC_RUN_LIVE_RELIABILITY").is_none() { + return; + } + + let raw = CCCCClient::discover().expect("discover live daemon"); + let created = raw + .call( + "group_create", + object([ + ("title", json!("cccc-sdk 0.0.2 live reliability test")), + ("by", json!("user")), + ]), + ) + .expect("create disposable group"); + let group_id = created["group_id"].as_str().expect("group_id").to_owned(); + + let result = (|| { + raw.call( + "group_start", + object([("group_id", json!(group_id)), ("by", json!("user"))]), + )?; + raw.call( + "actor_add", + object([ + ("group_id", json!(group_id)), + ("actor_id", json!("sdk-live")), + ("runtime", json!("custom")), + ("runner", json!("headless")), + ("command", json!(["/usr/bin/true"])), + ("by", json!("user")), + ]), + )?; + let client = IdentityBoundClient::new(raw.clone(), LiveTestIdentity("user"))?; + let peer = IdentityBoundClient::new(raw.clone(), LiveTestIdentity("sdk-live"))?; + let fresh = peer + .persistent_inbox(&group_id, "sdk-live", EmptyCursorStore) + .poll(1)?; + if !fresh.cursor.event_id.is_empty() || !fresh.cursor.ts.is_empty() { + return Err(cccc_sdk::Error::ReconciliationRequired( + "fresh native daemon cursor was not normalized to empty strings".into(), + )); + } + + let send_key = format!("cccc-sdk-live:{group_id}:send"); + let first = client.send_idempotent(&group_id, "live idempotency probe", &send_key)?; + let replay = client.reconcile_send(&group_id, "live idempotency probe", &send_key)?; + if first.event.id != replay.event.id { + return Err(cccc_sdk::Error::ReconciliationRequired(format!( + "send key produced two events: {} and {}", + first.event.id, replay.event.id + ))); + } + if !replay.replayed { + return Err(cccc_sdk::Error::ReconciliationRequired( + "native duplicate=true was not exposed as replayed=true".into(), + )); + } + + let second_key = format!("cccc-sdk-live:{group_id}:send:second"); + let second = + client.send_idempotent(&group_id, "live remote-ahead cursor probe", &second_key)?; + + let reply_key = format!("cccc-sdk-live:{group_id}:reply"); + let first_reply = peer.reply_idempotent( + &group_id, + &first.event.id, + "live reply idempotency probe", + &reply_key, + )?; + let replay_reply = peer.reconcile_reply( + &group_id, + &first.event.id, + "live reply idempotency probe", + &reply_key, + )?; + if first_reply.event.id != replay_reply.event.id { + return Err(cccc_sdk::Error::ReconciliationRequired(format!( + "reply key produced two events: {} and {}", + first_reply.event.id, replay_reply.event.id + ))); + } + if !replay_reply.replayed { + return Err(cccc_sdk::Error::ReconciliationRequired( + "native reply duplicate=true was not exposed as replayed=true".into(), + )); + } + + let first_cursor = InboxCursor { + event_id: first.event.id.clone(), + ts: first.event.ts.clone(), + }; + let second_cursor = InboxCursor { + event_id: second.event.id, + ts: second.event.ts, + }; + let second_mark_key = format!("cccc-sdk-live:{group_id}:mark:{}", second_cursor.event_id); + peer.mark_read_idempotent(&group_id, "sdk-live", &second_cursor, &second_mark_key)?; + let first_mark_key = format!("cccc-sdk-live:{group_id}:mark:{}", first_cursor.event_id); + let reconciled = + peer.reconcile_mark_read(&group_id, "sdk-live", &first_cursor, &first_mark_key)?; + if reconciled.wrote { + return Err(cccc_sdk::Error::ReconciliationRequired( + "remote-ahead read cursor repeated an older mark_read".into(), + )); + } + Ok::<(), cccc_sdk::Error>(()) + })(); + + let cleanup = raw.call( + "group_delete", + object([("group_id", json!(group_id)), ("by", json!("user"))]), + ); + result.expect("live reliability assertions"); + cleanup.expect("delete disposable group"); +} diff --git a/scripts/check_sdk_hardening.py b/scripts/check_sdk_hardening.py new file mode 100644 index 0000000..bc3d9a3 --- /dev/null +++ b/scripts/check_sdk_hardening.py @@ -0,0 +1,131 @@ +#!/usr/bin/env python3 +from __future__ import annotations + +import ast +import hashlib +import json +import re +from pathlib import Path + + +ROOT = Path(__file__).resolve().parents[1] +FIXTURE = ROOT / "spec" / "SDK_DAEMON_TARGET_0_4_33.json" +FIXTURE_SHA256 = "616f0c81a73204c5478becfdfe671e4b28d70e232447263c0f594dc53ab78d51" + + +def read(relative: str) -> str: + return (ROOT / relative).read_text(encoding="utf-8") + + +def python_method(source: str, name: str) -> ast.FunctionDef | None: + tree = ast.parse(source) + return next( + (node for node in ast.walk(tree) if isinstance(node, ast.FunctionDef) and node.name == name), + None, + ) + + +def main() -> int: + errors: list[str] = [] + + fixture_bytes = FIXTURE.read_bytes() + if hashlib.sha256(fixture_bytes).hexdigest() != FIXTURE_SHA256: + errors.append("0.4.33 daemon fixture hash changed without an explicit contract review") + fixture = json.loads(fixture_bytes) + if fixture.get("target", {}).get("version") != "0.4.33": + errors.append("daemon fixture must target exactly 0.4.33") + + completion = fixture.get("operations", {}).get("web_model_runtime_complete_turn", {}) + required = {"group_id", "actor_id", "turn_id", "delivery_id"} + if not required.issubset(set(completion.get("required_args", []))): + errors.append("completion fixture is missing required daemon arguments") + request_args = completion.get("request", {}).get("args", {}) + result = completion.get("completion_response", {}).get("result", {}) + if ( + completion.get("replay_key") != "delivery_id" + or result.get("delivery_id") != request_args.get("delivery_id") + ): + errors.append("completion fixture does not preserve delivery_id across replay") + + python_source = read("python/src/cccc_sdk/client_0430_runtime_ops.py") + method = python_method(python_source, "web_model_runtime_complete_turn") + if method is None: + errors.append("Python completion wrapper is missing") + else: + parameters = { + arg.arg + for arg in (*method.args.posonlyargs, *method.args.args, *method.args.kwonlyargs) + } + if not required.issubset(parameters): + errors.append("Python completion wrapper does not require delivery_id") + if not any( + isinstance(node, ast.Dict) + and any(isinstance(key, ast.Constant) and key.value == "delivery_id" for key in node.keys) + for node in ast.walk(method) + ): + errors.append("Python completion wrapper does not map delivery_id") + + ts_types = read("ts/src/types.ts") + ts_runtime = read("ts/src/client_0430_runtime_ops.ts") + if not re.search(r"\bdeliveryId\s*:\s*string\s*;", ts_types): + errors.append("TypeScript deliveryId is missing or optional") + if not re.search(r"delivery_id\s*:\s*options\.deliveryId\b", ts_runtime): + errors.append("TypeScript completion wrapper does not map deliveryId") + + transport = read("ts/src/transport.ts") + transport_markers = { + "remainingTimeout(deadline)": "TypeScript transport does not share one deadline", + "signal?.removeEventListener('abort', onAbort)": "TypeScript handshake abort cleanup is missing", + "assertBufferedLineLimit(remaining)": "TypeScript handshake remainder is not byte-capped", + "initialBuffer: Buffer": "TypeScript stream buffering is not byte-accurate", + } + for marker, message in transport_markers.items(): + if marker not in transport: + errors.append(message) + if "removeAllListeners" in transport: + errors.append("TypeScript transport still removes unrelated socket listeners") + + reliable = read("rust/src/reliable.rs") + rust_markers = { + 'alias = "duplicate"': "Rust does not map duplicate to replayed", + 'deserialize_with = "deserialize_nullable_string"': "Rust cursor is not nullable-wire compatible", + '"message_read_status"': "Rust remote-ahead reconciliation lacks read-status verification", + '"ledger_window"': "Rust notification reconciliation lacks ledger-order fallback", + "NamedTempFile::new_in": "Rust cursor writes do not use a unique same-directory temp file", + ".persist(&self.path)": "Rust cursor writes do not atomically replace the checkpoint", + } + for marker, message in rust_markers.items(): + if marker not in reliable: + errors.append(message) + if 'tempfile = "=3.20.0"' not in read("rust/Cargo.toml"): + errors.append("Rust tempfile dependency is not pinned for the declared MSRV") + + workflows = { + name: read(f".github/workflows/{name}") + for name in ("python-integration.yml", "ts-ci.yml", "rust-ci.yml") + } + install = "cargo install cccc --version '=0.4.33' --locked" + for name, workflow in workflows.items(): + if install not in workflow: + errors.append(f"{name} does not test the exact native 0.4.33 daemon") + rust_ci = workflows["rust-ci.yml"] + for marker, message in ( + ('toolchain: "1.74.0"', "Rust CI does not enforce the declared 1.74 MSRV"), + ("runs-on: windows-latest", "Rust CI does not test Windows cursor replacement"), + ('CCCC_RUN_LIVE_RELIABILITY: "1"', "Rust CI does not execute live reliability assertions"), + ("--test live_reliability", "Rust CI does not run the native reliability test"), + ): + if marker not in rust_ci: + errors.append(message) + + if errors: + for error in errors: + print(f"ERROR: {error}") + return 1 + + print("SDK hardening contract OK: completion replay, stream safety, reliable cursor, CI matrix") + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/scripts/check_specs_against_cccc.sh b/scripts/check_specs_against_cccc.sh index f911332..f98a742 100755 --- a/scripts/check_specs_against_cccc.sh +++ b/scripts/check_specs_against_cccc.sh @@ -2,18 +2,31 @@ set -euo pipefail CCCC_REPO="${1:-../cccc}" +CCCC_REF="${2:-}" SRC="${CCCC_REPO%/}/docs/standards" ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)" DST="${ROOT}/spec" -if [[ ! -d "${SRC}" ]]; then +if [[ -n "${CCCC_REF}" ]] && ! git -C "${CCCC_REPO}" rev-parse --verify "${CCCC_REF}^{commit}" >/dev/null 2>&1; then + echo "error: unknown CCCC git ref: ${CCCC_REF}" >&2 + exit 2 +fi + +if [[ -z "${CCCC_REF}" && ! -d "${SRC}" ]]; then echo "error: cannot find CCCC specs at: ${SRC}" >&2 exit 2 fi status=0 for name in CCCS_V1.md CCCC_DAEMON_IPC_V1.md CCCC_CONTEXT_OPS_V1.md; do - if ! cmp -s "${SRC}/${name}" "${DST}/${name}"; then + if [[ -n "${CCCC_REF}" ]]; then + if git -C "${CCCC_REPO}" show "${CCCC_REF}:docs/standards/${name}" | cmp -s - "${DST}/${name}"; then + continue + fi + echo "error: spec/${name} has drifted from CCCC core" >&2 + diff -u "${DST}/${name}" <(git -C "${CCCC_REPO}" show "${CCCC_REF}:docs/standards/${name}") || true + status=1 + elif ! cmp -s "${SRC}/${name}" "${DST}/${name}"; then echo "error: spec/${name} has drifted from CCCC core" >&2 diff -u "${DST}/${name}" "${SRC}/${name}" || true status=1 @@ -24,4 +37,8 @@ if [[ "${status}" -ne 0 ]]; then exit "${status}" fi -echo "All mirrored CCCC standards match core." +if [[ -n "${CCCC_REF}" ]]; then + echo "All mirrored CCCC standards match ${CCCC_REPO}@${CCCC_REF}." +else + echo "All mirrored CCCC standards match core." +fi diff --git a/scripts/sync_specs_from_cccc.sh b/scripts/sync_specs_from_cccc.sh index d83411e..1cffd6c 100755 --- a/scripts/sync_specs_from_cccc.sh +++ b/scripts/sync_specs_from_cccc.sh @@ -2,13 +2,38 @@ set -euo pipefail CCCC_REPO="${1:-../cccc}" -SRC="${CCCC_REPO%/}/docs/standards" +CCCC_REF="${2:-}" ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)" DST="${ROOT}/spec" +TMP="" + +cleanup() { + if [[ -n "${TMP}" && -d "${TMP}" ]]; then + rm -rf -- "${TMP}" + fi +} +trap cleanup EXIT + +if [[ -n "${CCCC_REF}" ]]; then + if ! git -C "${CCCC_REPO}" rev-parse --verify "${CCCC_REF}^{commit}" >/dev/null 2>&1; then + echo "error: unknown CCCC git ref: ${CCCC_REF}" >&2 + exit 2 + fi + TMP="$(mktemp -d)" + git -C "${CCCC_REPO}" archive "${CCCC_REF}" \ + docs/standards/CCCS_V1.md \ + docs/standards/CCCC_DAEMON_IPC_V1.md \ + docs/standards/CCCC_CONTEXT_OPS_V1.md | tar -x -C "${TMP}" + SRC="${TMP}/docs/standards" + SOURCE_LABEL="${CCCC_REPO}@${CCCC_REF}" +else + SRC="${CCCC_REPO%/}/docs/standards" + SOURCE_LABEL="${SRC}" +fi if [[ ! -d "${SRC}" ]]; then echo "error: cannot find CCCC specs at: ${SRC}" >&2 - echo "hint: pass the CCCC repo path explicitly: ./scripts/sync_specs_from_cccc.sh /path/to/cccc" >&2 + echo "hint: ./scripts/sync_specs_from_cccc.sh /path/to/cccc [git-ref]" >&2 exit 2 fi @@ -17,5 +42,4 @@ cp -f "${SRC}/CCCS_V1.md" "${DST}/CCCS_V1.md" cp -f "${SRC}/CCCC_DAEMON_IPC_V1.md" "${DST}/CCCC_DAEMON_IPC_V1.md" cp -f "${SRC}/CCCC_CONTEXT_OPS_V1.md" "${DST}/CCCC_CONTEXT_OPS_V1.md" -echo "Synced specs from ${SRC} -> ${DST}" - +echo "Synced specs from ${SOURCE_LABEL} -> ${DST}" diff --git a/spec/ADAPTATION_PLAN.md b/spec/ADAPTATION_PLAN.md index da0d7f7..85bfe1c 100644 --- a/spec/ADAPTATION_PLAN.md +++ b/spec/ADAPTATION_PLAN.md @@ -4,7 +4,8 @@ This plan is based on an operation-by-operation audit of current CCCC `main` (the v0.4.34 release-candidate line plus subsequent commits), refreshed on 2026-08-11 at core revision `7943724bf1025265d0716b2b181b3890efe24051`. The SDK source tree targets that contract; Python, TypeScript, and Rust package -publication remains a separate release process. +publication remains a separate release process. This audit was revalidated on +2026-08-12 without changing the pinned core revision. The goal is contract alignment, not one wrapper per daemon implementation detail. Public SDK methods should represent stable capabilities that an @@ -57,6 +58,15 @@ TypeScript expose explicit outcome-unknown errors; Python provides the same distinction, including malformed post-write response envelopes, while retaining `DaemonUnavailableError` compatibility. TypeScript cancellation also covers the TCP connection phase instead of beginning only after the stream handshake. +The connection and handshake now share one deadline, buffered stream items stay +as bytes until a complete line is validated, and cleanup removes only listeners +owned by the SDK. + +The 0.4.33 Web Model wait/complete surface remains available on the newer SDK +line. Completion requires a caller-stable `delivery_id`; exact replay behavior +is pinned in `SDK_DAEMON_TARGET_0_4_33.json`. Rust additionally exposes a +least-privilege identity-bound reliable messaging adapter with durable, +monotonic inbox checkpoints and explicit reconciliation of unknown outcomes. ### Current core-side parity blockers @@ -84,17 +94,21 @@ additional SDK-used fields, including actor profile scope/owner, advanced current core handlers and tests, but the authoritative standard should be expanded before they are described as normative cross-implementation v1. -### Validation evidence (2026-08-11) +### Validation evidence (2026-08-12) -- Python: 77 contract/transport tests, source compilation, sdist, and wheel +- Python: 78 contract/transport tests, source compilation, sdist, and wheel build passed. -- TypeScript: 114 tests, strict source and exported-option fixture typechecks, +- TypeScript: 119 tests, strict source and exported-option fixture typechecks, build, and npm package dry-run passed. -- Rust: formatting, warning-free clippy, 12 tests, and locked crate packaging - passed. +- Rust: formatting, warning-free clippy, 21 tests, locked crate packaging, + declared Rust 1.74 MSRV checking, and Windows target checking passed. The + exact 0.4.33 daemon reliability round trip remains an opt-in live test and is + required by CI. - All three mirrored standards match core revision `7943724bf1025265d0716b2b181b3890efe24051` byte-for-byte; scheduled CI now detects future drift. +- An isolated CCCC 0.4.34-rc2 Rust daemon accepted the Python, TypeScript, and + Rust compatibility probes on 2026-08-12. - In the 2026-08-08 live probe, a current Rust daemon bound to the IPv6 loopback wrote a connectable `::1` descriptor, and Python, TypeScript, and Rust clients all discovered it and diff --git a/spec/SDK_DAEMON_TARGET_0_4_33.json b/spec/SDK_DAEMON_TARGET_0_4_33.json new file mode 100644 index 0000000..04e13e4 --- /dev/null +++ b/spec/SDK_DAEMON_TARGET_0_4_33.json @@ -0,0 +1,168 @@ +{ + "schema_version": 1, + "target": { + "product": "cccc", + "version": "0.4.33", + "cargo_install_spec": "cccc@=0.4.33", + "package_sha256": "f715a6947f9fb58fc2a24c21e8c009cf5178e4382a3ccf35d1813f5e7f3f263c", + "daemon_package": "cccc-pair-daemon@=0.4.33", + "daemon_package_sha256": "6c890a08f8cdcdb3a722ce59ef8f5694b2b2587101db4870a446a061f0986270", + "daemon_source_commit": "43be897b75af0eb25a3d4728248bb359f9999204", + "daemon_package_vcs_dirty": true, + "standards_git_ref": "v0.4.33", + "standards_git_commit": "f6c511f1abb37fa49245754646316a4b0bbc0c14" + }, + "daemon_source_sha256": { + "crates/cccc-daemon/src/ops/runtime_state.rs": "d2243bf4bae984bd090ea0a431f5675f792cd0a18bdd2571bce8eb03b93d542f", + "crates/cccc-daemon/src/ops/messaging.rs": "6ba0dacd8c6b53d85c51e51c02aa4de56ff0eb2141b1e90c7c8603d6fb264e01", + "crates/cccc-daemon/src/ops/messaging_inbox.rs": "a2c7afe2591bed741592e2b29ed9446e05ac94c6271ea7d558f7ccf812e46044", + "crates/cccc-daemon/src/ops/messaging_status.rs": "d50c3816f28a372e4fe7ae2699637825704105b91352dd7621c8b8deea2280e1" + }, + "operations": { + "events_stream": { + "advertised_capability": true, + "operation_probe_supported": false, + "probe_error_code": "unknown_op", + "note": "The published native 0.4.33 daemon advertises events_stream but does not dispatch the operation. SDK stream transport remains fixture-tested and assertCompatible fails closed." + }, + "web_model_runtime_complete_turn": { + "required_args": ["group_id", "actor_id", "turn_id", "delivery_id"], + "event_selector_any_of": ["event_ids", "latest_event_id"], + "replay_key": "delivery_id", + "sdk_parameters": { + "python": ["group_id", "actor_id", "turn_id", "delivery_id"], + "typescript": ["groupId", "actorId", "turnId", "deliveryId"] + }, + "request": { + "v": 1, + "op": "web_model_runtime_complete_turn", + "args": { + "group_id": "g1", + "actor_id": "web-1", + "by": "web-1", + "turn_id": "turn-1", + "delivery_id": "runtime:turn-1", + "event_ids": ["e1"], + "status": "done" + } + }, + "missing_delivery_id_response": { + "v": 1, + "ok": false, + "result": {}, + "error": { + "code": "missing_delivery_id", + "message": "delivery_id is required", + "details": {} + } + }, + "completion_response": { + "v": 1, + "ok": true, + "result": { + "status": "done", + "turn_id": "turn-1", + "delivery_id": "runtime:turn-1", + "cursor_committed": true, + "cursor": {"event_id": "e1", "ts": ""}, + "read_event": { + "v": 1, + "id": "receipt-1", + "ts": "2026-08-05T01:00:01Z", + "kind": "chat.read", + "group_id": "g1", + "scope_key": "", + "by": "web-1", + "data": { + "actor_id": "web-1", + "event_id": "e1", + "turn_id": "turn-1", + "event_ids": ["e1"], + "status": "done", + "delivery_id": "runtime:turn-1", + "cursor_committed": true + } + }, + "ack_events": [], + "processed_event_ids": ["e1"], + "followup_delivery_scheduled": false, + "summary": "" + } + } + }, + "inbox_list": { + "cursor_event_id_nullable": true, + "cursor_ts_is_empty": true, + "fresh_response": { + "v": 1, + "ok": true, + "result": { + "messages": [], + "cursor": {"event_id": null, "ts": ""} + } + }, + "remote_ahead_response": { + "v": 1, + "ok": true, + "result": { + "messages": [], + "cursor": {"event_id": "e8", "ts": ""} + } + } + }, + "message_read_status": { + "covered_response": { + "v": 1, + "ok": true, + "result": { + "event_id": "e7", + "read_status": {"a1": true} + } + } + }, + "send": { + "replay_result_field": "duplicate", + "initial_response": { + "v": 1, + "ok": true, + "result": { + "event": { + "v": 1, + "id": "e1", + "ts": "2026-08-05T01:00:00Z", + "kind": "chat.message", + "group_id": "g1", + "scope_key": "", + "by": "workload:aquant", + "data": {"text": "hello", "client_id": "send:stable-1"} + }, + "delivery": {"accepted": true, "state": "queued"} + } + }, + "duplicate_response": { + "v": 1, + "ok": true, + "result": { + "event": { + "v": 1, + "id": "e1", + "ts": "2026-08-05T01:00:00Z", + "kind": "chat.message", + "group_id": "g1", + "scope_key": "", + "by": "workload:aquant", + "data": {"text": "hello", "client_id": "send:stable-1"} + }, + "delivery": { + "accepted": true, + "state": "duplicate", + "targeted": 0, + "online": 0, + "queued": 0 + }, + "duplicate": true + } + } + } + } +} diff --git a/ts/DESIGN.md b/ts/DESIGN.md index b39923c..19ad200 100644 --- a/ts/DESIGN.md +++ b/ts/DESIGN.md @@ -25,6 +25,7 @@ Out of scope: - `src/transport.ts`: endpoint discovery, socket I/O, events stream handshake. - `src/client.ts`: high-level SDK methods. +- `src/client_0430_runtime_ops.ts`: retained Web Model wait/complete operations. - `src/client_0434_ops.ts`: current terminal/Web Model contract additions. - `src/types.ts`: IPC-facing option and payload types. - `src/errors.ts`: typed error hierarchy. @@ -42,6 +43,8 @@ Out of scope: - Connection-establishment failures -> `DaemonConnectionError` and one safe endpoint rediscovery for auto-discovered clients. - Failures after exchange begins -> `OutcomeUnknownError` and no automatic replay. +- Connection/response/handshake phases share one caller deadline; stream bytes + are capped before UTF-8 decoding. - Oversized requests -> `RequestTooLargeError` before connecting. - Daemon `ok:false` responses -> `DaemonAPIError` with `code/message/details/raw`. - Compatibility failures -> `IncompatibleDaemonError`. diff --git a/ts/README.md b/ts/README.md index 2a37bfd..96ebf02 100644 --- a/ts/README.md +++ b/ts/README.md @@ -229,8 +229,27 @@ const recovered = await client.webModelRuntimeRecoverTurn({ actorId: 'web-model', eventIds: ['e_xxx'], }); + +// Complete a runtime-owned turn. Reuse the exact deliveryId after an unknown +// outcome so the daemon can replay the same completion receipt. +const acquired = await client.webModelRuntimeWaitNextTurn({ + groupId, + actorId: 'web-model', +}); +const turn = acquired.turn as { turn_id: string; event_ids: string[] }; +const deliveryId = `worker:${turn.turn_id}`; +await client.webModelRuntimeCompleteTurn({ + groupId, + actorId: 'web-model', + turnId: turn.turn_id, + deliveryId, + eventIds: turn.event_ids, +}); ``` +`deliveryId` is required and is the completion replay key. A retry for the +same acquired turn must reuse the same value. + `termResize()` sends the standard `term_resize` operation. For current Rust daemon builds that still expose `terminal_resize`, the SDK falls back only after receiving a structured `unknown_op`; transport failures are never diff --git a/ts/__tests__/client_0430_contract.test.ts b/ts/__tests__/client_0430_contract.test.ts index 8eafba7..e2a1f81 100644 --- a/ts/__tests__/client_0430_contract.test.ts +++ b/ts/__tests__/client_0430_contract.test.ts @@ -1,10 +1,22 @@ import { describe, it } from 'node:test'; import assert from 'node:assert/strict'; +import { readFileSync } from 'node:fs'; import { CCCCClient } from '../src/client.js'; import { IncompatibleDaemonError } from '../src/errors.js'; type CallCapture = { op: string; args?: Record }; +const targetFixture = JSON.parse( + readFileSync(new URL('../../spec/SDK_DAEMON_TARGET_0_4_33.json', import.meta.url), 'utf-8') +) as { + operations: { + web_model_runtime_complete_turn: { + required_args: string[]; + request: { args: Record }; + }; + }; +}; + async function makeClient(calls: CallCapture[]): Promise { const client = await CCCCClient.create({ endpoint: { transport: 'tcp', host: '127.0.0.1', port: 1, path: '' }, @@ -332,4 +344,29 @@ describe('cccc 0.4.33 JSON op alignment', () => { assert.equal(calls[1]?.args?.['text'], 'Include omissions'); assert.equal(calls[1]?.args?.['source_text'], 'Current meeting notes'); }); + + it('requires and reuses deliveryId for Web Model completion replay', async () => { + const calls: CallCapture[] = []; + const client = await makeClient(calls); + const contract = targetFixture.operations.web_model_runtime_complete_turn; + const args = contract.request.args; + const options = { + groupId: String(args['group_id']), + actorId: String(args['actor_id']), + turnId: String(args['turn_id']), + deliveryId: String(args['delivery_id']), + eventIds: args['event_ids'] as string[], + status: String(args['status']) as 'done', + }; + + await client.webModelRuntimeCompleteTurn(options); + await client.webModelRuntimeCompleteTurn(options); + + assert.deepEqual(calls[0], calls[1]); + assert.equal(calls[0]?.op, 'web_model_runtime_complete_turn'); + for (const required of contract.required_args) { + assert.ok(required in (calls[0]?.args ?? {}), `missing daemon arg: ${required}`); + } + assert.equal(calls[0]?.args?.['delivery_id'], args['delivery_id']); + }); }); diff --git a/ts/__tests__/transport.test.ts b/ts/__tests__/transport.test.ts index caabf04..09e747c 100644 --- a/ts/__tests__/transport.test.ts +++ b/ts/__tests__/transport.test.ts @@ -6,6 +6,7 @@ import * as fs from 'node:fs/promises'; import * as net from 'node:net'; import { Readable } from 'node:stream'; import { + callDaemon, discoverEndpoint, defaultHome, openEventsStream, @@ -13,6 +14,52 @@ import { MAX_LINE_SIZE, DEFAULT_TIMEOUT_MS, } from '../src/transport.js'; +import type { DaemonEndpoint, DaemonRequest } from '../src/types.js'; + +interface TestServer { + endpoint: DaemonEndpoint; + sockets: Set; + close(): Promise; +} + +async function startServer(onConnection: (socket: net.Socket) => void): Promise { + const sockets = new Set(); + const server = net.createServer((socket) => { + sockets.add(socket); + socket.once('close', () => sockets.delete(socket)); + onConnection(socket); + socket.resume(); + }); + await new Promise((resolve, reject) => { + server.once('error', reject); + server.listen(0, '127.0.0.1', resolve); + }); + const address = server.address(); + if (address === null || typeof address === 'string') { + throw new Error('test server did not bind a TCP port'); + } + return { + endpoint: { + transport: 'tcp', + host: '127.0.0.1', + port: address.port, + path: '', + }, + sockets, + close: async () => { + for (const socket of sockets) socket.destroy(); + await new Promise((resolve, reject) => { + server.close((error) => error ? reject(error) : resolve()); + }); + }, + }; +} + +const streamRequest: DaemonRequest = { + v: 1, + op: 'events_stream', + args: { group_id: 'g1', by: 'user' }, +}; describe('defaultHome', () => { it('returns CCCC_HOME env if set', () => { @@ -251,6 +298,61 @@ describe('openEventsStream abort handling', () => { }); }); +describe('daemon response deadlines', () => { + it('times out after TCP connect when the daemon never responds', async () => { + const server = await startServer(() => undefined); + const started = Date.now(); + try { + await assert.rejects( + callDaemon(server.endpoint, { v: 1, op: 'ping', args: {} }, 50), + /Response timeout/, + ); + assert.ok(Date.now() - started < 1_000, 'response timeout should be bounded'); + } finally { + await server.close(); + } + }); +}); + +describe('event stream handshake safety', () => { + it('rejects a daemon that accepts the socket but never handshakes', async () => { + const server = await startServer(() => undefined); + const started = Date.now(); + try { + await assert.rejects( + openEventsStream(server.endpoint, streamRequest, 50), + /Handshake timeout/, + ); + assert.ok(Date.now() - started < 1_000, 'handshake timeout should be bounded'); + } finally { + await server.close(); + } + }); + + it('honors AbortSignal while waiting for the handshake', async () => { + const server = await startServer(() => undefined); + const controller = new AbortController(); + try { + const pending = openEventsStream(server.endpoint, streamRequest, 5_000, controller.signal); + setTimeout(() => controller.abort(), 20); + await assert.rejects(pending, /aborted/); + const closeDeadline = Date.now() + 1_000; + while (server.sockets.size !== 0 && Date.now() < closeDeadline) { + await new Promise((resolve) => setTimeout(resolve, 10)); + } + assert.equal(server.sockets.size, 0, 'aborting the handshake must close the socket'); + } finally { + await server.close(); + } + }); + + it('applies the byte cap to data buffered with the handshake', async () => { + const socket = Readable.from([]) as unknown as net.Socket; + const lines = readLines(socket, Buffer.alloc(MAX_LINE_SIZE + 1, 0x78)); + await assert.rejects(lines.next(), /Stream line exceeds MAX_LINE_SIZE/); + }); +}); + describe('readLines', () => { it('preserves UTF-8 code points split across socket chunks', async () => { const encoded = Buffer.from('{"text":"中文"}\n', 'utf8'); diff --git a/ts/src/client_0430_ops.ts b/ts/src/client_0430_ops.ts index 23f883c..a961278 100644 --- a/ts/src/client_0430_ops.ts +++ b/ts/src/client_0430_ops.ts @@ -1,12 +1,15 @@ import { installCCCC0430AdminOps, type CCCC0430AdminOps } from './client_0430_admin_ops.js'; import { installCCCC0430AssistantOps, type CCCC0430AssistantOps } from './client_0430_assistant_ops.js'; import { installCCCC0430MemoryOps, type CCCC0430MemoryOps } from './client_0430_memory_ops.js'; +import { installCCCC0430RuntimeOps, type CCCC0430RuntimeOps } from './client_0430_runtime_ops.js'; import type { CCCC0430Client } from './client_0430_shared.js'; -export interface CCCC0430Ops extends CCCC0430AdminOps, CCCC0430AssistantOps, CCCC0430MemoryOps {} +export interface CCCC0430Ops + extends CCCC0430AdminOps, CCCC0430AssistantOps, CCCC0430MemoryOps, CCCC0430RuntimeOps {} export function installCCCC0430Ops(proto: CCCC0430Client & Partial): void { installCCCC0430AdminOps(proto); installCCCC0430AssistantOps(proto); installCCCC0430MemoryOps(proto); + installCCCC0430RuntimeOps(proto); } diff --git a/ts/src/client_0430_runtime_ops.ts b/ts/src/client_0430_runtime_ops.ts new file mode 100644 index 0000000..c9e85f2 --- /dev/null +++ b/ts/src/client_0430_runtime_ops.ts @@ -0,0 +1,46 @@ +import { compactRecord, type CCCC0430Client } from './client_0430_shared.js'; +import type { + WebModelRuntimeCompleteTurnOptions, + WebModelRuntimeWaitNextTurnOptions, +} from './types.js'; + +export interface CCCC0430RuntimeOps { + webModelRuntimeWaitNextTurn( + options: WebModelRuntimeWaitNextTurnOptions + ): Promise>; + webModelRuntimeCompleteTurn( + options: WebModelRuntimeCompleteTurnOptions + ): Promise>; +} + +const runtimeOps: CCCC0430RuntimeOps & ThisType = { + async webModelRuntimeWaitNextTurn(options) { + return this.call('web_model_runtime_wait_next_turn', { + group_id: options.groupId, + actor_id: options.actorId, + by: options.by ?? options.actorId, + limit: Math.min(Math.max(Math.trunc(options.limit ?? 20), 1), 20), + kind_filter: options.kindFilter ?? 'all', + }); + }, + + async webModelRuntimeCompleteTurn(options) { + return this.call('web_model_runtime_complete_turn', compactRecord({ + group_id: options.groupId, + actor_id: options.actorId, + by: options.by ?? options.actorId, + turn_id: options.turnId, + delivery_id: options.deliveryId, + event_ids: options.eventIds, + latest_event_id: options.latestEventId, + status: options.status ?? 'done', + summary: options.summary, + })); + }, +}; + +export function installCCCC0430RuntimeOps( + proto: CCCC0430Client & Partial +): void { + Object.assign(proto, runtimeOps); +} diff --git a/ts/src/index.ts b/ts/src/index.ts index 1986916..d5abc2b 100644 --- a/ts/src/index.ts +++ b/ts/src/index.ts @@ -155,6 +155,9 @@ export type { TerminalSinceOptions, TerminalSnapshotOptions, TerminalClearOptions, + WebModelRuntimeWaitNextTurnOptions, + WebModelRuntimeCompletionStatus, + WebModelRuntimeCompleteTurnOptions, WebModelDeliveryMode, WebModelDeliveryPreferencesGetOptions, WebModelDeliveryPreferencesUpdateOptions, diff --git a/ts/src/transport.ts b/ts/src/transport.ts index cdd6b8d..1b9c57d 100644 --- a/ts/src/transport.ts +++ b/ts/src/transport.ts @@ -121,9 +121,12 @@ function connect( return new Promise((resolve, reject) => { const socket = new net.Socket(); let settled = false; + let timer: ReturnType | undefined; const cleanup = () => { - socket.removeAllListeners(); + if (timer !== undefined) clearTimeout(timer); + socket.removeListener('connect', onConnect); + socket.removeListener('error', onError); signal?.removeEventListener('abort', onAbort); }; @@ -154,22 +157,58 @@ function connect( return; } signal?.addEventListener('abort', onAbort, { once: true }); - socket.setTimeout(timeoutMs); + timer = setTimeout(onTimeout, timeoutMs); + socket.once('connect', onConnect); socket.once('error', onError); - socket.once('timeout', onTimeout); - - if (endpoint.transport === 'tcp') { - socket.connect(endpoint.port, endpoint.host, onConnect); - } else if (endpoint.transport === 'unix') { - socket.connect(endpoint.path, onConnect); - } else { - rejectAndDestroy( - new DaemonConnectionError(`Invalid endpoint transport: ${endpoint.transport}`), - ); + + try { + if (endpoint.transport === 'tcp') { + socket.connect(endpoint.port, endpoint.host); + } else if (endpoint.transport === 'unix') { + socket.connect(endpoint.path); + } else { + rejectAndDestroy( + new DaemonConnectionError(`Invalid endpoint transport: ${endpoint.transport}`), + ); + } + } catch (error) { + rejectAndDestroy(new DaemonConnectionError( + error instanceof Error ? error.message : String(error), + )); } }); } +function remainingTimeout(deadline: number): number { + return Math.max(1, deadline - Date.now()); +} + +function appendChunk(buffer: Buffer, chunk: Buffer): Buffer { + return buffer.length === 0 ? chunk : Buffer.concat([buffer, chunk]); +} + +function assertBufferedLineLimit(buffer: Buffer): void { + let start = 0; + let newlineIndex: number; + while ((newlineIndex = buffer.indexOf(0x0a, start)) !== -1) { + if (newlineIndex - start > MAX_LINE_SIZE) { + throw new DaemonUnavailableError( + `Stream line exceeds MAX_LINE_SIZE (${MAX_LINE_SIZE} bytes)`, + ); + } + start = newlineIndex + 1; + } + if (buffer.length - start > MAX_LINE_SIZE) { + throw new DaemonUnavailableError( + `Stream line exceeds MAX_LINE_SIZE (${MAX_LINE_SIZE} bytes)`, + ); + } +} + +function decodeLine(bytes: Buffer): string { + return new TextDecoder('utf-8', { fatal: true }).decode(bytes); +} + // ============================================================ // IPC calls // ============================================================ @@ -189,6 +228,7 @@ export async function callDaemon( request: DaemonRequest, timeoutMs: number = DEFAULT_TIMEOUT_MS ): Promise { + const deadline = Date.now() + timeoutMs; const line = JSON.stringify(request) + '\n'; if (Buffer.byteLength(line, 'utf8') > MAX_REQUEST_SIZE) { throw new RequestTooLargeError(`Daemon request exceeds ${MAX_REQUEST_SIZE} bytes`); @@ -196,19 +236,31 @@ export async function callDaemon( const socket = await connect(endpoint, timeoutMs); return new Promise((resolve, reject) => { - let buffer = Buffer.alloc(0); - let resolved = false; + let buffer: Buffer = Buffer.alloc(0); + let settled = false; + let timer: ReturnType | undefined; const cleanup = () => { - socket.removeAllListeners(); + if (timer !== undefined) clearTimeout(timer); + socket.removeListener('data', onData); + socket.removeListener('error', onError); + socket.removeListener('close', onClose); }; - socket.on('data', (chunk: Buffer) => { - buffer = Buffer.concat([buffer, chunk]); + const fail = (error: Error) => { + if (settled) return; + settled = true; + cleanup(); + socket.destroy(); + reject(error); + }; + const onData = (chunk: Buffer) => { + if (settled) return; + buffer = appendChunk(buffer, chunk); const newlineIndex = buffer.indexOf(0x0a); - if (newlineIndex !== -1 && !resolved) { - resolved = true; + if (newlineIndex !== -1) { + settled = true; const responseBytes = buffer.subarray(0, newlineIndex); cleanup(); socket.destroy(); @@ -238,49 +290,34 @@ export async function callDaemon( } if (newlineIndex === -1 && buffer.length > MAX_LINE_SIZE) { - resolved = true; - cleanup(); - socket.destroy(); - reject(new OutcomeUnknownError(request.op, 'Response too large')); + fail(new OutcomeUnknownError(request.op, 'Response too large')); } - }); - - socket.on('error', (err) => { - if (!resolved) { - resolved = true; - cleanup(); - socket.destroy(); - reject(new OutcomeUnknownError(request.op, err.message)); - } - }); + }; - socket.on('close', () => { - if (!resolved) { - resolved = true; - cleanup(); - socket.destroy(); - reject(new OutcomeUnknownError(request.op, 'Connection closed unexpectedly')); - } - }); + const onError = (err: Error) => fail(new OutcomeUnknownError(request.op, err.message)); + const onClose = () => fail( + new OutcomeUnknownError(request.op, 'Connection closed unexpectedly'), + ); - socket.once('timeout', () => { - if (!resolved) { - resolved = true; - cleanup(); - socket.destroy(); - reject(new OutcomeUnknownError(request.op, 'Response timeout')); - } - }); + socket.on('data', onData); + socket.once('error', onError); + socket.once('close', onClose); + timer = setTimeout( + () => fail(new OutcomeUnknownError(request.op, 'Response timeout')), + remainingTimeout(deadline), + ); // Send request. - socket.write(line, (err) => { - if (err && !resolved) { - resolved = true; - cleanup(); - socket.destroy(); - reject(new OutcomeUnknownError(request.op, `Write failed: ${err.message}`)); - } - }); + try { + socket.write(line, (err) => { + if (err) fail(new OutcomeUnknownError(request.op, `Write failed: ${err.message}`)); + }); + } catch (error) { + fail(new OutcomeUnknownError( + request.op, + error instanceof Error ? error.message : String(error), + )); + } }); } @@ -311,6 +348,7 @@ export async function openEventsStream( timeoutMs: number = DEFAULT_TIMEOUT_MS, signal?: AbortSignal, ): Promise { + const deadline = Date.now() + timeoutMs; if (signal?.aborted) { throw new DaemonUnavailableError('Event stream aborted'); } @@ -330,107 +368,95 @@ export async function openEventsStream( handshake: DaemonResponse; remainingBuffer: Buffer; }>((resolve, reject) => { - let buffer = Buffer.alloc(0); - let resolved = false; + let buffer: Buffer = Buffer.alloc(0); + let settled = false; + let timer: ReturnType | undefined; const cleanup = () => { - socket.removeAllListeners(); + if (timer !== undefined) clearTimeout(timer); + socket.removeListener('data', onData); + socket.removeListener('error', onError); + socket.removeListener('close', onClose); signal?.removeEventListener('abort', onAbort); }; - const onAbort = () => { - if (!resolved) { - resolved = true; - cleanup(); - socket.destroy(); - reject(new DaemonUnavailableError('Event stream aborted')); - } + const fail = (error: Error) => { + if (settled) return; + settled = true; + cleanup(); + socket.destroy(); + reject(error); }; + const onAbort = () => fail(new DaemonUnavailableError('Event stream aborted')); + const onData = (chunk: Buffer) => { - buffer = Buffer.concat([buffer, chunk]); + if (settled) return; + buffer = appendChunk(buffer, chunk); const newlineIndex = buffer.indexOf(0x0a); - if (newlineIndex !== -1 && !resolved) { - resolved = true; - cleanup(); + if (newlineIndex !== -1) { const responseBytes = buffer.subarray(0, newlineIndex); const remaining = buffer.subarray(newlineIndex + 1); if (responseBytes.length > MAX_LINE_SIZE) { - socket.destroy(); - reject(new OutcomeUnknownError(request.op, 'Handshake response too large')); + fail(new OutcomeUnknownError(request.op, 'Handshake response too large')); return; } try { const parsed: unknown = JSON.parse(responseBytes.toString('utf8')); if (parsed === null || typeof parsed !== 'object' || Array.isArray(parsed)) { - socket.destroy(); - reject(new OutcomeUnknownError(request.op, 'Handshake must be a JSON object')); + fail(new OutcomeUnknownError(request.op, 'Handshake must be a JSON object')); return; } const handshake = parsed as DaemonResponse; if (handshake.v !== 1) { - socket.destroy(); - reject(new IncompatibleDaemonError( + fail(new IncompatibleDaemonError( `Daemon stream handshake uses unsupported IPC version: ${String(handshake.v)}`, )); return; } + assertBufferedLineLimit(remaining); + settled = true; + cleanup(); resolve({ handshake, remainingBuffer: remaining, }); - } catch { - socket.destroy(); - reject(new OutcomeUnknownError(request.op, 'Invalid handshake JSON')); + } catch (error) { + fail(error instanceof DaemonUnavailableError + ? error + : new OutcomeUnknownError(request.op, 'Invalid handshake JSON')); } } if (newlineIndex === -1 && buffer.length > MAX_LINE_SIZE) { - resolved = true; - cleanup(); - socket.destroy(); - reject(new OutcomeUnknownError(request.op, 'Handshake response too large')); + fail(new OutcomeUnknownError(request.op, 'Handshake response too large')); } }; + const onError = (err: Error) => fail(new OutcomeUnknownError(request.op, err.message)); + const onClose = () => fail( + new OutcomeUnknownError(request.op, 'Connection closed during handshake'), + ); + socket.on('data', onData); signal?.addEventListener('abort', onAbort, { once: true }); - socket.once('error', (err) => { - if (!resolved) { - resolved = true; - cleanup(); - socket.destroy(); - reject(new OutcomeUnknownError(request.op, err.message)); - } - }); - socket.once('close', () => { - if (!resolved) { - resolved = true; - cleanup(); - socket.destroy(); - reject(new OutcomeUnknownError(request.op, 'Connection closed during handshake')); - } - }); - socket.once('timeout', () => { - if (!resolved) { - resolved = true; - cleanup(); - socket.destroy(); - reject(new OutcomeUnknownError(request.op, 'Handshake timeout')); - } - }); - socket.write(line, (error) => { - if (error && !resolved) { - resolved = true; - cleanup(); - socket.destroy(); - reject(new OutcomeUnknownError(request.op, `Write failed: ${error.message}`)); - } - }); + socket.once('error', onError); + socket.once('close', onClose); + timer = setTimeout( + () => fail(new OutcomeUnknownError(request.op, 'Handshake timeout')), + remainingTimeout(deadline), + ); + try { + socket.write(line, (error) => { + if (error) fail(new OutcomeUnknownError(request.op, `Write failed: ${error.message}`)); + }); + } catch (error) { + fail(new OutcomeUnknownError( + request.op, + error instanceof Error ? error.message : String(error), + )); + } }); - // Remove timeout after handshake. - socket.setTimeout(0); - return { socket, handshake, initialBuffer: remainingBuffer }; } @@ -445,53 +471,49 @@ export async function* readLines( socket: net.Socket, initialBuffer: string | Buffer = '' ): AsyncGenerator { - const decoder = new TextDecoder('utf-8', { fatal: true }); let buffer = typeof initialBuffer === 'string' - ? initialBuffer - : decoder.decode(initialBuffer, { stream: true }); + ? Buffer.from(initialBuffer, 'utf8') + : initialBuffer; - const ensureBounded = (line: string): void => { - if (Buffer.byteLength(line, 'utf8') > MAX_LINE_SIZE) { + // Handle lines from initial buffer. + let newlineIndex: number; + while ((newlineIndex = buffer.indexOf(0x0a)) !== -1) { + if (newlineIndex > MAX_LINE_SIZE) { throw new DaemonUnavailableError( `Stream line exceeds MAX_LINE_SIZE (${MAX_LINE_SIZE} bytes)`, ); } - }; - - // Handle lines from initial buffer. - let newlineIndex: number; - while ((newlineIndex = buffer.indexOf('\n')) !== -1) { - const line = buffer.slice(0, newlineIndex); - buffer = buffer.slice(newlineIndex + 1); - ensureBounded(line); + const line = decodeLine(buffer.subarray(0, newlineIndex)); + buffer = buffer.subarray(newlineIndex + 1); if (line.trim()) { yield line; } } + assertBufferedLineLimit(buffer); // Continue reading from socket. for await (const chunk of socket) { - buffer += decoder.decode(chunk as Buffer, { stream: true }); - - if (buffer.indexOf('\n') === -1 && Buffer.byteLength(buffer, 'utf8') > MAX_LINE_SIZE) { - throw new DaemonUnavailableError( - `Stream line exceeds MAX_LINE_SIZE (${MAX_LINE_SIZE} bytes)` - ); - } - - while ((newlineIndex = buffer.indexOf('\n')) !== -1) { - const line = buffer.slice(0, newlineIndex); - buffer = buffer.slice(newlineIndex + 1); - ensureBounded(line); + const bytes = Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk as Uint8Array); + buffer = appendChunk(buffer, bytes); + + while ((newlineIndex = buffer.indexOf(0x0a)) !== -1) { + if (newlineIndex > MAX_LINE_SIZE) { + throw new DaemonUnavailableError( + `Stream line exceeds MAX_LINE_SIZE (${MAX_LINE_SIZE} bytes)`, + ); + } + const line = decodeLine(buffer.subarray(0, newlineIndex)); + buffer = buffer.subarray(newlineIndex + 1); if (line.trim()) { yield line; } } + assertBufferedLineLimit(buffer); } - buffer += decoder.decode(); - if (buffer.trim()) { - ensureBounded(buffer); - yield buffer; + if (buffer.length > 0) { + assertBufferedLineLimit(buffer); + const line = decodeLine(buffer); + if (line.trim()) yield line; } } diff --git a/ts/src/types.ts b/ts/src/types.ts index a783701..b58732f 100644 --- a/ts/src/types.ts +++ b/ts/src/types.ts @@ -697,6 +697,34 @@ export interface TerminalSnapshotOptions { by?: string; } +/** Read the next coalesced turn for one Web Model runtime actor. */ +export interface WebModelRuntimeWaitNextTurnOptions { + groupId: string; + actorId: string; + by?: string; + limit?: number; + kindFilter?: 'all' | 'chat' | 'notify'; +} + +export type WebModelRuntimeCompletionStatus = + | 'done' + | 'partial' + | 'failed' + | 'cancelled'; + +/** Complete a previously acquired Web Model turn with a stable replay key. */ +export interface WebModelRuntimeCompleteTurnOptions { + groupId: string; + actorId: string; + turnId: string; + deliveryId: string; + eventIds?: string[]; + latestEventId?: string; + status?: WebModelRuntimeCompletionStatus; + summary?: string; + by?: string; +} + export type WebModelDeliveryMode = 'standard' | 'image_compat'; export interface WebModelDeliveryPreferencesGetOptions {