From 6e36f43725632b7265c63feda565adfc3a02cef9 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Wed, 30 Sep 2026 01:48:08 +0800 Subject: [PATCH] perf(runtime): batch exact source fingerprint reads Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../2026-09-28-retirement-cadence.md | 20 +++++ .../2026-09-28-retirement-cadence.zh-CN.md | 15 ++++ loopx/contract.py | 2 +- loopx/control_plane/effect_runtime.py | 16 +++- .../{file_text_reads.py => file_reads.py} | 50 ++++++++++--- .../test_public_boundary_parallel_reads.py | 22 ++++-- .../test_runtime_source_read_batching.py | 75 +++++++++++++++++++ 7 files changed, 178 insertions(+), 22 deletions(-) rename loopx/control_plane/runtime/{file_text_reads.py => file_reads.py} (58%) create mode 100644 tests/control_plane/test_runtime_source_read_batching.py diff --git a/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-28-retirement-cadence.md b/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-28-retirement-cadence.md index 429200431b..350d5f9ae4 100644 --- a/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-28-retirement-cadence.md +++ b/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-28-retirement-cadence.md @@ -311,3 +311,23 @@ metadata preservation and unchanged provider state. This is a bounded capacity repair, not unlimited graph capacity, stable latency evidence, D2 qualification or permission to change the default provider. Full-Goal summary/list/detail adoption and sustained observation remain separate work. + +### Packaged-source fingerprint cost + +The next B-lane increment overlaps source-byte reads through the existing bounded +ordered file reader. It preserves every relative name and raw byte in the digest, +metadata invalidation, request-scoped reuse, and failure/retry behavior; the +Python adapter gains no state-policy owner or persistent cache. On macOS arm64, +Python 3.13.13 and Node 24.21.0, alternating nine fresh CLI processes per arm +(after startup warm-up) against detached File/SQLite copies measured status +medians of 1.985→1.874 s and 2.106→1.937 s. The same source directory contained +246 TS/JSON files (3,003,024 bytes); fingerprint-stage medians were 249.7→51.7 ms +and 228.3→55.8 ms. Effect processes and data were isolated; OS caches were not +flushed. Complete responses differed only at enumerated observation-time paths. + +This is a bounded caller-cost improvement, not provider throughput or D2/default +qualification. A repeated same-process microbenchmark with fingerprint memoization +explicitly cleared regressed from 9.3 to 17.7 ms; normal unchanged requests retain +memoization. Do not extrapolate either workload to every platform or an installed +fleet. Whole-Goal payload/consumer work and sustained operation remain open; this +increment authorizes no legacy-writer deletion or UI data truncation. diff --git a/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-28-retirement-cadence.zh-CN.md b/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-28-retirement-cadence.zh-CN.md index 8158b39a6d..58160a27d6 100644 --- a/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-28-retirement-cadence.zh-CN.md +++ b/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-28-retirement-cadence.zh-CN.md @@ -233,3 +233,18 @@ Todo/租约集合中位数为 File 输入 32.4→27.5 ms、SQLite 输入 32.9 真实 File/SQLite CLI 验证精确计数、跨远端记录的推断后继、metadata 保留与 provider 状态不变。这只修复有界容量,不证明任意规模、稳定延迟、D2 验收, 也不授予默认 provider 切换;整 Goal summary/list/detail 消费和持续观察仍待推进。 + +### 打包源码指纹成本 + +B 阶段下一增量复用已有的有界、有序文件读取器,并发读取源码字节。摘要仍包含 +全部相对文件名和原始字节,保留 metadata 失效、请求内复用及失败/重试行为; +Python 适配层不新增状态规则或持久缓存。macOS arm64、Python 3.13.13、Node 24.21.0 +上,预热后每组交替运行九个独立 CLI 进程,对隔离 File/SQLite 副本测得 status +中位数为 1.985→1.874 秒、2.106→1.937 秒。同一源码目录包含 246 个 TS/JSON 文件 +(3,003,024 字节);指纹阶段中位数为 249.7→51.7 毫秒、228.3→55.8 毫秒。 +数据与 Effect 进程均隔离,没有清空 OS 缓存;完整响应仅列明的观测时间字段不同。 + +这是有界的调用方成本改善,不是 provider 吞吐或 D2/默认项验收。显式清除指纹 +缓存后,同进程反复计算的微基准从 9.3 退化到 17.7 毫秒;正常未变更请求仍复用 +缓存。两种负载都不能外推为所有平台或已安装用户的结果。整 Goal 大包/消费者与 +持续运行仍待推进,本增量不授权删除旧 writer 或裁剪 UI 数据。 diff --git a/loopx/contract.py b/loopx/contract.py index 7c214833c5..fe92f694fa 100644 --- a/loopx/contract.py +++ b/loopx/contract.py @@ -29,7 +29,7 @@ classify_index_duplicate_records, index_identity, ) -from .control_plane.runtime.file_text_reads import iter_utf8_file_reads +from .control_plane.runtime.file_reads import iter_utf8_file_reads from .control_plane.todos.active_state_editing import COMPLETED_WORK_ARCHIVE_HEADING from .control_plane.todos.authoring_scope import todo_contract_diagnostics from .history import ( diff --git a/loopx/control_plane/effect_runtime.py b/loopx/control_plane/effect_runtime.py index 2c256983c9..7a7d55d73a 100644 --- a/loopx/control_plane/effect_runtime.py +++ b/loopx/control_plane/effect_runtime.py @@ -12,7 +12,7 @@ import time import uuid from collections.abc import Iterator, Mapping -from contextlib import ExitStack, contextmanager +from contextlib import ExitStack, closing, contextmanager from contextvars import ContextVar from dataclasses import dataclass from functools import lru_cache @@ -21,6 +21,7 @@ from typing import IO, Any from ..file_lock import process_is_alive +from .runtime.file_reads import iter_binary_file_reads from .content_digest import BARE_SHA256_PATTERN EFFECT_RUNTIME_REQUEST_SCHEMA_VERSION = "loopx_effect_runtime_request_v0" @@ -279,9 +280,16 @@ def _runtime_fingerprint_for_snapshot( ) -> str: digest = hashlib.sha256() source_root = Path(root) - for relative, *_metadata in snapshot: - digest.update(relative.encode("utf-8")) - digest.update((source_root / relative).read_bytes()) + paths = (source_root / relative for relative, *_metadata in snapshot) + # Reads may finish out of order; hash the same relative names and original + # bytes in snapshot order. No disk cache or skipped freshness check. + with closing(iter_binary_file_reads(paths)) as reads: + for (relative, *_metadata), read in zip(snapshot, reads, strict=True): + if read.error is not None: + raise read.error + assert read.data is not None + digest.update(relative.encode("utf-8")) + digest.update(read.data) return digest.hexdigest() diff --git a/loopx/control_plane/runtime/file_text_reads.py b/loopx/control_plane/runtime/file_reads.py similarity index 58% rename from loopx/control_plane/runtime/file_text_reads.py rename to loopx/control_plane/runtime/file_reads.py index bf2f9c65e8..304a2150ce 100644 --- a/loopx/control_plane/runtime/file_text_reads.py +++ b/loopx/control_plane/runtime/file_reads.py @@ -1,4 +1,4 @@ -"""Bounded ordered UTF-8 reads; classification stays with the caller. +"""Bounded ordered file reads; classification stays with the caller. This is a filesystem adapter, not an authorization or scan-result cache. The caller must exclude private inputs before submitting them. Every admitted path @@ -8,11 +8,12 @@ from __future__ import annotations from collections import deque -from collections.abc import Iterable, Iterator +from collections.abc import Callable, Generator, Iterable from concurrent.futures import Future, ThreadPoolExecutor from dataclasses import dataclass from itertools import chain, islice from pathlib import Path +from typing import TypeVar @dataclass(frozen=True) @@ -29,9 +30,40 @@ def _read_utf8(path: Path) -> Utf8FileRead: return Utf8FileRead(path, None, error) +@dataclass(frozen=True) +class BinaryFileRead: + path: Path + data: bytes | None + error: OSError | None + + +def _read_bytes(path: Path) -> BinaryFileRead: + try: + return BinaryFileRead(path, path.read_bytes(), None) + except OSError as error: + return BinaryFileRead(path, None, error) + + +_Read = TypeVar("_Read") + + +def iter_binary_file_reads( + paths: Iterable[Path], *, max_workers: int = 8 +) -> Generator[BinaryFileRead, None, None]: + """Read original bytes without decoding or newline normalization.""" + yield from _ordered_reads(paths, _read_bytes, max_workers) + + def iter_utf8_file_reads( paths: Iterable[Path], *, max_workers: int = 8 -) -> Iterator[Utf8FileRead]: +) -> Generator[Utf8FileRead, None, None]: + """Read UTF-8 text with the same ordered, bounded filesystem lifetime.""" + yield from _ordered_reads(paths, _read_utf8, max_workers) + + +def _ordered_reads( + paths: Iterable[Path], read: Callable[[Path], _Read], max_workers: int +) -> Generator[_Read, None, None]: """Overlap disk waits with at most ``max_workers`` pending reads. Results retain input order, including failures. Unlike ``Executor.map`` on @@ -47,14 +79,14 @@ def iter_utf8_file_reads( return if max_workers == 1 or len(first_paths) == 1: for path in chain(first_paths, iterator): - yield _read_utf8(path) + yield read(path) return with ThreadPoolExecutor(max_workers=max_workers, thread_name_prefix="loopx-file-read") as pool: - pending: deque[Future[Utf8FileRead]] = deque( - pool.submit(_read_utf8, path) for path in first_paths + pending: deque[Future[_Read]] = deque( + pool.submit(read, path) for path in first_paths ) while pending: yield pending.popleft().result() - path = next(iterator, None) - if path is not None: - pending.append(pool.submit(_read_utf8, path)) + next_path = next(iterator, None) + if next_path is not None: + pending.append(pool.submit(read, next_path)) diff --git a/tests/control_plane/test_public_boundary_parallel_reads.py b/tests/control_plane/test_public_boundary_parallel_reads.py index b5042097dc..8f2c05fc06 100644 --- a/tests/control_plane/test_public_boundary_parallel_reads.py +++ b/tests/control_plane/test_public_boundary_parallel_reads.py @@ -7,16 +7,18 @@ import pytest from loopx import contract -from loopx.control_plane.runtime.file_text_reads import iter_utf8_file_reads +from loopx.control_plane.runtime.file_reads import iter_binary_file_reads, iter_utf8_file_reads +@pytest.mark.parametrize("binary", [False, True]) def test_reads_overlap_but_results_remain_in_input_order( - tmp_path: Path, monkeypatch: pytest.MonkeyPatch + tmp_path: Path, monkeypatch: pytest.MonkeyPatch, binary: bool ) -> None: paths = [tmp_path / f"{index}.md" for index in range(4)] for index, path in enumerate(paths): path.write_text(str(index), encoding="utf-8") - original = Path.read_text + original = Path.read_bytes if binary else Path.read_text + reader = iter_binary_file_reads if binary else iter_utf8_file_reads release_first, other_completed = Event(), Event() calls: list[Path] = [] lock = Lock() @@ -31,9 +33,9 @@ def read(path: Path, *args, **kwargs) -> str: other_completed.set() return result - monkeypatch.setattr(Path, "read_text", read) + monkeypatch.setattr(Path, "read_bytes" if binary else "read_text", read) with ThreadPoolExecutor(max_workers=1) as consumer: - result = consumer.submit(lambda: list(iter_utf8_file_reads(paths, max_workers=2))) + result = consumer.submit(lambda: list(reader(paths, max_workers=2))) try: assert other_completed.wait(5), "I/O is still serial" assert not result.done() @@ -41,13 +43,17 @@ def read(path: Path, *args, **kwargs) -> str: release_first.set() reads = result.result(timeout=5) assert [item.path for item in reads] == paths - assert [item.text for item in reads] == ["0", "1", "2", "3"] + if binary: + assert [item.data for item in reads] == [b"0", b"1", b"2", b"3"] + else: + assert [item.text for item in reads] == ["0", "1", "2", "3"] assert all(item.error is None for item in reads) assert sorted(calls) == paths # each path opened exactly once +@pytest.mark.parametrize("reader", [iter_utf8_file_reads, iter_binary_file_reads]) @pytest.mark.parametrize("workers", [1, 3, 8]) -def test_input_consumption_and_pending_results_are_bounded(tmp_path: Path, workers: int) -> None: +def test_input_consumption_and_pending_results_are_bounded(tmp_path: Path, workers: int, reader) -> None: paths = [tmp_path / f"{index:02}.md" for index in range(20)] for path in paths: path.write_text("public", encoding="utf-8") @@ -58,7 +64,7 @@ def inputs(): submitted.append(path) yield path - reads = iter_utf8_file_reads(inputs(), max_workers=workers) + reads = reader(inputs(), max_workers=workers) try: assert next(reads).path == paths[0] assert submitted == paths[:workers] diff --git a/tests/control_plane/test_runtime_source_read_batching.py b/tests/control_plane/test_runtime_source_read_batching.py new file mode 100644 index 0000000000..1cd1f167c6 --- /dev/null +++ b/tests/control_plane/test_runtime_source_read_batching.py @@ -0,0 +1,75 @@ +from __future__ import annotations + +import hashlib +from pathlib import Path + +import pytest + +from loopx.control_plane import effect_runtime as runtime + + +def serial_fingerprint(root: Path) -> str: + """Independent reference: names and raw bytes, not decoded source text.""" + digest = hashlib.sha256() + for path in sorted(p for p in root.rglob("*") if p.suffix in {".ts", ".json"}): + digest.update(path.relative_to(root).as_posix().encode("utf-8")) + digest.update(path.read_bytes()) + return digest.hexdigest() + + +def test_fingerprint_preserves_exact_bytes_and_observes_source_changes(tmp_path, monkeypatch): + monkeypatch.setattr(runtime, "_control_plane_root", lambda: tmp_path) + files = {"a.ts": b"// line\r\n", "nested/b.json": '{"text":"中文🙂"}'.encode(), + "z.ts": b"// bytes\n\xff", "ignored.py": b"not runtime source"} + for name, data in files.items(): + path = tmp_path / name + path.parent.mkdir(exist_ok=True) + path.write_bytes(data) + before = runtime._runtime_fingerprint() + assert before == serial_fingerprint(tmp_path) + (tmp_path / "a.ts").write_bytes(b"// line\n") + modified = runtime._runtime_fingerprint() + assert modified == serial_fingerprint(tmp_path) and modified != before + (tmp_path / "new.ts").write_bytes(b"export {};\n") + added = runtime._runtime_fingerprint() + assert added == serial_fingerprint(tmp_path) and added != modified + (tmp_path / "nested/b.json").unlink() + removed = runtime._runtime_fingerprint() + assert removed == serial_fingerprint(tmp_path) and removed != added + (tmp_path / "ignored.py").write_bytes(b"still not runtime source") + assert runtime._runtime_fingerprint() == removed + + +@pytest.mark.parametrize("error", [PermissionError("denied"), RuntimeError("unexpected")]) +def test_source_read_failure_is_not_a_partial_fingerprint(tmp_path, monkeypatch, error): + monkeypatch.setattr(runtime, "_control_plane_root", lambda: tmp_path) + for name in ("a.ts", "b.ts", "c.ts"): + (tmp_path / name).write_bytes(b"export {};\n") + original = Path.read_bytes + + def read(path): + if path.name == "b.ts": + raise error + return original(path) + + monkeypatch.setattr(Path, "read_bytes", read) + with pytest.raises(type(error), match=str(error)): + runtime._runtime_fingerprint() + monkeypatch.setattr(Path, "read_bytes", original) + assert runtime._runtime_fingerprint() == serial_fingerprint(tmp_path) + + +def test_deleted_source_retries_the_new_topology(tmp_path, monkeypatch): + monkeypatch.setattr(runtime, "_control_plane_root", lambda: tmp_path) + removed = tmp_path / "a.ts" + removed.write_bytes(b"before") + (tmp_path / "b.json").write_bytes(b"{}") + original = Path.read_bytes + + def read(path): + if path == removed and path.exists(): + path.unlink() + return original(path) + + monkeypatch.setattr(Path, "read_bytes", read) + assert runtime._runtime_fingerprint() == serial_fingerprint(tmp_path)