Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -327,3 +327,29 @@ 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 B-lane increment overlaps source-byte reads through the existing bounded,
ordered file reader. It preserves relative names, raw bytes, metadata
invalidation, request-scoped memoization and failure/retry behavior. The Python
filesystem adapter gains no state-policy owner or persistent cache.

Current validation compared baseline `0538bf1631a7` with this implementation on
macOS arm64, Python 3.13.13 and Node 24.21.0. Each arm ran nine alternating fresh
CLI processes after one startup warm-up, against the same disposable synthetic
File/SQLite fixtures. Effect processes were isolated; OS caches were not flushed.
The source snapshot contained 249 TS/JSON files (3,065,799 bytes). Fingerprint
stage medians were 123.5→53.8 ms for File and 111.3→48.0 ms for SQLite.
Whole `status` medians were 1.032→1.054 s and 1.019→1.010 s; sampled p95 values
were 1.745→1.104 s and 1.114→1.086 s (with nine samples, p95 is the maximum).
Twenty full-response pairs differed only at explicitly enumerated observation
timestamps; malformed-registry rejection was unchanged.

This supports a bounded cold-caller cost improvement, not a general status
speedup, provider throughput or D2/default qualification. A warm same-process
microbenchmark with fingerprint memoization explicitly cleared regressed from
6.7 to 12.8 ms; normal unchanged requests retain memoization. Thread scheduling
costs more when all bytes are already hot. Neither workload establishes a fleet
latency guarantee. Whole-Goal payload/consumer work and sustained operation
remain open; this increment authorizes no legacy-writer deletion or UI truncation.
Original file line number Diff line number Diff line change
Expand Up @@ -243,3 +243,23 @@ Todo/租约集合中位数为 File 输入 32.4→27.5 ms、SQLite 输入 32.9
真实 File/SQLite CLI 验证精确计数、跨远端记录的推断后继、metadata 保留与
provider 状态不变。这只修复有界容量,不证明任意规模、稳定延迟、D2 验收,
也不授予默认 provider 切换;整 Goal summary/list/detail 消费和持续观察仍待推进。

### 打包源码指纹成本

B 阶段增量复用已有的有界、有序文件读取器,并发读取源码字节。摘要保留全部
相对文件名和原始字节、metadata 失效、请求内缓存及失败/重试行为;Python
文件系统适配层不新增状态规则或持久缓存。

当前复核对比基线 `0538bf1631a7` 与本实现,使用 macOS arm64、Python 3.13.13、
Node 24.21.0。启动预热一次后,每组交替运行九个独立 CLI 进程,读取相同的可丢弃
合成 File/SQLite fixture;Effect 进程隔离,没有清空 OS 缓存。源码快照包含
249 个 TS/JSON 文件(3,065,799 字节)。指纹阶段中位数分别为 123.5→53.8 毫秒、
111.3→48.0 毫秒。完整 `status` 中位数为 1.032→1.054 秒、1.019→1.010 秒;
样本 p95 为 1.745→1.104 秒、1.114→1.086 秒(九个样本的 p95 就是最大值)。
二十组完整响应仅明确列出的观测时间字段不同,畸形 registry 的拒绝结果不变。

这支持有界的冷调用成本改善,不证明通用 status 提速、provider 吞吐或 D2/默认项
验收。显式清除指纹缓存后,同进程热缓存微基准从 6.7 退化到 12.8 毫秒;正常未变更
请求仍复用缓存。字节已经在热缓存中时,线程调度成本更高。两种负载都不能外推为
用户群体的延迟保证。整 Goal 大包/消费者与持续运行仍待推进,本增量不授权删除
旧 writer 或裁剪 UI 数据。
2 changes: 1 addition & 1 deletion loopx/contract.py
Original file line number Diff line number Diff line change
Expand Up @@ -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 (
Expand Down
16 changes: 12 additions & 4 deletions loopx/control_plane/effect_runtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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"
Expand Down Expand Up @@ -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()


Expand Down
Original file line number Diff line number Diff line change
@@ -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
Expand All @@ -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)
Expand All @@ -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
Expand All @@ -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))
22 changes: 14 additions & 8 deletions tests/control_plane/test_public_boundary_parallel_reads.py
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand All @@ -31,23 +33,27 @@ 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()
finally:
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")
Expand All @@ -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]
Expand Down
75 changes: 75 additions & 0 deletions tests/control_plane/test_runtime_source_read_batching.py
Original file line number Diff line number Diff line change
@@ -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)
Loading