From 4edb21c852626bc50d18e47a4157835a8c2d77b3 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 27 Sep 2026 13:11:30 +0800 Subject: [PATCH 1/5] refactor(authority): own bounded shadow drain in TypeScript Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../local_authority_shadow_adapter.py | 522 +----------------- .../coordination/shadow_drain.ts | 175 ++++++ .../coordination/shadow_drain_files.ts | 152 +++++ .../coordination/shadow_drain_plan.ts | 4 +- .../coordination/shadow_lock_host.py | 39 ++ .../control_plane/effect_runtime_handlers.ts | 4 +- loopx/control_plane/effect_runtime_io.ts | 35 +- .../testing/authority_e2e_rows_stage2c2.py | 58 +- .../testing/shadow_drain_fault_process.ts | 13 + tests/control_plane/shadow_e2e_fixture.py | 58 +- .../test_local_authority_shadow_drain.py | 162 +----- .../test_local_authority_shadow_outbox.py | 10 - .../test_shadow_drain_adversarial.py | 6 +- .../test_shadow_drain_native_plan.py | 139 ++--- .../effect_runtime_io.test.ts | 12 + tests/control_plane_ts/shadow_drain.test.ts | 119 ++++ 16 files changed, 648 insertions(+), 860 deletions(-) create mode 100644 loopx/control_plane/coordination/shadow_drain.ts create mode 100644 loopx/control_plane/coordination/shadow_drain_files.ts create mode 100644 loopx/control_plane/coordination/shadow_lock_host.py create mode 100644 loopx/control_plane/testing/shadow_drain_fault_process.ts create mode 100644 tests/control_plane_ts/shadow_drain.test.ts diff --git a/loopx/control_plane/coordination/local_authority_shadow_adapter.py b/loopx/control_plane/coordination/local_authority_shadow_adapter.py index 5ace9dc6d1..52b10976d7 100644 --- a/loopx/control_plane/coordination/local_authority_shadow_adapter.py +++ b/loopx/control_plane/coordination/local_authority_shadow_adapter.py @@ -6,38 +6,29 @@ from __future__ import annotations -import time -from collections.abc import Iterator, Mapping -from contextlib import contextmanager -from dataclasses import asdict, dataclass, field +import sys +from collections.abc import Mapping +from dataclasses import asdict, dataclass, field, fields from pathlib import Path from typing import Any -from ...file_lock import ( - LockAcquireTimeoutError, - exclusive_cross_runtime_file_lock, - try_exclusive_file_lock, -) -from ...history import load_registry +from ..projects.registry_codec import load_registry from ...paths import resolve_runtime_root from ...registry import find_registry_goal from ..effect_runtime import effect_runtime_result from . import local_authority_shadow_outbox as outbox from .coordination_state_contract_generated import ( - LOCAL_AUTHORITY_SHADOW_COMMIT_ENTRY_REQUEST_SCHEMA, LOCAL_AUTHORITY_SHADOW_COMMIT_ENTRY_RESULT_SCHEMA, LOCAL_AUTHORITY_SHADOW_READ_REQUEST_SCHEMA, LOCAL_AUTHORITY_SHADOW_READ_RESULT_SCHEMA, LOCAL_AUTHORITY_SHADOW_TRANSACTION_EVIDENCE_SCHEMA, ) from .local_authority_shadow_projection import ( - PARTITIONS, - TODO_PARTITION, head_digest, todo_partition_projection, ) from .runtime_shadow import resolve_coordination_runtime_shadow_config, capture_todo_archive_dependencies -from .shadow_management import read_shadow_capture_binding +from .shadow_management import read_shadow_capture_binding, shadow_management_state_path from .runtime_shadow import local_authority_shadow_summary @@ -79,17 +70,6 @@ def effective_runtime_root( INLINE_DRAIN_LOCK_TIMEOUT_SECONDS = 0.25 CLI_DRAIN_LOCK_TIMEOUT_SECONDS = 5.0 RETENTION_PRESSURE_BYTES = 8 * 1024 * 1024 -_COMMIT_ENTRY_OUTCOMES = { - "delivered", - "replayed", - "ambiguous_reconciled", - "ambiguous_unproved", - "unavailable", - "failed", - "protocol_mismatch", - "conflict_retry_required", -} -_SETTLED_OUTCOMES = {"delivered", "replayed", "ambiguous_reconciled"} _SEED_WRITE_CLASSES = {"seed", "reseed_after_crash_gap"} _EVIDENCE_V1_OUTCOMES = { "delivered", @@ -188,121 +168,6 @@ def project(state_text: str) -> dict[str, Any]: return project -def primary_lock_is_free(target: Path) -> bool: - """Probe a partition's Python primary lock once without waiting.""" - - try: - with try_exclusive_file_lock( - target, operation="local_authority_shadow_drain_probe" - ) as held: - return held is not None - except OSError: - return False - - -@dataclass(frozen=True) -class _GoalSources: - goal: dict[str, Any] | None - state_path: Path - lease_dir: Path - - -def _goal_sources( - registry: dict[str, Any], - *, - runtime_root: Path, - goal_id: str, -) -> _GoalSources: - from ...state_refresh import resolve_goal_state - - goal, _project, state_path = resolve_goal_state( - registry=registry, - goal_id=goal_id, - project_override=None, - state_file_override=None, - ) - return _GoalSources( - goal=goal, - state_path=state_path, - lease_dir=outbox.lease_directory(runtime_root, goal_id), - ) - - -@contextmanager -def _primary_lock_if_free( - partition: str, - *, - runtime_root: Path, - goal_id: str, - sources: _GoalSources, -) -> Iterator[bool]: - """Hold the partition's primary lock only if it is free right now.""" - - if partition == TODO_PARTITION: - from .legacy_writer_fence import legacy_coordination_todo_lock_path - - try: - with ( - exclusive_cross_runtime_file_lock( - legacy_coordination_todo_lock_path( - runtime_root=runtime_root, goal_id=goal_id - ), - timeout_seconds=0.0, - operation="local_authority_shadow_drain_resolve", - ), - exclusive_cross_runtime_file_lock( - sources.state_path, - timeout_seconds=0.0, - operation="local_authority_shadow_drain_resolve", - ), - ): - yield True - except LockAcquireTimeoutError: - yield False - return - from ..work_items.task_lease import task_lease_lock_path - - target = task_lease_lock_path(runtime_root=runtime_root, goal_id=goal_id) - try: - with exclusive_cross_runtime_file_lock( - target, - timeout_seconds=0.0, - operation="local_authority_shadow_drain_resolve", - ): - yield True - except LockAcquireTimeoutError: - yield False - - -def _commit_entry_request( - *, runtime_root: Path, goal_id: str, entry: outbox.OutboxEntry, -) -> dict[str, Any]: - """Select witnessed disk evidence; TS owns projection and source resolution.""" - return { - "schema_version": LOCAL_AUTHORITY_SHADOW_COMMIT_ENTRY_REQUEST_SCHEMA, - "runtime_root": str(runtime_root), "goal_id": goal_id, - "entry_id": entry.entry_id, "partition": entry.partition, "seq": entry.seq, - "capture_lineage_id": entry.prepared.get("capture_lineage_id"), - "prepared_sha256": outbox.raw_bytes_digest(entry.prepared_path.read_bytes()), - "committed_sha256": outbox.raw_bytes_digest(entry.committed_path.read_bytes()) - if entry.committed_path else None, - } - - -def _valid_commit_entry_result(result: object, entry: outbox.OutboxEntry) -> bool: - if not isinstance(result, dict): - return False - return ( - result.get("schema_version") - == LOCAL_AUTHORITY_SHADOW_COMMIT_ENTRY_RESULT_SCHEMA - and result.get("outcome") in _COMMIT_ENTRY_OUTCOMES - and result.get("entry_id") == entry.entry_id - and result.get("partition") == entry.partition - and result.get("seq") == entry.seq - and isinstance(result.get("no_op"), bool) - ) - - def read_local_authority_shadow( *, runtime_root: Path, @@ -338,258 +203,6 @@ def read_local_authority_shadow( return dict(result) -class _DrainBudget: - def __init__(self, *, max_entries: int, budget_seconds: float) -> None: - self._max_entries = max(1, max_entries) - self._deadline = time.monotonic() + max(0.0, budget_seconds) - self.consumed = 0 - - def exhausted(self) -> bool: - return self.consumed >= self._max_entries or time.monotonic() >= self._deadline - - @property - def remaining_entries(self) -> int: - return max(0, self._max_entries - self.consumed) - - def can_reclaim(self, count: int) -> bool: - return ( - self.consumed + count <= self._max_entries - and time.monotonic() < self._deadline - ) - - -class _PartitionDrainer: - """Prove under M, release for the TS transaction, then reacquire before cleanup.""" - - def __init__( - self, - *, - registry_path: Path, - runtime_root: Path, - goal_id: str, - partition: str, - sources: _GoalSources, - result: DrainResult, - budget: _DrainBudget, - lock_timeout_seconds: float, - capture_lineage_id: str, - ) -> None: - self._runtime_root = runtime_root - self._goal_id = goal_id - self._partition = partition - self._sources = sources - self._result = result - self._budget = budget - self._lock_timeout = lock_timeout_seconds - self._directory = outbox.partition_directory(runtime_root, goal_id, partition) - self._lineage: str | None = capture_lineage_id - self.last_delivered_digest: str | None = None - - def _lock(self) -> Any: - return exclusive_cross_runtime_file_lock( - outbox.drain_lock_target(self._runtime_root, self._goal_id), - timeout_seconds=self._lock_timeout, - operation="local_authority_shadow_drain", - ) - - def _binding(self) -> dict[str, Any]: - view = read_shadow_capture_binding(self._runtime_root, self._goal_id) - if view["status"] != "active": - raise outbox.OutboxError( - str(view.get("reason_code") or "bootstrap_required"), - "shadow capture has no active binding", - ) - binding = dict(view["binding"]) - lineage = str(binding["capture_lineage_id"]) - if self._lineage is not None and self._lineage != lineage: - raise outbox.OutboxError( - "stale_generation", "drain belongs to an earlier lineage" - ) - self._lineage = lineage - return binding - - def _reconcile( - self, *, acknowledgement: dict[str, Any] | None = None, - ) -> tuple[list[outbox.OutboxEntry], int]: - binding = self._binding() - # Malformed cursor bytes remain evidence, even if the candidate is unavailable. - cursor = outbox.read_cursor(self._directory) - entries = outbox.list_entries(self._directory, allow_committed_only=True) - files: dict[str, list[tuple[Path, str]]] = {} - observations = [] - for entry in entries: - observation: dict[str, Any] = { - "entry_id": entry.entry_id, "seq": entry.seq, - "prepared": bool(entry.prepared), - "capture_lineage_id": entry.prepared.get("capture_lineage_id"), - } - files[entry.entry_id] = [] - for path, key in ( - (entry.prepared_path, "prepared_sha256"), - (entry.committed_path, "committed_sha256"), - ): - digest = ( - outbox.raw_bytes_digest(path.read_bytes()) - if path is not None and path.exists() else None - ) - observation[key] = digest - if digest is not None and path is not None: - files[entry.entry_id].append((path, digest)) - observations.append(observation) - plan = effect_runtime_result( - "coordination.runtime_shadow.plan_drain", - { - "schema_version": "loopx_shadow_drain_plan_request_v0", - "runtime_root": str(self._runtime_root), "goal_id": self._goal_id, - "partition": self._partition, - "capture_lineage_id": binding["capture_lineage_id"], - "store_identity": binding["store_identity"], - "source_root_digest": binding["source_root_digest"], - "cursor": cursor, "entries": observations, - "remaining_entries": self._budget.remaining_entries, - "budget_open": self._budget.can_reclaim(0), - "acknowledgement": acknowledgement, - }, - timeout=15.0, - ) - if not isinstance(plan, dict) or plan.get("schema_version") != "loopx_shadow_drain_plan_result_v0": - raise outbox.OutboxError("shadow_drain_result_invalid", "invalid drain plan") - view = plan.get("view") - if isinstance(view, dict): - if self._result.cursor_before is None: - self._result.cursor_before = view.get("cursor") - self._record_view(view) - if plan.get("status") != "planned": - raise outbox.OutboxError( - str(plan.get("reason_code") or "shadow_drain_result_invalid"), - "native drain plan rejected observations", - ) - # Time can expire during the native read; the plan cannot extend the budget. - if not self._budget.can_reclaim(0): - self._result.budget_exhausted = True - return [], int(plan["next_seq"]) - self._result.budget_exhausted |= plan["budget_exhausted"] - if plan["history_present"]: - with _primary_lock_if_free( - self._partition, runtime_root=self._runtime_root, - goal_id=self._goal_id, sources=self._sources, - ) as held: - if not held: - raise outbox.OutboxError("primary_writer_busy", "primary writer is in flight") - self._binding() - if ( - outbox.read_cursor(self._directory) != cursor - or outbox.list_entries(self._directory, allow_committed_only=True) != entries - ): - raise outbox.OutboxError("outbox_file_changed", "outbox changed during proof") - outbox.verify_observed_files(item for batch in files.values() for item in batch) - if plan["cursor_update"] is not None: - outbox.write_cursor( - self._directory, partition=self._partition, **plan["cursor_update"], - ) - self._result.reclaimed_residue += outbox.reclaim_verified_files([ - item for entry_id in plan["reclaim_entry_ids"] for item in files[entry_id] - ]) - for replay in plan["replay_entries"]: - summary = dict(replay) - self._result.no_op += int(summary.pop("no_op")) - self._result.entries.append(summary) - self._result.replayed += 1 - self._budget.consumed += 1 - pending_ids = set(plan["pending_entry_ids"]) - return [entry for entry in entries if entry.entry_id in pending_ids], int(plan["next_seq"]) - - def _record_view(self, view: dict[str, Any]) -> None: - self._result.candidate_readback_verified = True - self._result.store_identity = view.get("store_identity") - self._result.provider_revision = view.get("provider_revision") - self._result.last_cursor = view.get("cursor") - self._result.cursor_after = view.get("cursor") - self._result.head_digest = view.get("head_digest") - - def run(self) -> None: - while not self._budget.exhausted(): - with self._lock(): - pending, next_seq = self._reconcile() - if not pending: - return - if self._budget.exhausted(): - self._result.budget_exhausted = True - return - entry = pending[0] - if not entry.is_committed: - # Preserve the legacy flock busy signal. This host probe - # makes no source decision; TS rechecks under shared locks. - with _primary_lock_if_free( - self._partition, runtime_root=self._runtime_root, - goal_id=self._goal_id, sources=self._sources, - ) as held: - if not held: - raise outbox.OutboxError("primary_writer_busy", "primary writer is in flight") - if entry.seq != next_seq: - raise outbox.OutboxError( - "outbox_sequence_gap", "pending sequence is not continuous" - ) - request = _commit_entry_request( - runtime_root=self._runtime_root, - goal_id=self._goal_id, - entry=entry, - ) - # TS owns M for every public commit, including retries. Never re-enter M across RPC. - raw = effect_runtime_result( - "coordination.runtime_shadow.commit_entry", request, timeout=15.0 - ) - self._budget.consumed += 1 - if not _valid_commit_entry_result(raw, entry): - raise outbox.OutboxError( - "shadow_commit_entry_result_invalid", "invalid commit result" - ) - if raw["outcome"] not in _SETTLED_OUTCOMES: - self._result.stopped_at = { - "partition": entry.partition, - "seq": entry.seq, - "entry_id": entry.entry_id, - "outcome": raw["outcome"], - "reason_code": "outbox_source_unproved" if raw.get("reason_code") == "source_transaction_unproved" - else raw.get("reason_code"), - } - return - resolution = raw["resolution"] - digest = raw["partition_digest"] - with self._lock(): - self._reconcile(acknowledgement={ - "entry_id": entry.entry_id, "seq": entry.seq, - "cursor": raw.get("cursor"), - "provider_revision": raw.get("provider_revision"), - "store_identity": raw.get("store_identity"), - "no_op": raw.get("no_op"), "partition_digest": digest, - }) - summary = { - "entry_id": entry.entry_id, - "partition": entry.partition, - "seq": entry.seq, - "resolution": resolution, - "outcome": raw["outcome"], - "reason_code": raw.get("reason_code"), - "cursor": raw.get("cursor"), - "provider_revision": raw.get("provider_revision"), - "partition_digest": digest, - } - self._result.entries.append(summary) - if raw["outcome"] == "delivered": - self._result.delivered += 1 - elif raw["outcome"] == "replayed": - self._result.replayed += 1 - else: - self._result.reconciled += 1 - if raw["no_op"]: - self._result.no_op += 1 - elif digest is not None: - self.last_delivered_digest = digest - if outbox.list_entries(self._directory, allow_committed_only=True): - self._result.budget_exhausted = True - - def _drain_prelude( result: DrainResult, *, @@ -622,69 +235,6 @@ def _drain_prelude( return registry, resolved -def _drain_partitions( - result: DrainResult, - *, - registry: dict[str, Any], - registry_path: Path, - runtime_root: Path, - goal_id: str, - max_entries: int, - budget_seconds: float, - lock_timeout_seconds: float, -) -> None: - """Drain partitions through the shared management lock and TS commit owner.""" - - sources = _goal_sources(registry, runtime_root=runtime_root, goal_id=goal_id) - binding_view = read_shadow_capture_binding(runtime_root, goal_id) - if binding_view["status"] != "active": - raise outbox.OutboxError( - str(binding_view.get("reason_code") or "bootstrap_required"), - "drain requires an active capture lineage", - ) - capture_lineage_id = str(binding_view["binding"]["capture_lineage_id"]) - budget = _DrainBudget(max_entries=max_entries, budget_seconds=budget_seconds) - for partition in PARTITIONS: - if result.stopped_at is not None: - break - drainer = _PartitionDrainer( - registry_path=registry_path, - runtime_root=runtime_root, - goal_id=goal_id, - partition=partition, - sources=sources, - result=result, - budget=budget, - lock_timeout_seconds=lock_timeout_seconds, - capture_lineage_id=capture_lineage_id, - ) - drainer.run() - - -def _settle_drain_outcome(result: DrainResult) -> None: - if result.stopped_at is not None: - result.outcome = "stopped" - result.reason_code = str( - result.stopped_at.get("reason_code") or result.stopped_at["outcome"] - ) - else: - result.outcome = ( - "drained" - if result.drained_count or result.budget_exhausted - else "nothing_pending" - ) - - -def _count_backlog(result: DrainResult, runtime_root: Path, goal_id: str) -> None: - summary_after = outbox.outbox_summary(runtime_root, goal_id) - result.pending_after = sum( - int(item["committed_pending"]) for item in summary_after.values() - ) - result.prepared_only_after = sum( - int(item["prepared_only"]) for item in summary_after.values() - ) - - def drain_local_authority_shadow_outbox( *, registry_path: Path, @@ -694,59 +244,40 @@ def drain_local_authority_shadow_outbox( budget_seconds: float = INLINE_DRAIN_BUDGET_SECONDS, lock_timeout_seconds: float = INLINE_DRAIN_LOCK_TIMEOUT_SECONDS, ) -> DrainResult: - """Deliver pending outbox entries to the candidate store, one transaction each. + """Transport one native drain batch; never replay a timed-out invocation. - The drain lock is per goal. A held lock means another drainer is already - at work, so the caller's write stays ``pending`` instead of waiting on it. + A lost response leaves candidate commit/cleanup progress unknown. The next + explicit drain recovers it from durable receipts, not Python memory. """ - result = DrainResult(goal_id=goal_id) prelude = _drain_prelude( result, registry_path=registry_path, runtime_root=runtime_root, goal_id=goal_id ) if prelude is None: return result - registry, resolved_root = prelude - binding = read_shadow_capture_binding(resolved_root, goal_id) - if binding["status"] != "active": - runtime_enabled = resolve_coordination_runtime_shadow_config(find_registry_goal(registry, goal_id)).enabled - requires_bootstrap = (runtime_enabled or binding["status"] in {"inactive", "hold"} - or outbox.outbox_root(resolved_root, goal_id).exists()) - result.outcome = "stopped" if requires_bootstrap else "nothing_pending" - result.reason_code = ( - str(binding.get("reason_code") or "bootstrap_required") - if result.outcome == "stopped" - else None - ) + _registry, resolved_root = prelude + # No activation or persisted capture state: avoid starting the TS runtime + # for ordinary feature-off writes. Existing state is interpreted only by TS. + if (not result.config_enabled + and not shadow_management_state_path(resolved_root, goal_id).exists() + and not outbox.outbox_root(resolved_root, goal_id).exists()): return result try: - _drain_partitions( - result, - registry=registry, - registry_path=registry_path, - runtime_root=resolved_root, - goal_id=goal_id, - max_entries=max_entries, - budget_seconds=budget_seconds, - lock_timeout_seconds=lock_timeout_seconds, + raw = effect_runtime_result( + "coordination.runtime_shadow.drain", + {"schema_version": "loopx_shadow_drain_v0", "runtime_root": str(resolved_root), + "goal_id": goal_id, "python_executable": sys.executable, + "config_enabled": result.config_enabled, "max_entries": max_entries, + "budget_seconds": budget_seconds, "lock_timeout_seconds": lock_timeout_seconds}, + timeout=max(15.0, budget_seconds + 15.0), retry_safe=False, ) - except LockAcquireTimeoutError: - result.outcome = "drain_deferred" - result.reason_code = "drain_lock_busy" - except outbox.OutboxError as error: - result.outcome = "stopped" - result.reason_code = error.reason_code + if not isinstance(raw, dict) or raw.get("schema_version") != "loopx_shadow_drain_v0" or raw.get("goal_id") != goal_id: + raise ValueError("invalid native drain result") + return DrainResult(**{item.name: raw[item.name] for item in fields(DrainResult)}) except Exception: result.outcome = "stopped" - result.reason_code = "shadow_drain_failed" - else: - _settle_drain_outcome(result) - try: - _count_backlog(result, resolved_root, goal_id) - except Exception: - result.outcome = "stopped" - result.reason_code = result.reason_code or "outbox_status_unavailable" - return result + result.reason_code = "shadow_drain_outcome_unknown" + return result class _CandidateMissing(Exception): @@ -986,7 +517,6 @@ def valid_evidence_v1(result: object, *, goal_id: str) -> bool: "RETENTION_PRESSURE_BYTES", "DrainResult", "capture_evidence", - "primary_lock_is_free", "todo_partition_projector", "drain_local_authority_shadow_outbox", "local_authority_shadow_status", diff --git a/loopx/control_plane/coordination/shadow_drain.ts b/loopx/control_plane/coordination/shadow_drain.ts new file mode 100644 index 0000000000..0c1b4ac586 --- /dev/null +++ b/loopx/control_plane/coordination/shadow_drain.ts @@ -0,0 +1,175 @@ +/** One bounded drain invocation. The host supplies configuration, never source + * resolution, receipt proof or cleanup decisions. Candidate state is not authority. */ +import {performance} from "node:perf_hooks"; +import {isAbsolute, join} from "node:path"; +import {existsSync} from "node:fs"; +import type {JsonObject} from "../effect_program.ts"; +import {durableWriteJson, withFileMutationLock} from "../effect_runtime_io.ts"; +import {EffectRuntimeLockTimeoutError} from "../effect_runtime_errors.ts"; +import {requireJsonObject} from "../runtime_decode.ts"; +import {canonicalAuthoritySha256, hasExactAuthorityKeys, requireAuthorityStoreId} from "./authority_store_codec.ts"; +import {readShadowManagementState, requireShadowCaptureBinding, shadowMaintenanceLockPath} from "./shadow_management.ts"; +import {readShadowDrainPlan, SHADOW_DRAIN_PLAN_REQUEST_SCHEMA} from "./shadow_drain_plan.ts"; +import {deliverShadowEntry, SHADOW_ENTRY_DELIVERY_REQUEST_SCHEMA} from "./shadow_entry_delivery.ts"; +import {outboxPartitionDirectory, LOCAL_AUTHORITY_SHADOW_DRAIN_CURSOR_SCHEMA} from "./local_authority_shadow_outbox.ts"; +import {ShadowLineageError, type LocalAuthorityShadowDependencies} from "./local_authority_shadow.ts"; +import {DrainKernelLockHost, drainInventory, reclaimDrainFiles, verifyDrainFiles, withDrainPrimary, type DrainPartition, type DrainEntry} from "./shadow_drain_files.ts"; + +export const SHADOW_DRAIN_SCHEMA = "loopx_shadow_drain_v0"; +interface Request {runtime_root: string; goal_id: string; python_executable: string; config_enabled: boolean; + max_entries: number; budget_seconds: number; lock_timeout_seconds: number} +interface Dependencies extends LocalAuthorityShadowDependencies { + /** Scheduling-only fault seam; never accepted from a public request. */ + afterEffect?: (phase: "after_proof" | "before_commit" | "after_commit" | "after_cursor" | "after_unlink") => Promise; +} +function decode(value: unknown): Request { + const r = requireJsonObject(value, "drain request"); + if (!hasExactAuthorityKeys(r, ["schema_version", "runtime_root", "goal_id", "python_executable", "config_enabled", + "max_entries", "budget_seconds", "lock_timeout_seconds"]) || r.schema_version !== SHADOW_DRAIN_SCHEMA || + typeof r.runtime_root !== "string" || !isAbsolute(r.runtime_root) || r.runtime_root.includes("\0") || + typeof r.python_executable !== "string" || !isAbsolute(r.python_executable) || r.python_executable.includes("\0") || + typeof r.config_enabled !== "boolean" || !Number.isSafeInteger(r.max_entries) || Number(r.max_entries) < 1 || + [r.budget_seconds, r.lock_timeout_seconds].some(n => typeof n !== "number" || !Number.isFinite(n) || n < 0)) + throw new ShadowLineageError("shadow_drain_request_invalid"); + requireAuthorityStoreId(r.goal_id, "goal id"); + return r as unknown as Request; +} +function errorCode(error: unknown): string { + if (error instanceof EffectRuntimeLockTimeoutError) return "drain_lock_busy"; + const e = error as {reason_code?: string; code?: string}; + return e.reason_code ?? e.code ?? "shadow_drain_failed"; +} + +export async function drainShadowOutbox(value: unknown, dependencies: Dependencies = {}): Promise { + const r = decode(value), root = r.runtime_root, goal = r.goal_id; + const result = {goal_id: goal, outcome: "nothing_pending", config_enabled: r.config_enabled, + delivered: 0, replayed: 0, reconciled: 0, no_op: 0, reseeded: 0, reclaimed_residue: 0, + pending_after: 0, prepared_only_after: 0, in_flight_partitions: [] as string[], budget_exhausted: false, + stopped_at: null as JsonObject | null, reason_code: null as string | null, + store_identity: null as unknown, provider_revision: null as unknown, last_cursor: null as unknown, + cursor_before: null as unknown, cursor_after: null as unknown, head_digest: null as unknown, + candidate_readback_verified: null as boolean | null, entries: [] as JsonObject[]}; + const kernel = new DrainKernelLockHost(r.python_executable); + let consumed = 0; + const deadline = performance.now() + r.budget_seconds * 1000; + const timeOpen = () => performance.now() < deadline; + const finish = (): JsonObject => ({schema_version: SHADOW_DRAIN_SCHEMA, ...result, + ok: ["drained", "nothing_pending"].includes(result.outcome) && result.stopped_at === null, + drained_count: result.delivered + result.replayed + result.reconciled}); + const observe = (view: JsonObject) => { + result.candidate_readback_verified = true; + result.cursor_before ??= view.cursor; + result.store_identity = view.store_identity; result.provider_revision = view.provider_revision; + result.last_cursor = view.cursor; result.cursor_after = view.cursor; result.head_digest = view.head_digest; + }; + const locked = (operation: () => Promise): Promise => + withFileMutationLock(shadowMaintenanceLockPath(root, goal), operation, + Math.min(r.lock_timeout_seconds * 1000, Math.max(0, deadline - performance.now()))); + try { + const initial = await readShadowManagementState(root, goal); + if (initial?.status !== "active") { + if (r.config_enabled || initial !== null || existsSync(join(root, "authority-shadow", "outbox", goal))) { + result.outcome = "stopped"; result.reason_code = initial && initial.status !== "inactive" ? "shadow_management_in_progress" : "bootstrap_required"; + } + return finish(); + } + const lineage = (await requireShadowCaptureBinding(root, goal)).capture_lineage_id; + const binding = async () => { + const active = await requireShadowCaptureBinding(root, goal); + if (active.capture_lineage_id !== lineage) throw new ShadowLineageError("stale_generation"); + return active; + }; + const reconcile = async (partition: DrainPartition, acknowledgement: JsonObject | null = null): Promise => { + const active = await binding(), before = await drainInventory(root, goal, partition); + const plan = await readShadowDrainPlan({schema_version: SHADOW_DRAIN_PLAN_REQUEST_SCHEMA, runtime_root: root, goal_id: goal, + partition, capture_lineage_id: lineage, store_identity: active.store_identity, source_root_digest: active.source_root_digest, + cursor: before.cursor, entries: before.entries, remaining_entries: Math.max(0, r.max_entries - consumed), + budget_open: timeOpen(), acknowledgement}, dependencies); + if (plan.view !== null && typeof plan.view === "object") observe(plan.view as JsonObject); + if (plan.status !== "planned") throw new ShadowLineageError(String(plan.reason_code ?? "shadow_drain_result_invalid")); + await dependencies.afterEffect?.("after_proof"); + // Proof time and OS lock acquisition cannot extend the caller's effect budget. + if (!timeOpen()) {result.budget_exhausted = true; return [];} + result.budget_exhausted ||= plan.budget_exhausted === true; + if (plan.history_present) await withDrainPrimary(root, goal, partition, active, kernel, async () => { + await binding(); + const current = await drainInventory(root, goal, partition); + if (canonicalAuthoritySha256({cursor: current.cursor, entries: current.entries}) !== + canonicalAuthoritySha256({cursor: before.cursor, entries: before.entries})) throw new ShadowLineageError("outbox_file_changed"); + await verifyDrainFiles([...before.files.values()].flat()); + if (!timeOpen()) {result.budget_exhausted = true; return;} + if (plan.cursor_update !== null) { + await durableWriteJson(join(outboxPartitionDirectory(root, goal, partition), "drain-cursor.json"), { + schema_version: LOCAL_AUTHORITY_SHADOW_DRAIN_CURSOR_SCHEMA, partition, + ...plan.cursor_update as JsonObject, updated_at: new Date().toISOString(), + }); + await dependencies.afterEffect?.("after_cursor"); + } + const files = (plan.reclaim_entry_ids as string[]).flatMap(id => before.files.get(id) ?? []); + result.reclaimed_residue += await reclaimDrainFiles(files, () => dependencies.afterEffect?.("after_unlink") ?? Promise.resolve()); + for (const entry of plan.replay_entries as JsonObject[]) { + const {no_op: noOp, ...summary} = entry; + result.no_op += Number(noOp === true); result.entries.push(summary); result.replayed++; consumed++; + } + }); + if (!timeOpen()) {result.budget_exhausted = true; return [];} + const pending = new Set(plan.pending_entry_ids as string[]); + const entries = before.entries.filter(e => pending.has(e.entry_id)); + if (entries.length && entries[0].seq !== plan.next_seq) throw new ShadowLineageError("outbox_sequence_gap"); + return entries; + }; + for (const partition of ["todos", "leases"] as const) { + while (timeOpen() && consumed < r.max_entries && result.stopped_at === null) { + const entry = await locked(async () => { + const pending = await reconcile(partition); + if (!pending.length || !timeOpen() || consumed >= r.max_entries) return null; + const next = pending[0]; + if (!next.prepared) throw new ShadowLineageError("outbox_file_invalid"); + if (next.committed_sha256 === null) + await withDrainPrimary(root, goal, partition, await binding(), kernel, async () => {}); + return next; + }); + if (entry === null) break; + await dependencies.afterEffect?.("before_commit"); + if (!timeOpen()) {result.budget_exhausted = true; break;} + const raw = await deliverShadowEntry({schema_version: SHADOW_ENTRY_DELIVERY_REQUEST_SCHEMA, runtime_root: root, goal_id: goal, + partition, entry_id: entry.entry_id, seq: entry.seq, capture_lineage_id: entry.capture_lineage_id, + prepared_sha256: entry.prepared_sha256, committed_sha256: entry.committed_sha256}, dependencies); + consumed++; await dependencies.afterEffect?.("after_commit"); + if (!["delivered", "replayed", "ambiguous_reconciled"].includes(String(raw.outcome))) { + result.stopped_at = {partition, seq: entry.seq, entry_id: entry.entry_id, outcome: raw.outcome, + reason_code: raw.reason_code === "source_transaction_unproved" ? "outbox_source_unproved" : raw.reason_code}; + break; + } + await locked(() => reconcile(partition, {entry_id: entry.entry_id, seq: entry.seq, cursor: raw.cursor, + provider_revision: raw.provider_revision, store_identity: raw.store_identity, no_op: raw.no_op, partition_digest: raw.partition_digest})); + result.entries.push({entry_id: entry.entry_id, partition, seq: entry.seq, resolution: raw.resolution, + outcome: raw.outcome, reason_code: raw.reason_code, cursor: raw.cursor, provider_revision: raw.provider_revision, + partition_digest: raw.partition_digest}); + if (raw.outcome === "delivered") result.delivered++; + else if (raw.outcome === "replayed") result.replayed++; + else result.reconciled++; + if (raw.no_op === true) result.no_op++; + } + if (result.stopped_at !== null) break; + } + if (result.stopped_at !== null) { + result.outcome = "stopped"; result.reason_code = String(result.stopped_at.reason_code ?? result.stopped_at.outcome); + } else { + result.budget_exhausted ||= !timeOpen(); + result.outcome = consumed || result.budget_exhausted ? "drained" : "nothing_pending"; + } + } catch (error) { + result.reason_code = errorCode(error); + result.outcome = result.reason_code === "drain_lock_busy" ? "drain_deferred" : "stopped"; + } finally {await kernel.close();} + try { + for (const partition of ["todos", "leases"] as const) { + const inventory = await drainInventory(root, goal, partition); + result.pending_after += inventory.entries.filter(e => e.committed_sha256 !== null).length; + result.prepared_only_after += inventory.entries.filter(e => e.committed_sha256 === null).length; + } + result.budget_exhausted ||= consumed >= r.max_entries && result.pending_after + result.prepared_only_after > 0; + } catch (error) {result.outcome = "stopped"; result.reason_code ??= errorCode(error);} + return finish(); +} diff --git a/loopx/control_plane/coordination/shadow_drain_files.ts b/loopx/control_plane/coordination/shadow_drain_files.ts new file mode 100644 index 0000000000..38238237d2 --- /dev/null +++ b/loopx/control_plane/coordination/shadow_drain_files.ts @@ -0,0 +1,152 @@ +/** Filesystem effects for native drain. Receipt proof, not a cursor or filename, + * authorizes cleanup. All callers hold the shadow maintenance guard. */ +import {spawn, type ChildProcessWithoutNullStreams} from "node:child_process"; +import {createInterface} from "node:readline"; +import {lstat, readFile, readdir, unlink} from "node:fs/promises"; +import {join} from "node:path"; +import {fileURLToPath} from "node:url"; +import type {JsonObject} from "../effect_program.ts"; +import {withFileMutationLock} from "../effect_runtime_io.ts"; +import {EffectRuntimeLockTimeoutError} from "../effect_runtime_errors.ts"; +import {requireJsonObject} from "../runtime_decode.ts"; +import {ShadowLineageError} from "./local_authority_shadow.ts"; +import {outboxEntryIdentity, OUTBOX_ENTRY_FILE_PATTERN} from "./local_authority_shadow_identity.ts"; +import {sha256Digest, readOutboxCursor, outboxPartitionDirectory} from "./local_authority_shadow_outbox.ts"; +import {legacyCoordinationTodoLockPath, taskLeaseLockPath} from "./legacy_writer_lock_paths.ts"; +import {readShadowBootstrapSourcePath, type ShadowCaptureBinding} from "./shadow_management.ts"; +import {LOCAL_AUTHORITY_SHADOW_OUTBOX_ENTRY_SCHEMA, LOCAL_AUTHORITY_SHADOW_OUTBOX_COMMIT_SCHEMA} from "./coordination_state_contract.generated.ts"; + +export type DrainPartition = "todos" | "leases"; +export interface DrainEntry { + entry_id: string; seq: number; prepared: boolean; capture_lineage_id: string | null; + prepared_sha256: string | null; committed_sha256: string | null; +} +export interface DrainInventory {cursor: JsonObject | null; entries: DrainEntry[]; files: Map} +function ensure(value: unknown): asserts value {if (!value) throw new ShadowLineageError("outbox_file_invalid");} + +export async function drainInventory(root: string, goal: string, partition: DrainPartition): Promise { + const directory = outboxPartitionDirectory(root, goal, partition); + const entries = new Map(); + const files = new Map(); + let names: string[]; + try { names = await readdir(directory); } + catch (error) {if ((error as NodeJS.ErrnoException).code === "ENOENT") names = []; else throw error;} + const records = new Map(); + for (const name of names.sort()) { + const path = join(directory, name), info = await lstat(path); + ensure(info.isFile() && !info.isSymbolicLink()); + if (name === "drain-cursor.json") continue; + const match = OUTBOX_ENTRY_FILE_PATTERN.exec(name); + ensure(match !== null && Number(match[1]) > 0); + const seq = Number(match[1]), id = match[2], phase = match[3]; + const entry = entries.get(seq) ?? {entry_id: id, seq, prepared: false, capture_lineage_id: null, + prepared_sha256: null, committed_sha256: null}; + ensure(entry.entry_id === id); + const bytes = await readFile(path), digest = sha256Digest(bytes); + let record: JsonObject; + try {record = requireJsonObject(JSON.parse(new TextDecoder("utf-8", {fatal: true}).decode(bytes)), "outbox record");} + catch {throw new ShadowLineageError("outbox_file_invalid");} + ensure(record.entry_id === id && typeof record.capture_lineage_id === "string" && record.capture_lineage_id.length > 0); + if (phase === "prepared") { + const source = requireJsonObject(record.source, "source"), writer = requireJsonObject(record.writer, "writer"); + const ref = typeof source.bytes_digest === "string" && source.bytes_digest ? source.bytes_digest + : typeof source.event_id === "string" && source.event_id ? `event:${source.event_id}` + : typeof record.partition_digest === "string" && record.partition_digest ? `seed:${record.partition_digest}` : null; + ensure(record.schema_version === LOCAL_AUTHORITY_SHADOW_OUTBOX_ENTRY_SCHEMA && record.seq === seq && + record.goal_id === goal && record.partition === partition && + ["python", "typescript"].includes(String(writer.runtime)) && typeof writer.write_class === "string" && writer.write_class.length > 0 && + ["markdown_active_state", "state_event_log", "task_lease_record"].includes(String(source.kind)) && + typeof record.source_root_digest === "string" && /^sha256:[0-9a-f]{64}$/.test(record.source_root_digest) && ref !== null && + outboxEntryIdentity(goal, partition, seq, ref, record.capture_lineage_id, record.source_root_digest) === id); + entry.prepared = true; entry.capture_lineage_id = record.capture_lineage_id; entry.prepared_sha256 = digest; + } else { + ensure(record.schema_version === LOCAL_AUTHORITY_SHADOW_OUTBOX_COMMIT_SCHEMA); + entry.committed_sha256 = digest; + } + entries.set(seq, entry); records.set(`${id}:${phase}`, record); + files.set(id, [...(files.get(id) ?? []), [path, digest]]); + } + for (const entry of entries.values()) { + const marker = records.get(`${entry.entry_id}:committed`); + ensure(!entry.prepared || marker === undefined || marker.capture_lineage_id === entry.capture_lineage_id); + } + return {cursor: await readOutboxCursor(directory, partition), entries: [...entries.values()].sort((a,b) => a.seq-b.seq), files}; +} + +export async function verifyDrainFiles(files: [string, string][]): Promise { + for (const [path, digest] of files) { + const info = await lstat(path); + if (!info.isFile() || info.isSymbolicLink() || sha256Digest(await readFile(path)) !== digest) + throw new ShadowLineageError("outbox_file_changed"); + } +} +export async function reclaimDrainFiles(files: [string, string][], afterUnlink?: () => Promise): Promise { + await verifyDrainFiles(files); + // Prepared first deliberately permits committed-only residue after a crash. + for (const [path] of [...files].sort(([a],[b]) => Number(b.endsWith(".prepared.json"))-Number(a.endsWith(".prepared.json")))) { + await unlink(path); await afterUnlink?.(); + } + return files.length; +} + +/** One lazy OS-lock process per batch. It holds no locks between sections. */ +export class DrainKernelLockHost { + private readonly python: string; + private child: ChildProcessWithoutNullStreams | null = null; + private replies: AsyncIterator | null = null; + private exited: Promise | null = null; + constructor(python: string) {this.python = python;} + + private async command(line: string): Promise { + if (!this.child) { + this.child = spawn(this.python, ["-m", "loopx.control_plane.coordination.shadow_lock_host"], { + stdio: ["pipe", "pipe", "pipe"], env: {...process.env, PYTHONPATH: fileURLToPath(new URL("../../../", import.meta.url))}, + }); + this.exited = new Promise(resolve => this.child!.once("close", () => resolve())); + this.child.on("error", () => {}); // The closed reply stream reports an unavailable host. + this.child.stdin.on("error", () => {}); + this.child.stderr.resume(); + this.replies = createInterface({input: this.child.stdout})[Symbol.asyncIterator](); + } + let timer: ReturnType | undefined; + try { + this.child.stdin.write(line + "\n"); + const reply = await Promise.race([ + this.replies!.next().then(item => item.done ? "failed" : item.value), + this.exited!.then(() => "failed"), + new Promise(resolve => {timer = setTimeout(() => resolve("failed"), 5000);}), + ]); + if (reply === "failed") throw new ShadowLineageError("shadow_lock_host_unavailable"); + return reply; + } finally {clearTimeout(timer);} + } + + async withLocks(paths: string[], operation: () => Promise): Promise { + const state = await this.command(JSON.stringify(paths)); + if (state !== "held") throw new ShadowLineageError(state === "busy" ? "primary_writer_busy" : "shadow_lock_host_unavailable"); + try {return await operation();} + finally { + if (await this.command("release") !== "released") throw new ShadowLineageError("shadow_lock_host_unavailable"); + } + } + + async close(): Promise { + if (this.child) { + this.child.stdin.end(); this.child.kill(); await this.exited; + } + } +} + +/** Maintain the existing M → primary marker → kernel-lock order, with no wait + * for a busy primary writer. Killing the TS process closes the pipe as well. */ +export async function withDrainPrimary(root: string, goal: string, partition: DrainPartition, + binding: ShadowCaptureBinding, kernel: DrainKernelLockHost, operation: () => Promise): Promise { + const paths = partition === "todos" ? [legacyCoordinationTodoLockPath(root, goal), await readShadowBootstrapSourcePath(root, goal, binding)] + : [taskLeaseLockPath({runtime_root: root, goal_id: goal})]; + const lock = (i: number): Promise => i === paths.length ? kernel.withLocks(paths, operation) + : withFileMutationLock(paths[i], () => lock(i + 1), 0); + try {return await lock(0);} catch (error) { + if (error instanceof EffectRuntimeLockTimeoutError) throw new ShadowLineageError("primary_writer_busy"); + throw error; + } +} diff --git a/loopx/control_plane/coordination/shadow_drain_plan.ts b/loopx/control_plane/coordination/shadow_drain_plan.ts index 9f5ad13277..b446d79a23 100644 --- a/loopx/control_plane/coordination/shadow_drain_plan.ts +++ b/loopx/control_plane/coordination/shadow_drain_plan.ts @@ -180,8 +180,8 @@ export function planShadowDrain(value: unknown, rawView: unknown): JsonObject { })}; } -/** Same one read RPC as the old proof path, but history never crosses into Python. - * The caller holds M; taking it again here would deadlock across runtimes. */ +/** Native proof read for the batch owner; complete history stays in TypeScript. + * The caller holds M; taking it again here would deadlock. */ export async function readShadowDrainPlan(value: unknown, dependencies: LocalAuthorityShadowDependencies = {}): Promise { let verifiedView: JsonObject | null = null; diff --git a/loopx/control_plane/coordination/shadow_lock_host.py b/loopx/control_plane/coordination/shadow_lock_host.py new file mode 100644 index 0000000000..ded22785a2 --- /dev/null +++ b/loopx/control_plane/coordination/shadow_lock_host.py @@ -0,0 +1,39 @@ +"""Kernel-lock effect for the native drainer; EOF releases every held lock. + +The TS caller owns marker locks and all drain decisions. Keep OS-specific flock / +Windows locking in the existing Python adapter and use the caller's interpreter. +""" +from contextlib import ExitStack +import json +from pathlib import Path +import sys + +from ...file_lock import try_exclusive_file_lock + + +def main() -> None: + for line in sys.stdin: + paths = json.loads(line) + if not isinstance(paths, list) or not paths or any( + not isinstance(path, str) or not Path(path).is_absolute() for path in paths + ): + raise ValueError("absolute lock targets required") + with ExitStack() as stack: + for path in paths: + held = stack.enter_context(try_exclusive_file_lock( + Path(path), operation="local_authority_shadow_drain" + )) + if held is None: + print("busy", flush=True) + break + else: + print("held", flush=True) + # EOF (including parent death) always releases the kernel locks. + if sys.stdin.readline() != "release\n": + return + stack.close() + print("released", flush=True) + + +if __name__ == "__main__": + main() diff --git a/loopx/control_plane/effect_runtime_handlers.ts b/loopx/control_plane/effect_runtime_handlers.ts index ce15ff6852..f892db08ed 100644 --- a/loopx/control_plane/effect_runtime_handlers.ts +++ b/loopx/control_plane/effect_runtime_handlers.ts @@ -2,7 +2,7 @@ import {planStateEventReplay} from "./goals/state_event_replay.ts"; import {projectTodoSummary} from "./todos/summary_projection.ts"; import {admitAutomationStart, confirmAutomationStart, manageAutomationCadence, projectCadenceSchedule} from "./quota/automation_cadence.ts"; import {deliverShadowEntry} from "./coordination/shadow_entry_delivery.ts"; -import {readShadowDrainPlan} from "./coordination/shadow_drain_plan.ts"; +import {drainShadowOutbox} from "./coordination/shadow_drain.ts"; import {readCanonicalSnapshotPage} from "./coordination/canonical_snapshot_page.ts"; import {manageLocalAuthorityArchive} from "./coordination/local_authority_archive.ts"; import {selectPeriodicReportProgress, selectPeriodicReportApprovalRetry} from "./capabilities/periodic_report_progress.ts"; @@ -617,7 +617,7 @@ export function createEffectRuntimeHandlers( ["coordination.local_authority_shadow.record", recordLocalAuthorityShadow], ["coordination.runtime_shadow.commit_entry", deliverShadowEntry], ["coordination.runtime_shadow.outbox_read", readLocalAuthorityShadow], - ["coordination.runtime_shadow.plan_drain", readShadowDrainPlan], + ["coordination.runtime_shadow.drain", drainShadowOutbox], [ "effect.program_from_ordered_steps", (params) => effectProgramFromOrderedSteps( diff --git a/loopx/control_plane/effect_runtime_io.ts b/loopx/control_plane/effect_runtime_io.ts index a5be4f73cc..d994d0a8fe 100644 --- a/loopx/control_plane/effect_runtime_io.ts +++ b/loopx/control_plane/effect_runtime_io.ts @@ -237,16 +237,16 @@ export async function releaseFileMutationLockClaim( await removeCreatedFile(claim.claimPath, claim.identity); } -async function reclaimStaleMutationLock(path: string): Promise { +async function reclaimStaleMutationLock(path: string): Promise { const identity = await readFileIdentity(path); - if (!identity) return; + if (!identity) return false; const owner = await readMutationLockOwner(path); - if (owner && processIsAlive(owner.pid)) return; + if (owner && processIsAlive(owner.pid)) return false; if (!owner) { try { - if (Date.now() - (await stat(path)).mtimeMs < INVALID_LOCK_STALE_MS) return; + if (Date.now() - (await stat(path)).mtimeMs < INVALID_LOCK_STALE_MS) return false; } catch { - return; + return false; } } const targetPath = path.slice(0, -".ts-effect.lock".length); @@ -254,29 +254,29 @@ async function reclaimStaleMutationLock(path: string): Promise { targetPath, owner?.token ?? INVALID_LOCK_CLAIM_TOKEN, ); - if (!claim) return; + if (!claim) return false; const stalePath = `${path}.stale.${randomUUID()}`; try { const current = await readMutationLockOwner(path); - if (owner && (!current || current.token !== owner.token)) return; - if (current && processIsAlive(current.pid)) return; + if (owner && (!current || current.token !== owner.token)) return false; + if (current && processIsAlive(current.pid)) return false; if (!current) { try { if (Date.now() - (await stat(path)).mtimeMs < INVALID_LOCK_STALE_MS) { - return; + return false; } } catch { - return; + return false; } } // The lock pathname is not a compare-and-swap primitive. Holding the // token claim serializes compliant writers; the identity check additionally // prevents a replacement inode from being retired after a stale read. - if (!sameFileIdentity(identity, await readFileIdentity(path))) return; + if (!sameFileIdentity(identity, await readFileIdentity(path))) return false; await rename(path, stalePath); } catch (error) { if ((error as NodeJS.ErrnoException).code !== "ENOENT") throw error; - return; + return false; } finally { if (claim) { try { @@ -287,6 +287,7 @@ async function reclaimStaleMutationLock(path: string): Promise { } } await rm(stalePath, { force: true }); + return true; } export interface FileMutationLock { @@ -310,6 +311,7 @@ export async function acquireFileMutationLock( const lockPath = `${targetPath}.ts-effect.lock`; const token = randomUUID(); const deadline = Date.now() + timeoutMs; + let retriedAfterReclaim = false; while (true) { try { const handle = await open(lockPath, "wx", 0o600); @@ -354,7 +356,14 @@ export async function acquireFileMutationLock( return { targetPath, lockPath, token }; } catch (error) { if ((error as NodeJS.ErrnoException).code !== "EEXIST") throw error; - await reclaimStaleMutationLock(lockPath); + const reclaimed = await reclaimStaleMutationLock(lockPath); + // Removing a dead owner is progress, not waiting for a live owner. Allow + // one immediate acquisition even for a zero-wait caller. Bound the retry + // so repeated replacement cannot extend that caller's lock budget. + if (reclaimed && !retriedAfterReclaim) { + retriedAfterReclaim = true; + continue; + } if (Date.now() >= deadline) { throw new EffectRuntimeLockTimeoutError(); } diff --git a/loopx/control_plane/testing/authority_e2e_rows_stage2c2.py b/loopx/control_plane/testing/authority_e2e_rows_stage2c2.py index 50e3cbac80..cbaa96ac9d 100644 --- a/loopx/control_plane/testing/authority_e2e_rows_stage2c2.py +++ b/loopx/control_plane/testing/authority_e2e_rows_stage2c2.py @@ -111,48 +111,47 @@ from loopx.control_plane.coordination import local_authority_shadow_adapter as adapter from loopx.control_plane.coordination import local_authority_shadow_outbox as outbox from loopx.control_plane.todos import active_state_editing -window, state = sys.argv[1], pathlib.Path(sys.argv[2]).resolve() -def pause(): - print('BARRIER ' + json.dumps({'window': window}), flush=True) +window, state = sys.argv[1], pathlib.Path(sys.argv[2]) +def pause(payload=None): + print('BARRIER ' + json.dumps(payload or {}), flush=True) time.sleep(40) raise RuntimeError('parent failed to terminate at persistence barrier') actual_rpc = adapter.effect_runtime_result def rpc(method, request, **kwargs): - if method == 'coordination.runtime_shadow.commit_entry' and window == 'before_commit': - pause() - result = actual_rpc(method, request, **kwargs) - if method == 'coordination.runtime_shadow.commit_entry' and window == 'after_commit': - pause() - return result + if method == 'coordination.runtime_shadow.drain' and window in {'before_commit', 'after_commit', 'after_cursor', 'between_unlinks'}: + import subprocess + child = subprocess.Popen(['node', '--no-warnings', '--experimental-strip-types', + 'loopx/control_plane/testing/shadow_drain_fault_process.ts', window], + stdin=subprocess.PIPE, stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True) + child.stdin.write(json.dumps(request)); child.stdin.close() + for line in child.stdout: + if line.startswith('BARRIER '): + print(line, end='', flush=True) + sys.stdin.readline() # Parent requests death at the observed barrier. + child.kill(); child.wait(timeout=10) + print('REAPED', flush=True) + time.sleep(40) + raise RuntimeError('parent failed to terminate crash worker') + raise RuntimeError('native barrier missing: ' + child.stderr.read()) + return actual_rpc(method, request, **kwargs) adapter.effect_runtime_result = rpc -actual_cursor = outbox.write_cursor -def cursor(*args, **kwargs): - result = actual_cursor(*args, **kwargs) - if window == 'after_cursor': - pause() - return result -outbox.write_cursor = cursor actual_json = outbox.durable_write_json def write_json(path, value): - if window == 'before_marker' and path.name.endswith('.committed.json'): - pause() + if window == 'before_marker' and path.name.endswith('.committed.json'): pause() return actual_json(path, value) outbox.durable_write_json = write_json actual_replace = active_state_editing.os.replace def replace(source, target): - is_primary = pathlib.Path(target).resolve() == state - if is_primary and window == 'before_replace': - pause() + is_primary = pathlib.Path(target) == state + if is_primary and window == 'before_replace': pause() result = actual_replace(source, target) - if is_primary and window == 'after_replace': - pause() + if is_primary and window == 'after_replace': pause() return result active_state_editing.os.replace = replace actual_unlink = pathlib.Path.unlink def unlink(path, *args, **kwargs): result = actual_unlink(path, *args, **kwargs) - if window == 'between_unlinks' and path.name.endswith('.prepared.json'): - pause() + if window == 'between_unlinks' and path.name.endswith('.prepared.json'): pause() return result pathlib.Path.unlink = unlink raise SystemExit(main(sys.argv[3:])) @@ -425,6 +424,7 @@ def crash_cli(workspace: GoalWorkspace, window: str, *args: str) -> None: command, cwd=REPO_ROOT, env=cli_env(workspace), + stdin=subprocess.PIPE, stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True, encoding="utf-8", errors="replace", @@ -436,6 +436,14 @@ def crash_cli(workspace: GoalWorkspace, window: str, *args: str) -> None: if readable: line = process.stdout.readline() finally: + if line.startswith("BARRIER "): + barrier = json.loads(line.removeprefix("BARRIER ")) + if barrier.get("native_pid"): + assert process.stdin is not None and process.stdout is not None + process.stdin.write("terminate_native\n") + process.stdin.flush() + readable, _, _ = select.select([process.stdout], [], [], 10) + expect(bool(readable) and process.stdout.readline().strip() == "REAPED", "native owner must be reaped") process.kill() _, stderr = process.communicate(timeout=10) expect(line.startswith("BARRIER "), f"{window}: the CLI did not reach its persistence window: {stderr[-200:]}") diff --git a/loopx/control_plane/testing/shadow_drain_fault_process.ts b/loopx/control_plane/testing/shadow_drain_fault_process.ts new file mode 100644 index 0000000000..06b7164d65 --- /dev/null +++ b/loopx/control_plane/testing/shadow_drain_fault_process.ts @@ -0,0 +1,13 @@ +/** Private-process test driver, never installed as an RPC effect. */ +import {drainShadowOutbox} from "../coordination/shadow_drain.ts"; +let input = ""; +for await (const bytes of process.stdin) input += bytes; +const request = JSON.parse(input); +const phase = process.argv[2] === "between_unlinks" ? "after_unlink" : process.argv[2]; +const result = await drainShadowOutbox(request, {afterEffect: async observed => { + if (observed === phase) { + process.stdout.write(`BARRIER ${JSON.stringify({native_pid: process.pid, request})}\n`); + await new Promise(() => {setInterval(() => {}, 1000);}); + } +}}); +process.stdout.write(JSON.stringify(result) + "\n"); diff --git a/tests/control_plane/shadow_e2e_fixture.py b/tests/control_plane/shadow_e2e_fixture.py index 85471c8f83..caf9e4225e 100644 --- a/tests/control_plane/shadow_e2e_fixture.py +++ b/tests/control_plane/shadow_e2e_fixture.py @@ -9,6 +9,8 @@ import sys from dataclasses import dataclass +from loopx.control_plane.testing.authority_e2e_rows_stage2c2 import CRASH_WORKER + REPO = Path(__file__).resolve().parents[2] @@ -67,6 +69,7 @@ def crash(self, window: str, *args: str) -> dict: *self.arguments(*args), ], cwd=REPO, + stdin=subprocess.PIPE, stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True, @@ -81,6 +84,12 @@ def crash(self, window: str, *args: str) -> dict: stdout, stderr = child.communicate(timeout=10) raise AssertionError(f"No process barrier: {line}{stdout}\n{stderr}") payload = json.loads(line.removeprefix("BARRIER ")) + if payload.get("native_pid"): + assert child.stdin is not None + child.stdin.write("terminate_native\n") + child.stdin.flush() + readable, _, _ = select.select([child.stdout], [], [], 10) + assert readable and child.stdout.readline().strip() == "REAPED", "native owner must be reaped" child.kill() child.communicate(timeout=10) assert child.returncode == -9 @@ -130,52 +139,3 @@ def workspace(path: Path, *, bootstrap: bool = True) -> ShadowWorkspace: boot = result.cli("coordination-shadow", "bootstrap", "--execute")["bootstrap"] assert boot["status"] == "applied", boot return result - - -CRASH_WORKER = r""" -import json, pathlib, sys, time -from loopx.cli import main -from loopx.control_plane.coordination import local_authority_shadow_adapter as adapter -from loopx.control_plane.coordination import local_authority_shadow_outbox as outbox -from loopx.control_plane.todos import active_state_editing -window, state = sys.argv[1], pathlib.Path(sys.argv[2]) -def pause(payload=None): - print('BARRIER ' + json.dumps(payload or {}), flush=True) - time.sleep(40) - raise RuntimeError('parent failed to terminate at persistence barrier') -actual_rpc = adapter.effect_runtime_result -def rpc(method, request, **kwargs): - if method == 'coordination.runtime_shadow.commit_entry' and window == 'before_commit': - pause({'request': request}) - result = actual_rpc(method, request, **kwargs) - if method == 'coordination.runtime_shadow.commit_entry' and window == 'after_commit': - pause({'request': request, 'result': result}) - return result -adapter.effect_runtime_result = rpc -actual_cursor = outbox.write_cursor -def cursor(*args, **kwargs): - result = actual_cursor(*args, **kwargs) - if window == 'after_cursor': pause() - return result -outbox.write_cursor = cursor -actual_json = outbox.durable_write_json -def write_json(path, value): - if window == 'before_marker' and path.name.endswith('.committed.json'): pause() - return actual_json(path, value) -outbox.durable_write_json = write_json -actual_replace = active_state_editing.os.replace -def replace(source, target): - is_primary = pathlib.Path(target) == state - if is_primary and window == 'before_replace': pause() - result = actual_replace(source, target) - if is_primary and window == 'after_replace': pause() - return result -active_state_editing.os.replace = replace -actual_unlink = pathlib.Path.unlink -def unlink(path, *args, **kwargs): - result = actual_unlink(path, *args, **kwargs) - if window == 'between_unlinks' and path.name.endswith('.prepared.json'): pause() - return result -pathlib.Path.unlink = unlink -raise SystemExit(main(sys.argv[3:])) -""" diff --git a/tests/control_plane/test_local_authority_shadow_drain.py b/tests/control_plane/test_local_authority_shadow_drain.py index ae2b3f9df2..1206db67c1 100644 --- a/tests/control_plane/test_local_authority_shadow_drain.py +++ b/tests/control_plane/test_local_authority_shadow_drain.py @@ -1,3 +1,4 @@ +# Native persistence-window regressions live in shadow_drain.test.ts and the real-process CLI suite. from __future__ import annotations import json @@ -186,35 +187,6 @@ def test_drain_delivers_each_committed_entry_once_in_order_and_verifies_readback assert again.ok is True -def test_drain_replays_when_store_committed_but_cursor_was_not_written( - tmp_path: Path, monkeypatch: pytest.MonkeyPatch -) -> None: - registry, state, runtime_root = _fixture(tmp_path) - capture = _record_todo_write(registry, state, runtime_root, "Crash after store commit") - - def crash(*_args: object, **_kwargs: object) -> None: - raise OSError("simulated crash before the drain cursor landed") - - monkeypatch.setattr(outbox, "write_cursor", crash) - first = _drain(registry, runtime_root) - assert first.outcome == "stopped" - assert first.reason_code == "shadow_drain_failed" - assert first.pending_after == 1 - monkeypatch.undo() - - calls = _commit_entry_calls(monkeypatch) - second = _drain(registry, runtime_root) - assert second.ok is True - assert "coordination.runtime_shadow.commit_entry" not in calls - assert (second.delivered, second.replayed) == (0, 1) - assert second.entries[0]["entry_id"] == capture.outcome.entry_id - assert second.entries[0]["cursor"] == "2" - assert second.pending_after == 0 - assert second.candidate_readback_verified is True - view = adapter.read_local_authority_shadow(runtime_root=runtime_root, goal_id=GOAL_ID, scan_limit=10) - assert len(view["scan"]["transactions"]) == 2 - - def test_drain_defers_when_another_drainer_holds_the_lock(tmp_path: Path) -> None: registry, state, runtime_root = _fixture(tmp_path) _record_todo_write(registry, state, runtime_root, "Pending behind a drainer") @@ -251,42 +223,6 @@ def test_drain_batch_is_bounded_and_reports_what_it_left(tmp_path: Path) -> None assert view["head"]["partitions"]["todos"]["seq"] == 3 -def test_drain_stops_in_order_on_real_candidate_corruption( - tmp_path: Path, monkeypatch: pytest.MonkeyPatch -) -> None: - registry, state, runtime_root = _fixture(tmp_path) - for index in range(3): - _record_todo_write(registry, state, runtime_root, f"Ordered {index}") - real = adapter.effect_runtime_result - calls: list[str] = [] - saved: list[bytes] = [] - candidate = next((runtime_root / "authority-shadow" / "file-v0").glob("authority-store-*.json")) - - def corrupt_before_second_commit(method: str, params: object, **kwargs: object) -> object: - if method == "coordination.runtime_shadow.commit_entry": - calls.append(method) - if len(calls) == 2: - saved.append(candidate.read_bytes()) - candidate.write_text("{malformed candidate history") - # The real TypeScript handler and real FileAuthorityStore decide every result. - return real(method, params, **kwargs) - - monkeypatch.setattr(adapter, "effect_runtime_result", corrupt_before_second_commit) - result = _drain(registry, runtime_root) - assert result.outcome == "stopped" - assert result.delivered == 1 - assert result.stopped_at is not None and result.stopped_at["seq"] == 2 - assert result.stopped_at["outcome"] in {"failed", "unavailable"} - assert result.pending_after == 2 - assert [entry.seq for entry in outbox.list_entries(_todo_dir(runtime_root))] == [2, 3] - assert candidate.read_text() == "{malformed candidate history" - monkeypatch.undo() - candidate.write_bytes(saved[0]) - recovered = _drain(registry, runtime_root) - assert recovered.delivered == 2 - assert recovered.pending_after == 0 - - def test_prepared_only_entries_resolve_only_under_a_free_primary_lock(tmp_path: Path) -> None: registry, state, runtime_root = _fixture(tmp_path) proven = _record_todo_write(registry, state, runtime_root, "Marker lost after write", mark_committed=False) @@ -482,102 +418,6 @@ def test_capture_evidence_v1_reports_measured_facts_only(tmp_path: Path) -> None assert adapter.valid_evidence_v1(failed_evidence, goal_id=GOAL_ID) -def _commit_entry_calls(monkeypatch: pytest.MonkeyPatch) -> list[str]: - calls: list[str] = [] - real = adapter.effect_runtime_result - - def counting(method: str, params: object, **kwargs: object) -> object: - calls.append(method) - return real(method, params, **kwargs) - - monkeypatch.setattr(adapter, "effect_runtime_result", counting) - return calls - - -def test_crash_between_the_two_unlinks_leaves_residue_the_next_drain_reclaims( - tmp_path: Path, monkeypatch: pytest.MonkeyPatch -) -> None: - registry, state, runtime_root = _fixture(tmp_path) - capture = _record_todo_write(registry, state, runtime_root, "Retired but half-removed") - real_remove = outbox.reclaim_verified_files - - def crash_between_unlinks(files: object) -> None: - batch = list(files) - batch[0][0].unlink() - raise OSError("simulated crash between the prepared and committed unlinks") - - monkeypatch.setattr(outbox, "reclaim_verified_files", crash_between_unlinks) - first = _drain(registry, runtime_root) - assert first.outcome == "stopped" - assert first.reason_code == "shadow_drain_failed" - monkeypatch.setattr(outbox, "reclaim_verified_files", real_remove) - - # On disk: the cursor covers seq 1 and only the committed marker survives. - marker_name = outbox.entry_file_name(1, str(capture.outcome.entry_id), "committed") - names = sorted(path.name for path in _todo_dir(runtime_root).iterdir()) - assert names == sorted([marker_name, "drain-cursor.json"]) - assert outbox.read_cursor(_todo_dir(runtime_root))["last_seq"] == 1 - # The marker is retired residue, not corruption: listing stays valid. - assert len(outbox.list_entries(_todo_dir(runtime_root), allow_committed_only=True)) == 1 - assert [path.name for path in outbox.retired_residue(_todo_dir(runtime_root))] == [marker_name] - summary = outbox.outbox_summary(runtime_root, GOAL_ID)["todos"] - assert summary["invalid"] == "outbox_file_invalid" - assert len(outbox.retired_residue(_todo_dir(runtime_root))) == 1 - assert summary["committed_pending"] == 0 - status = adapter.local_authority_shadow_status(registry_path=registry, runtime_root=runtime_root, goal_id=GOAL_ID) - assert status["ok"] is False - assert status["outbox"]["todos"]["invalid"] == "outbox_file_invalid" - - calls = _commit_entry_calls(monkeypatch) - second = _drain(registry, runtime_root) - assert second.ok is True - assert second.outcome == "drained" - assert second.reclaimed_residue == 1 - assert (second.delivered, second.replayed) == (0, 1) - assert "coordination.runtime_shadow.commit_entry" not in calls - assert list(_todo_dir(runtime_root).iterdir()) == [_todo_dir(runtime_root) / "drain-cursor.json"] - view = adapter.read_local_authority_shadow(runtime_root=runtime_root, goal_id=GOAL_ID, scan_limit=5) - assert view["cursor"] == "2" - - # A later write mints seq 2 from the cursor, never reusing the retired seq. - later = _record_todo_write(registry, state, runtime_root, "After the reclaim") - assert later.outcome.seq == 2 - third = _drain(registry, runtime_root) - assert third.delivered == 1 - # The newly delivered transaction also removes its two verified files. - assert third.reclaimed_residue == 2 - - -def test_crash_after_the_cursor_but_before_any_unlink_is_reclaimed_without_a_store_call( - tmp_path: Path, monkeypatch: pytest.MonkeyPatch -) -> None: - registry, state, runtime_root = _fixture(tmp_path) - capture = _record_todo_write(registry, state, runtime_root, "Cursor written, files untouched") - real_remove = outbox.reclaim_verified_files - - def crash_before_unlinks(files: object) -> None: - raise OSError("simulated crash after the cursor write") - - monkeypatch.setattr(outbox, "reclaim_verified_files", crash_before_unlinks) - assert _drain(registry, runtime_root).outcome == "stopped" - monkeypatch.setattr(outbox, "reclaim_verified_files", real_remove) - names = sorted(path.name for path in _todo_dir(runtime_root).iterdir()) - entry_id = str(capture.outcome.entry_id) - assert names == [ - outbox.entry_file_name(1, entry_id, "committed"), - outbox.entry_file_name(1, entry_id, "prepared"), - "drain-cursor.json", - ] - assert len(outbox.list_entries(_todo_dir(runtime_root), allow_committed_only=True)) == 1 - - calls = _commit_entry_calls(monkeypatch) - result = _drain(registry, runtime_root) - assert result.ok is True - assert result.reclaimed_residue == 2 - assert "coordination.runtime_shadow.commit_entry" not in calls - assert list(_todo_dir(runtime_root).iterdir()) == [_todo_dir(runtime_root) / "drain-cursor.json"] - - def test_an_orphan_marker_above_the_cursor_is_still_corruption(tmp_path: Path) -> None: registry, state, runtime_root = _fixture(tmp_path) capture = _record_todo_write(registry, state, runtime_root, "Marker without its prepared file") diff --git a/tests/control_plane/test_local_authority_shadow_outbox.py b/tests/control_plane/test_local_authority_shadow_outbox.py index 19d6cebeec..792a879aa3 100644 --- a/tests/control_plane/test_local_authority_shadow_outbox.py +++ b/tests/control_plane/test_local_authority_shadow_outbox.py @@ -22,7 +22,6 @@ build_runtime_shadow_source_snapshot, ) from loopx.control_plane.coordination.shadow_management import require_shadow_primary_write_allowed -from loopx.file_lock import exclusive_file_lock from loopx.history import load_registry from loopx.registry import find_registry_goal @@ -338,15 +337,6 @@ def test_capture_failure_is_typed_and_preserves_the_primary_result(tmp_path: Pat assert state.read_text() == new_text -def test_primary_lock_probe_reports_held_locks(tmp_path: Path) -> None: - target = tmp_path / "ACTIVE_GOAL_STATE.md" - target.write_text("", encoding="utf-8") - assert adapter.primary_lock_is_free(target) is True - with exclusive_file_lock(target, timeout_seconds=1.0, operation="test_hold"): - assert adapter.primary_lock_is_free(target) is False - assert adapter.primary_lock_is_free(target) is True - - def test_prepared_records_must_bind_their_directory_identity_and_source(tmp_path: Path) -> None: registry, state, runtime_root = _fixture(tmp_path) original = state.read_text(encoding="utf-8") diff --git a/tests/control_plane/test_shadow_drain_adversarial.py b/tests/control_plane/test_shadow_drain_adversarial.py index 976b9ab387..74a1153c15 100644 --- a/tests/control_plane/test_shadow_drain_adversarial.py +++ b/tests/control_plane/test_shadow_drain_adversarial.py @@ -136,7 +136,11 @@ def test_native_markerless_resolution_requires_source_evidence( w.crash(window, "todo", "add", "--role", "agent", "--text", "Source proof is not a caller flag") directory = outbox.partition_directory(w.runtime, w.goal, "todos") [entry] = outbox.list_entries(directory) - request = adapter._commit_entry_request(runtime_root=w.runtime, goal_id=w.goal, entry=entry) + request = {"schema_version": "loopx_shadow_entry_delivery_request_v0", "runtime_root": str(w.runtime), + "goal_id": w.goal, "partition": entry.partition, "seq": entry.seq, "entry_id": entry.entry_id, + "capture_lineage_id": entry.prepared["capture_lineage_id"], + "prepared_sha256": outbox.raw_bytes_digest(entry.prepared_path.read_bytes()), + "committed_sha256": None} request["resolution"] = claimed_resolution before = {path.name: path.read_bytes() for path in directory.iterdir()} with pytest.raises(EffectRuntimeRejected, match="shadow_entry_selection_invalid"): diff --git a/tests/control_plane/test_shadow_drain_native_plan.py b/tests/control_plane/test_shadow_drain_native_plan.py index c2dcd81c07..c2b5801f31 100644 --- a/tests/control_plane/test_shadow_drain_native_plan.py +++ b/tests/control_plane/test_shadow_drain_native_plan.py @@ -1,137 +1,74 @@ -"""Native recovery plans execute only against the filesystem observations they proved.""" +"""The Python entry transports one batch, with no source/history decision owner. + +Filesystem mutation and failure-window cases moved to shadow_drain.test.ts and +real TS-process SIGKILL tests; patching retired Python internals proves nothing. +""" from __future__ import annotations from pathlib import Path +import json from typing import Any import pytest from loopx.control_plane.coordination import local_authority_shadow_adapter as adapter -from loopx.control_plane.coordination import local_authority_shadow_outbox as outbox from test_local_authority_shadow_drain import _drain, _fixture, _record_todo_write -PLAN = "coordination.runtime_shadow.plan_drain" -COMMIT = "coordination.runtime_shadow.commit_entry" +def test_disabled_without_capture_state_does_not_start_native_runtime(tmp_path, monkeypatch): + registry = tmp_path / "registry.json" + registry.write_text(json.dumps({"goals": [{"id": "goal-e2e"}]})) + def unexpected(*args, **kwargs): + raise AssertionError("feature-off drain must not invoke the runtime") -def inventory(directory: Path) -> dict[str, bytes]: - return {p.name: p.read_bytes() for p in directory.iterdir()} + monkeypatch.setattr(adapter, "effect_runtime_result", unexpected) + result = _drain(registry, tmp_path / "runtime") + assert result.ok and result.outcome == "nothing_pending" + assert not (tmp_path / "runtime").exists() -def test_adapter_exchanges_observations_and_effects_without_history( +def test_adapter_transports_one_batch_without_per_entry_rpc( tmp_path: Path, monkeypatch: pytest.MonkeyPatch, ) -> None: registry, state, runtime = _fixture(tmp_path) - _record_todo_write(registry, state, runtime, "Retained projection text is not drain transport") + for i in range(3): + _record_todo_write(registry, state, runtime, f"Private source text {i}") actual = adapter.effect_runtime_result - seen: list[tuple[str, dict[str, Any], dict[str, Any]]] = [] + seen = [] def rpc(method: str, request: dict, **kwargs: Any) -> dict: result = actual(method, request, **kwargs) - seen.append((method, request, result)) + seen.append((method, request, result, kwargs)) return result monkeypatch.setattr(adapter, "effect_runtime_result", rpc) result = _drain(registry, runtime) - assert result.ok and result.delivered == 1 - plans = [(request, response) for method, request, response in seen if method == PLAN] - assert len(plans) >= 2 # Preflight and post-commit receipt validation. - assert sum(method == COMMIT for method, _, _ in seen) == 1 - for request, response in plans: - assert "proof" not in response and "transactions" not in response - assert "head" not in response["view"] - for entry in request["entries"]: - assert set(entry) == { - "entry_id", "seq", "prepared", "capture_lineage_id", - "prepared_sha256", "committed_sha256", - } - assert any(request["acknowledgement"] is not None for request, _ in plans) - assert not any(method.endswith("outbox_read") for method, _, _ in seen) - - -@pytest.mark.parametrize("mutation", ["cursor", "entry", "new_entry"]) -def test_filesystem_change_after_native_proof_prevents_cleanup( - tmp_path: Path, monkeypatch: pytest.MonkeyPatch, mutation: str, -) -> None: - registry, state, runtime = _fixture(tmp_path) - _record_todo_write(registry, state, runtime, "Commit with a local concurrent change") - actual = adapter.effect_runtime_result - directory = outbox.partition_directory(runtime, "goal-e2e", "todos") - after_change: dict[str, bytes] = {} - - def rpc(method: str, request: dict, **kwargs: Any) -> dict: - result = actual(method, request, **kwargs) - if method == PLAN and request["acknowledgement"] is not None: - if mutation == "cursor": - outbox.write_cursor(directory, partition="todos", **result["cursor_update"]) - elif mutation == "entry": - path = next(directory.glob("*.prepared.json")) - path.write_bytes(path.read_bytes() + b"\n") - else: - _record_todo_write(registry, state, runtime, "Writer arrived after native proof") - after_change.update(inventory(directory)) - return result - - monkeypatch.setattr(adapter, "effect_runtime_result", rpc) - result = _drain(registry, runtime) - assert result.ok is False - # Formatting-only mutation is caught by byte revalidation even when dataclass - # equality sees the same decoded entry. Cursor/inventory changes fail earlier. - assert result.reason_code == "outbox_file_changed" - assert result.candidate_readback_verified is True - assert result.reclaimed_residue == 0 - assert inventory(directory) == after_change - - -def test_forged_commit_ack_keeps_residue_for_a_later_exact_recovery( - tmp_path: Path, monkeypatch: pytest.MonkeyPatch, -) -> None: - registry, state, runtime = _fixture(tmp_path) - _record_todo_write(registry, state, runtime, "Recover the committed transaction, not its corrupted ACK") - actual = adapter.effect_runtime_result - - def rpc(method: str, request: dict, **kwargs: Any) -> dict: - result = actual(method, request, **kwargs) - return {**result, "provider_revision": "foreign-revision"} if method == COMMIT else result + assert result.ok and result.delivered == 3 + assert len(seen) == 1 + method, request, response, kwargs = seen[0] + assert method == "coordination.runtime_shadow.drain" + assert kwargs["retry_safe"] is False + assert "projection" not in request and "entries" not in request + assert "proof" not in response and "transactions" not in response and "head" not in response - monkeypatch.setattr(adapter, "effect_runtime_result", rpc) - stopped = _drain(registry, runtime) - assert not stopped.ok and stopped.reason_code == "shadow_commit_entry_result_invalid" - assert stopped.reclaimed_residue == 0 - monkeypatch.setattr(adapter, "effect_runtime_result", actual) - recovered = _drain(registry, runtime) - assert recovered.ok and recovered.replayed == 1 and recovered.delivered == 0 - -def test_budget_expiring_during_native_read_cannot_authorize_late_cleanup( +def test_lost_batch_response_reports_unknown_and_explicit_retry_reads_durable_facts( tmp_path: Path, monkeypatch: pytest.MonkeyPatch, ) -> None: registry, state, runtime = _fixture(tmp_path) - _record_todo_write(registry, state, runtime, "Preserve evidence after time expires") + _record_todo_write(registry, state, runtime, "Retain this committed write") actual = adapter.effect_runtime_result - expired = False - can_reclaim = adapter._DrainBudget.can_reclaim - def rpc(method: str, request: dict, **kwargs: Any) -> dict: - nonlocal expired - result = actual(method, request, **kwargs) - if method == PLAN and request["acknowledgement"] is not None: - expired = True - return result + def lose_response(method: str, request: dict, **kwargs: Any) -> dict: + actual(method, request, **kwargs) + raise TimeoutError("lost response after successful native drain") - def budget(self: Any, count: int) -> bool: - return not expired and can_reclaim(self, count) - - monkeypatch.setattr(adapter, "effect_runtime_result", rpc) - monkeypatch.setattr(adapter._DrainBudget, "can_reclaim", budget) - result = _drain(registry, runtime) - directory = outbox.partition_directory(runtime, "goal-e2e", "todos") - assert result.ok and result.delivered == 1 - assert result.budget_exhausted and result.reclaimed_residue == 0 - assert len(outbox.list_entries(directory)) == 1 - assert outbox.read_cursor(directory) is None - expired = False + monkeypatch.setattr(adapter, "effect_runtime_result", lose_response) + failed = _drain(registry, runtime) + assert not failed.ok and failed.reason_code == "shadow_drain_outcome_unknown" monkeypatch.setattr(adapter, "effect_runtime_result", actual) recovered = _drain(registry, runtime) - assert recovered.ok and recovered.replayed == 1 + assert recovered.ok and recovered.outcome == "nothing_pending" + view = adapter.read_local_authority_shadow(runtime_root=runtime, goal_id="goal-e2e", scan_limit=10) + assert len(view["proof"]["transactions"]) == 2 diff --git a/tests/control_plane_ts/effect_runtime_io.test.ts b/tests/control_plane_ts/effect_runtime_io.test.ts index 8270cef243..81ff34bdf4 100644 --- a/tests/control_plane_ts/effect_runtime_io.test.ts +++ b/tests/control_plane_ts/effect_runtime_io.test.ts @@ -1,4 +1,5 @@ import assert from "node:assert/strict"; +import {spawnSync} from "node:child_process"; import { mkdtemp, open, rm, stat, utimes, writeFile } from "node:fs/promises"; import { tmpdir } from "node:os"; import { join } from "node:path"; @@ -17,6 +18,17 @@ async function workspace(t: TestContext): Promise { return root; } +test("zero-wait acquisition retries a reclaimed dead owner but preserves a live owner", async t => { + const root = await workspace(t), target = join(root, "state"); + const dead = spawnSync(process.execPath, ["-e", "process.exit(0)"]); + assert.equal(dead.status, 0); + await writeFile(`${target}.ts-effect.lock`, JSON.stringify({pid: dead.pid, token: "dead"})); + const acquired = await acquireFileMutationLock(target, process.pid, 0); + await assert.rejects(acquireFileMutationLock(target, process.pid, 0), {code: "mutation_lock_timeout"}); + assert.equal((await mutationLockOwner(target))?.token, acquired.token); + assert.equal(await releaseFileMutationLock(target, acquired.token), true); +}); + test("token-safe release cannot remove a replacement lock", async (t) => { const root = await workspace(t); const target = join(root, "state"); diff --git a/tests/control_plane_ts/shadow_drain.test.ts b/tests/control_plane_ts/shadow_drain.test.ts new file mode 100644 index 0000000000..9a41f16f14 --- /dev/null +++ b/tests/control_plane_ts/shadow_drain.test.ts @@ -0,0 +1,119 @@ +import assert from "node:assert/strict"; +import test from "node:test"; +import {execFileSync} from "node:child_process"; +import {readFile, readdir, writeFile, unlink} from "node:fs/promises"; +import {join} from "node:path"; +import type {JsonObject} from "../../loopx/control_plane/effect_program.ts"; +import {drainShadowOutbox, SHADOW_DRAIN_SCHEMA} from "../../loopx/control_plane/coordination/shadow_drain.ts"; +import {commitLocalAuthorityShadowEntry} from "../../loopx/control_plane/coordination/local_authority_shadow.ts"; +import {fixture, pendingEntry, todo, type ShadowFixture} from "./shadow_file_fixture.ts"; +import {productionScaleCoordinationFixture} from "./production_scale_coordination_fixture.ts"; + +const python = execFileSync("python3", ["-c", "import sys; print(sys.executable)"], {encoding: "utf8"}).trim(); +function request(f: ShadowFixture, extra: JsonObject = {}): JsonObject { + return {schema_version: SHADOW_DRAIN_SCHEMA, runtime_root: f.root, goal_id: "goal-a", python_executable: python, + config_enabled: true, max_entries: 256, budget_seconds: 30, lock_timeout_seconds: 1, ...extra}; +} +const directory = (f: ShadowFixture) => join(f.root, "authority-shadow/outbox/goal-a/todos"); +async function inventory(f: ShadowFixture) { + return Object.fromEntries(await Promise.all((await readdir(directory(f))).sort().map(async name => [name, await readFile(join(directory(f), name), "utf8")]))); +} + +test("complete batch retains mixed Todo metadata and returns no history across transport", async t => { + const f = await fixture(t), mixed = productionScaleCoordinationFixture("goal-a", "native"); + const todos = mixed.projection.todos as JsonObject[]; + await pendingEntry(f, 1, {handoff_mode: "hard_lease", todos}); + const result = await drainShadowOutbox(request(f)); + assert.equal(result.ok, true, JSON.stringify(result)); assert.equal(result.delivered, 1); + assert.equal(result.pending_after, 0); assert.equal(result.reclaimed_residue, 2); + assert.ok(Buffer.byteLength(JSON.stringify(result)) < 4096); + assert.equal(Object.hasOwn(result, "projection"), false); + const loaded = await f.store.loadAuthority(); + assert.equal(loaded.status, "loaded"); + if (loaded.status === "loaded") assert.deepEqual(loaded.head.todos, todos); + assert.equal((await drainShadowOutbox(request(f))).outcome, "nothing_pending"); +}); + +for (const phase of ["after_commit", "after_cursor", "after_unlink"] as const) { + test(`failure at ${phase} replays receipt without another commit`, async t => { + const f = await fixture(t); + await pendingEntry(f, 1, {handoff_mode: "hard_lease", todos: [todo()]}); + let failed = false; + const first = await drainShadowOutbox(request(f), {afterEffect: async p => { + if (p === phase && !failed) {failed = true; throw new Error("injected IO interruption");} + }}); + assert.equal(first.ok, false); assert.equal(failed, true); + const before = await f.store.scanCommitted(null, 10); + const second = await drainShadowOutbox(request(f)); + assert.equal(second.ok, true, JSON.stringify(second)); assert.equal(second.replayed, 1); assert.equal(second.delivered, 0); + assert.deepEqual(await f.store.scanCommitted(null, 10), before); + assert.deepEqual(await readdir(directory(f)), ["drain-cursor.json"]); + }); +} + +for (const mutation of ["entry", "cursor", "new_entry"] as const) { + test(`filesystem ${mutation} changed after proof cannot be cleaned`, async t => { + const f = await fixture(t); + const entry = await pendingEntry(f, 1, {handoff_mode: "hard_lease", todos: [todo()]}); + assert.equal((await commitLocalAuthorityShadowEntry(entry)).outcome, "delivered"); + let changed = false, after: Record = {}; + const result = await drainShadowOutbox(request(f), {afterEffect: async phase => { + if (phase !== "after_proof" || changed) return; + changed = true; + if (mutation === "entry") { + const path = join(directory(f), (await readdir(directory(f))).find(n => n.endsWith("prepared.json"))!); + await writeFile(path, (await readFile(path, "utf8")) + "\n"); + } else if (mutation === "cursor") await writeFile(join(directory(f), "drain-cursor.json"), "{}"); + else await pendingEntry(f, 2, {handoff_mode: "hard_lease", todos: [todo(), todo("second")]}); + after = await inventory(f); + }}); + assert.equal(result.ok, false); assert.equal(result.reclaimed_residue, 0); + assert.deepEqual(await inventory(f), after); + }); +} + +test("deadline expiring in proof keeps recoverable evidence", async t => { + const f = await fixture(t); + const entry = await pendingEntry(f, 1, {handoff_mode: "hard_lease", todos: [todo()]}); + await commitLocalAuthorityShadowEntry(entry); + const before = await inventory(f); + const result = await drainShadowOutbox(request(f, {budget_seconds: .05}), {afterEffect: async phase => { + if (phase === "after_proof") await new Promise(resolve => setTimeout(resolve, 60)); + }}); + assert.equal(result.budget_exhausted, true); assert.equal(result.reclaimed_residue, 0); + assert.deepEqual(await inventory(f), before); + assert.equal((await drainShadowOutbox(request(f))).replayed, 1); +}); + +test("one-entry budget leaves the remaining source writes intact", async t => { + const f = await fixture(t); + const one = {handoff_mode: "hard_lease", todos: [todo()]}; + await pendingEntry(f, 1, one); + await pendingEntry(f, 2, {handoff_mode: "hard_lease", todos: [todo(), todo("two")]}, {previousPartitionProjection: one}); + const first = await drainShadowOutbox(request(f, {max_entries: 1})); + assert.equal(first.delivered, 1, JSON.stringify(first)); assert.equal(first.pending_after, 1); + const second = await drainShadowOutbox(request(f)); + assert.equal(second.delivered, 1); assert.equal(second.pending_after, 0); +}); + +test("a corrupt candidate tail stops ordered delivery and remains repairable", async t => { + const f = await fixture(t), one = {handoff_mode: "hard_lease", todos: [todo()]}; + await pendingEntry(f, 1, one); + await pendingEntry(f, 2, {handoff_mode: "hard_lease", todos: [todo(), todo("two")]}, {previousPartitionProjection: one}); + let count = 0, saved: Buffer | null = null; + const result = await drainShadowOutbox(request(f), {afterEffect: async phase => { + if (phase === "before_commit" && ++count === 2) {saved = await readFile(f.store.path); await writeFile(f.store.path, "{corrupt");} + }}); + assert.equal(result.ok, false); assert.equal(result.delivered, 1); assert.equal(result.pending_after, 1); + assert.ok(saved); await writeFile(f.store.path, saved); + assert.equal((await drainShadowOutbox(request(f))).delivered, 1); +}); + +test("unproved committed-only residue never becomes successful cleanup", async t => { + const f = await fixture(t); + await pendingEntry(f, 1, {handoff_mode: "hard_lease", todos: [todo()]}); + const name = (await readdir(directory(f))).find(n => n.endsWith("prepared.json"))!; + await unlink(join(directory(f), name)); const before = await inventory(f); + assert.equal((await drainShadowOutbox(request(f))).ok, false); + assert.deepEqual(await inventory(f), before); +}); From 8744faf242b3dafd6aeebad0444df0e096e5bbb2 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 27 Sep 2026 13:11:47 +0800 Subject: [PATCH 2/5] docs(authority): reconcile remaining default cutover deliveries Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../2026-09-27-recovery-audit.md | 80 ++++++++++++++++--- .../2026-09-27-recovery-audit.zh-CN.md | 60 ++++++++++++-- .../reviewed-coordination-promotion.md | 27 +++++++ 3 files changed, 147 insertions(+), 20 deletions(-) diff --git a/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-27-recovery-audit.md b/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-27-recovery-audit.md index 1b8f81653e..74747ccb8b 100644 --- a/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-27-recovery-audit.md +++ b/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-27-recovery-audit.md @@ -1,6 +1,6 @@ # Local default cutover: recovery audit and remaining delivery scopes -- Audited baseline: `157ab7b11`, 2026-09-27, plus this delivery. +- Historical recovery baseline: `157ab7b11`; current inventory: `70b3cca01`, 2026-09-27. - Owner: overall roadmap #4574 R5/G2; shared authority D2/D3; TS T3/T4. - Supersedes the **current count**, not historical evidence, in the [September 24 reconciliation](2026-09-24-default-cutover-reconciliation.md). @@ -21,7 +21,7 @@ target before the final digest mismatch stopped recovery. There was also no independent read-only CLI proof of a restored store's complete retained history and receipt lookup. These are recovery gaps, not missing capture writers. -## Four scoped deliveries starting with this PR +## Historical recovery delivery allocation | Delivery | Observable result and remaining boundary | | --- | --- | @@ -30,7 +30,7 @@ and receipt lookup. These are recovery gaps, not missing capture writers. | **3. Whole-Goal activation and rollback integration** | Reconcile #5054's retained-source inventory, then exercise source drain, saved reviewed cutover, all retained command consumers and fenced recovery/rollback together. Bind a recovered copy through an explicit transition; do not revive a source lease or overwrite later writes. Delete only Python decisions whose callers have actually moved. | | **4. Default entrypoints and final bounded retirement** | New Goal creation, settings, installation, packaged frontend/Lark/CLI consistently use the qualified local profile. Existing Goals have explicit migration and disable/recovery paths. Remove last legacy business writers after their caller inventory and rollback constraints pass; retain rendering and Host IO. | -This is **four planned new delivery PRs including this one, three afterwards**, +At that recovery checkpoint, this was **four planned new delivery PRs including the recovery PR, three afterwards**, not a guarantee that no acceptance defect will require another PR. The original three *architectural packages* are not a decrementing PR counter. This PR closes one named recovery slice inside package 2; it does not close all of package 2. @@ -97,17 +97,73 @@ PostgreSQL 16 server passed, with no skipped checks in these suites. A detached audit. Recovery of the earlier timed-out SQLite destination passed without reissuing its committed operations. These checks do not claim active cutover. +## Current delivery inventory and native drain (`70b3cca01`) + +The current plan contains **seven delivery slots including this change**: four +already-open PRs and three scoped deliveries. It does not mean seven new PRs, +nor guarantee that seven merges suffice. Earlier counts treated whole-Goal +integration as a single PR before its recovery and execution gaps were bounded; +that was an architectural grouping, not a reliable PR commitment. + +| Slot | Existing work / observable completion | +| --- | --- | +| 1 | **#5173**, open: reviewed File↔SQLite selector/fence cutover, durable backup, full-history audit, recovery and retry. Integrate it; do not rebuild it. | +| 2 | **#5144**, open: managed Host execution lifetime/lease supervision. Attached Hosts still need an explicit cancellation boundary. | +| 3 | **#5054**, open: retire legacy Todo event projection/backfill/completion and isolate the experimental supervisor log. | +| 4 | **#4931**, open: SQLite retained-proof encoding/read cost; rerun the applicable formal D2 workload instead of equating an optimization with qualification. | +| 5 | **This change**: move the complete bounded source-outbox drain to TS, remove the Python sequencing/proof/cleanup coordinator and the unused per-entry planning RPC. Keep existing durable formats, full receipt verification and the kernel-lock adapter. | +| 6 | **Whole-Goal integration**: combine the accepted slices with retained consumer parity, interrupted cutover/rollback and post-cutover writes. This drain is one completed subitem, not completion of that scope. | +| 7 | **Default entrypoints and bounded Python retirement**: qualify creation/settings/install/frontend/Lark/CLI, migrate existing Goals explicitly, and delete business writers only after their real callers have moved. | + +#5169's content-aware idempotency work is adjacent and must be integrated without +rewriting it; it is not silently counted as another required default-cutover PR. +D1 consumer coverage, D2 capacity/ten-day natural soak and D3 cohort/maintainer +promotion remain **evidence gates**, outside the arithmetic. PostgreSQL service +identity, deployment and operational qualification remain a medium-term scope. + +TS now inventories witnessed source files, invokes the existing receipt planner +and transaction owner, checks the monotonic budget after proof, and performs +cursor/cleanup effects under M → primary marker → kernel-lock exclusion. The +Python facade sends one bounded request with no source projection/history and +does not automatically retry a lost response. `shadow_drain_outcome_unknown` +requires a later explicit receipt-based drain; it never asserts no commit. +Unconfigured Goals with no capture state still make no drain RPC. + +Real process-death validation found a shared lock defect: a zero-wait acquire +reclaimed a dead owner but returned timeout before trying the now-free path. +It now allows one immediate retry after proven reclamation; live owners still +reject immediately. This uses the shared lock owner rather than a drain-only +sleep or increased timeout. The two crash harnesses now share one scheduling +fixture and kill/reap the actual TS owner at the durable boundary. + +The runtime shadow remains a **File candidate**, not a promoted authority or a +new SQLite shadow provider. SQLite remains a supported canonical promotion and +archive-restore target. No provider selector, registry or active Goal is changed +by this delivery. The existing CLI and inline writer drain entrypoints adopt the +same owner; no frontend/Lark configuration contract changes. + +Validation for native drain: the new full-batch tests exercise the existing +mixed production-scale Todo fixture; real CLI SIGKILL, filesystem permission, +cursor tampering and File/SQLite reviewed-promotion tests cover the persistence +boundaries. Shared-store regression passes against isolated PostgreSQL 16. +A detached authorized snapshot supplies 1,101 complete Todo records; three +explicitly synthetic source writes produce a fresh four-transaction candidate. +Every original Todo JSON record survives drain and SQLite/File archive restore. +This is not replay of that snapshot's old transaction history or a live migration. + +On three local three-entry trials, median drain time is 1.34 s on the audited +base and 0.41 s here, with 11 facade RPCs reduced to one. The kernel-lock process +is lazy and reused within a batch; it releases locks between sections. These +small warm-runtime measurements are not D2 p95/capacity claims. The existing +CLI output-budget suite fails on both base and head at 14,514 versus 14,500 +characters for the crowded Turn JSON packet. No ceiling is raised; merge +qualification retains that failure rather than declaring all checks green. + ## Shared-runtime latency reconciliation (`96a3b90f4`) -The recovery delivery above is now merged as #5140. #5144 is the open managed -Host process supervision slice; attached hosts still need their declared -cancellation boundary. Whole-Goal activation/rollback and default entrypoint -cutover remain the two planned subsequent implementation PRs. Existing #5054 -(retirement) and #4931 (SQLite proof encoding) remain open and are not new work. -Thus the inventory is two planned implementation PRs plus those three existing -PRs, **before this newly reproduced latency repair**. This is an inventory, not -an unconditional completion count: #4224 D2 capacity/soak and D1/D3 evidence -remain gates, and failures may require additional scoped fixes. +The recovery delivery above merged as #5140. The subsequent latency repair is +also on the current audited main. Its observations below remain historical +validation, not another open delivery or a D2 qualification. The latency repair does not retire another Python owner or close D2. Isolated fixed File snapshots reproduce 9.9–10.5 second cold history verification, diff --git a/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-27-recovery-audit.zh-CN.md b/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-27-recovery-audit.zh-CN.md index 6c2fbf7705..93e9c8c67d 100644 --- a/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-27-recovery-audit.zh-CN.md +++ b/docs/architecture/rfcs/ledger/shared-goal-authority-state-provider-v0/2026-09-27-recovery-audit.zh-CN.md @@ -17,7 +17,7 @@ checkpoint/delta 格式升级和有界 Python 原型退役已在 main。尤其 # 内容可能先写入隔离目标,直到最终摘要不符才失败。此外缺少独立只读 CLI,证明恢复 目标的完整历史及回执查询仍正确。这是恢复缺口,不是尚未实现事件捕获。 -## 从本次开始的四个交付范围 +## 当时的恢复交付拆分(历史记录) | 交付 | 可观察结果及剩余边界 | | --- | --- | @@ -26,7 +26,7 @@ checkpoint/delta 格式升级和有界 Python 原型退役已在 main。尤其 # | **3. 整 Goal 激活及回退集成** | 对齐 #5054 的保留来源清单,联合验证来源 drain、保存的 reviewed cutover、全部保留命令消费者及 fenced recovery/rollback。通过明确转换绑定恢复副本,不复活旧租约,不覆盖后来写入。仅删除 caller 已迁走的 Python 决策。 | | **4. 默认入口与最后一批有界退役** | 新建 Goal、settings、安装及打包 frontend/Lark/CLI 一致使用合格本地 profile;存量 Goal 有明确迁移及停用/恢复路径。caller 清单与回退约束通过后,删除最后的旧业务 writer,保留渲染及 Host IO。 | -这是**包含本次在内四个规划新 PR,本次交付后剩三个**;不保证验收不会再发现需要 +当时的恢复检查点规划为**包含恢复 PR 在内四个新 PR,交付后剩三个**;不保证验收不会再发现需要 修复的缺陷。原来的三个“架构工作包”不是倒计时 PR 数。本次关闭第 2 包中的一个 具名恢复切片,没有把整个第 2 包标为完成。后续必须指出哪行真正交付,不能再重复 一个不变的“5–8”。 @@ -75,14 +75,58 @@ PostgreSQL 16 的 4 项跨 provider 检查全部通过,这些套件没有跳 独立真实来源快照通过 File、SQLite 恢复及独立审计;此前超时的 SQLite 目标也成功 恢复,没有重发已提交操作。这些结果不表示已经完成活跃 Goal 切换。 +## 当前交付清单与原生 drain(`70b3cca01`) + +当前规划有 **7 个交付槽位,包含本次:4 个已有开放 PR,加 3 个交付范围**。 +这不意味着还要新开 7 个 PR,也不能保证合并 7 个就足够。此前把整 Goal 整合 +当成一个 PR,但恢复与执行边界尚未拆清;那是架构工作包,不能作为准确倒计时。 + +| 项目 | 已有工作与完成标准 | +| --- | --- | +| 1 | **#5173**:受审 File↔SQLite selector/fence 切换、备份、完整历史核对、恢复和重试。整合已有实现,不重写。 | +| 2 | **#5144**:受管 Host 执行生命周期与租约监督;attached Host 仍需明确取消能力边界。 | +| 3 | **#5054**:退役旧 Todo 事件投影、回填与 completion 分支,独立实验 supervisor 日志。 | +| 4 | **#4931**:SQLite 历史证明编码与读取成本;需正式 D2 工作负载复验,优化代码不等于资格通过。 | +| 5 | **本次**:完整有界 source-outbox drain 归 TS,删除 Python 的顺序、证明、清理编排和无人调用的逐条规划 RPC;保留持久格式、回执证明、内核锁适配。 | +| 6 | **整 Goal 整合**:接通已交付切片,验证保留消费者、中断切换/回退及切换后的新写入。本次 drain 只完成其中一个子项。 | +| 7 | **默认入口与有界 Python 退役**:创建、设置、安装、前端、Lark、CLI 一致;已有 Goal 明确迁移;真实调用者迁走后再删业务 writer。 | + +#5169 的内容感知幂等是相邻工作,应避免重写,不把它悄悄加成另一个默认切换必需 PR。 +D1 消费者覆盖、D2 容量与至少十天自然 soak、D3 队列与维护者晋升仍是独立证据门。 +PostgreSQL 服务身份、部署和运维资格属于中期范围。 + +TS 统一读取来源文件见证、调用既有回执 planner 与事务 owner,在证明后重新检查 +单调时钟预算,并在 M → primary marker → kernel lock 的顺序下更新游标和清理。 +Python 只发送一次有界请求,不搬完整 projection/history,也不自动重试响应丢失的 +批次。`shadow_drain_outcome_unknown` 表示必须由下次显式 drain 读取回执恢复, +不能解释为“没有提交”。未启用且不存在捕获状态的 Goal 仍不产生 drain RPC。 + +真正杀进程的测试发现共享锁缺陷:零等待获取已经回收死进程锁,却立即报超时。 +现在仅在确认回收后立即再尝试一次;活进程持锁时仍马上退出。这是共享锁修复, +没有为 drain 增加 sleep 或放宽超时。两套崩溃 harness 也共用一个调度 fixture, +在持久化边界杀掉并回收真正执行决策的 TS 进程。 + +runtime shadow 仍是 **File 候选存储**,不是正式 authority,也没有新增 SQLite +shadow provider。SQLite 是 canonical 晋升和归档恢复目标。本次不切换任何活跃 +Goal、provider selector 或 registry。原有 CLI 与 writer 内联 drain 采用同一 owner; +没有前端/Lark 配置合同变化。 + +本次验证复用混合 production-scale Todo fixture;真实 CLI 的 SIGKILL、文件权限、 +游标篡改、File/SQLite 受审晋升覆盖持久边界;共享存储在隔离 PostgreSQL 16 上回归。 +授权隔离快照提供 1,101 个完整 Todo,三笔明确标记的合成来源写入形成新的四事务 +候选历史。原始 Todo JSON 经 drain 及 SQLite/File 归档恢复仍完整相等;这不是 +重放该快照的旧事务历史,也不是活跃 Goal 迁移。 + +本机三个事务的三次测量中位数由基线 1.34 秒降至本次 0.41 秒,facade RPC 从 +11 次降为 1 次。内核锁进程按需创建、批次内复用,各临界区之间释放锁。这是小 +负载热运行时测量,不是 D2 p95/容量结论。CLI 输出预算检查在 base/head 都因 +拥挤 Turn JSON 为 14,514 字符、超过 14,500 上限而失败;未提高预算,合并资格 +保留此失败,不能报告所有检查为绿。 + ## 共享运行时延迟核对(`96a3b90f4`) -上文恢复交付已作为 #5140 合入。#5144 是待合并的受管 Host 进程监督切片;attached -Host 仍需明确取消能力边界。整 Goal 激活/回退、默认入口切换仍是后续两个规划实现 -PR。已有 #5054(退役)、#4931(SQLite 证明编码)仍开放,不重复实现。因此清单是 -**两个规划实现 PR,加三个已有 PR,再加本次新复现的延迟修复**。这是工作清单, -不是无条件完成倒计时:#4224 D2 容量/soak、D1/D3 证据仍需验收;失败可产生新的 -有界修复,必须指出具体缺陷,不能重新复述一个固定区间。 +上文恢复交付 #5140,以及后续共享运行时延迟修复,都已进入当前核对的 main。 +下文保留当时的验证事实,不再把它们算作开放工作,也不将其当作 D2 资格。 本次不退役额外 Python owner,也不关闭 D2。隔离固定 File 快照复现了 9.9–10.5 秒 的冷历史校验,期间 ping 在原 10 秒预算内超时。热缓存掩盖问题,交替读取 Goal 又 diff --git a/docs/reference/reviewed-coordination-promotion.md b/docs/reference/reviewed-coordination-promotion.md index e1dc807cf9..fa63f9477a 100644 --- a/docs/reference/reviewed-coordination-promotion.md +++ b/docs/reference/reviewed-coordination-promotion.md @@ -166,3 +166,30 @@ lost acknowledgement after commit. These are synthetic fault injections, not a claim of arbitrary process-death or elapsed-soak coverage. Real local rehearsals must use read-only captured sources and disposable copies; never promote an active Goal merely to validate this refactor. + +## Drain before promotion + +`authority-shadow drain` and post-write inline drains use one bounded TS batch. +It verifies the complete retained lineage and exact source byte witnesses before +advancing a cursor or reclaiming outbox files. A cursor is a checkpoint, not +permission to delete; Todo metadata and original transaction receipts remain +part of the existing complete record contract. + +```bash +loopx --format json authority-shadow drain --goal-id example-goal \ + --max-entries 64 --budget-seconds 30 +``` + +`budget_exhausted=true` with pending entries means another bounded invocation is +needed. The budget controls admission of further effects; it does not cancel an +in-flight durable commit or replace the transport timeout. A live primary writer returns `primary_writer_busy` without waiting for +it. `shadow_drain_outcome_unknown` means the transport lost a trustworthy batch +result: inspect status and invoke drain again explicitly. The next invocation +uses persisted receipts; it must not recreate the business write, delete the +outbox, or assume that no candidate commit happened. Other proof failures require +repairing their reported cause before continuing. + +The runtime shadow is still a File candidate. Draining it does not select a +canonical provider, promote a Goal, bypass qualification, or create a SQLite +shadow. File/SQLite promotion uses the reviewed journey above. Feature-off +writers with no prior capture state do not start the drain runtime. From 47b29f8ec1f056fa0d634e1e0d97df2c54bea309 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 27 Sep 2026 23:16:23 +0800 Subject: [PATCH 3/5] test(authority): pause the real native drain at its commit boundary The management interleaving tests hooked the retired per-entry coordination.runtime_shadow.commit_entry RPC, which the batch drain no longer crosses, so the public writer finished without ever reaching the barrier. Route the same public write through the private native fault driver instead: it runs the production drainShadowOutbox, stops at the real before/after commit phase and resumes when the test releases it, so rollback, rebootstrap and the late commit/cursor stay real. Capture lineage and provider revision now come from durable readback rather than the retired RPC payload, and an explicit retry of the superseded batch still cannot deliver. The fault driver only gains an optional release path; without it the barrier remains terminal for the crash-window callers. The native drain test also reuses resolveTestPython instead of discovering a bare python3, matching the existing guard. Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../testing/shadow_drain_fault_process.ts | 8 ++- .../test_shadow_management_e2e.py | 69 +++++++++++++++---- tests/control_plane_ts/shadow_drain.test.ts | 4 +- 3 files changed, 64 insertions(+), 17 deletions(-) diff --git a/loopx/control_plane/testing/shadow_drain_fault_process.ts b/loopx/control_plane/testing/shadow_drain_fault_process.ts index 06b7164d65..d6a7e3ca62 100644 --- a/loopx/control_plane/testing/shadow_drain_fault_process.ts +++ b/loopx/control_plane/testing/shadow_drain_fault_process.ts @@ -1,13 +1,19 @@ /** Private-process test driver, never installed as an RPC effect. */ +import {existsSync} from "node:fs"; import {drainShadowOutbox} from "../coordination/shadow_drain.ts"; let input = ""; for await (const bytes of process.stdin) input += bytes; const request = JSON.parse(input); const phase = process.argv[2] === "between_unlinks" ? "after_unlink" : process.argv[2]; +// Without a release path the observed barrier is terminal (a crash window). +// With one, the same real batch continues after the caller resumes it, which is +// what a management interleaving needs. +const release = process.argv[3]; const result = await drainShadowOutbox(request, {afterEffect: async observed => { if (observed === phase) { process.stdout.write(`BARRIER ${JSON.stringify({native_pid: process.pid, request})}\n`); - await new Promise(() => {setInterval(() => {}, 1000);}); + if (release === undefined) await new Promise(() => {setInterval(() => {}, 1000);}); + else while (!existsSync(release)) await new Promise(resolve => {setTimeout(resolve, 10);}); } }}); process.stdout.write(JSON.stringify(result) + "\n"); diff --git a/tests/control_plane/test_shadow_management_e2e.py b/tests/control_plane/test_shadow_management_e2e.py index 8ab890dd48..7491afec69 100644 --- a/tests/control_plane/test_shadow_management_e2e.py +++ b/tests/control_plane/test_shadow_management_e2e.py @@ -102,23 +102,45 @@ def test_public_rollback_preserves_other_goal_and_replays_after_primary_changes( assert (runtime / "authority-shadow" / "file-v0" / "store-identity").read_bytes() == identity +_DRAIN_PHASE = {"before": "before_commit", "after": "after_commit"} + +# The public write keeps its own inline drain, but the pause now happens inside +# the real native batch: the private fault driver runs the same +# ``drainShadowOutbox`` the production RPC handler runs and stops at the actual +# commit boundary. Releasing it lets that same batch continue, so a commit or a +# cursor update attempted after the generation changed is still real. _DELAYED_WRITER = """ -import json, pathlib, sys, time +import json, pathlib, subprocess, sys, time from loopx.control_plane.coordination import local_authority_shadow_adapter as adapter from loopx.cli import main -barrier, release, timing = pathlib.Path(sys.argv[1]), pathlib.Path(sys.argv[2]), sys.argv[3] +barrier, release, phase = pathlib.Path(sys.argv[1]), pathlib.Path(sys.argv[2]), sys.argv[3] actual = adapter.effect_runtime_result def delayed(method, request, **kwargs): - if method != 'coordination.runtime_shadow.commit_entry': + if method != 'coordination.runtime_shadow.drain': return actual(method, request, **kwargs) - result = actual(method, request, **kwargs) if timing == 'after' else None - barrier.write_text(json.dumps({'request':request, 'result':result})) + child = subprocess.Popen( + ['node', '--no-warnings', '--experimental-strip-types', + 'loopx/control_plane/testing/shadow_drain_fault_process.ts', phase, str(release)], + stdin=subprocess.PIPE, stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True, + ) + child.stdin.write(json.dumps(request)) + child.stdin.close() + for line in child.stdout: + if line.startswith('BARRIER '): + barrier.write_text(line.removeprefix('BARRIER ')) + break + else: + raise RuntimeError('native barrier missing: ' + child.stderr.read()) deadline = time.monotonic() + 60 while not release.exists(): if time.monotonic() > deadline: + child.kill(); child.wait() raise RuntimeError('test scheduling barrier timed out') time.sleep(.01) - return result if timing == 'after' else actual(method, request, **kwargs) + remaining, errors = child.stdout.read(), child.stderr.read() + if child.wait() != 0: + raise RuntimeError(errors or remaining) + return json.loads(remaining.strip().splitlines()[-1]) adapter.effect_runtime_result = delayed raise SystemExit(main(sys.argv[4:])) """ @@ -130,7 +152,7 @@ def _paused_writer(tmp_path: Path, registry: Path, runtime: Path, timing: str) - args = _arguments(registry, runtime, "todo", "add", "--goal-id", "goal-a", "--role", "agent", "--text", "Transaction across a management boundary", "--claimed-by", "agent-a") child = subprocess.Popen( - [sys.executable, "-c", _DELAYED_WRITER, str(barrier), str(release), timing, *args], + [sys.executable, "-c", _DELAYED_WRITER, str(barrier), str(release), _DRAIN_PHASE[timing], *args], cwd=REPO_ROOT, stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True, ) deadline = time.monotonic() + 20 @@ -174,8 +196,19 @@ def test_late_real_commit_cannot_cross_rollback_and_rebootstrap(tmp_path: Path, child, release, barrier = _paused_writer(tmp_path, registry, runtime, timing) try: request = barrier["request"] - assert request["capture_lineage_id"] == first["capture_lineage_id"] - revision = first["provider_revision"] if timing == "before" else barrier["result"]["provider_revision"] + assert request["schema_version"] == "loopx_shadow_drain_v0" + assert request["goal_id"] == "goal-a" + assert request["runtime_root"] == str(runtime.resolve()) + # The paused batch still runs under the generation that is about to be + # superseded, and the revision it targets is durable readback. + management = read_shadow_management_state(runtime, "goal-a") + assert management is not None and management["status"] == "active" + assert management["binding"]["capture_lineage_id"] == first["capture_lineage_id"] + revision = json.loads(_candidate(runtime, "goal-a").read_text())["provider_revision"] + if timing == "before": + assert revision == first["provider_revision"] + else: + assert revision != first["provider_revision"] rollback = _cli(registry, runtime, "coordination-shadow", "rollback", "--goal-id", "goal-a", "--provider-revision", revision, "--execute")["rollback"] archived = Path(rollback["outbox_archive_path"]) @@ -184,15 +217,21 @@ def test_late_real_commit_cannot_cross_rollback_and_rebootstrap(tmp_path: Path, second = _bootstrap(registry, runtime) assert second["capture_lineage_id"] != first["capture_lineage_id"] candidate = _candidate(runtime, "goal-a").read_bytes() - _release(child, release) + payload = _release(child, release) + # The resumed batch reports this write from its own drain evidence; it + # must never claim a delivery into the generation that replaced it. + assert payload["coordination_runtime_shadow"]["outcome"] != "delivered" assert _candidate(runtime, "goal-a").read_bytes() == candidate assert _candidate(runtime, "goal-b").read_bytes() == other assert {str(path.relative_to(archived)): path.read_bytes() for path in archived.rglob("*") if path.is_file()} == retained active_outbox = runtime / "authority-shadow" / "outbox" / "goal-a" assert not list(active_outbox.rglob("drain-cursor.json")) assert not list(active_outbox.rglob("*.prepared.json")) - late = effect_runtime_result("coordination.runtime_shadow.commit_entry", request) + # An explicit retry of the superseded batch reads durable receipts and + # still cannot deliver into the generation that replaced it. + late = effect_runtime_result("coordination.runtime_shadow.drain", request) assert late["outcome"] not in {"delivered", "replayed", "reconciled"} + assert late["delivered"] == 0 and late["replayed"] == 0 assert _candidate(runtime, "goal-a").read_bytes() == candidate finally: if child.poll() is None: @@ -202,14 +241,16 @@ def test_late_real_commit_cannot_cross_rollback_and_rebootstrap(tmp_path: Path, def test_corrupt_history_after_real_commit_cannot_authorize_cursor_cleanup(tmp_path: Path) -> None: registry, runtime = _workspace(tmp_path) - _bootstrap(registry, runtime) + first = _bootstrap(registry, runtime) child, release, barrier = _paused_writer(tmp_path, registry, runtime, "after") try: - assert barrier["result"]["outcome"] in {"delivered", "replayed", "reconciled"} + # The commit already happened: the durable candidate moved past bootstrap. + assert barrier["request"]["goal_id"] == "goal-a" + candidate = _candidate(runtime, "goal-a") + assert json.loads(candidate.read_text())["provider_revision"] != first["provider_revision"] directory = runtime / "authority-shadow" / "outbox" / "goal-a" / "todos" entries = {path.name: path.read_bytes() for path in directory.glob("*.json")} assert any(name.endswith(".prepared.json") for name in entries) - candidate = _candidate(runtime, "goal-a") record = json.loads(candidate.read_text()) record["committed"][0]["provider_revision"] = "file:1:" + "0" * 24 candidate.write_text(json.dumps(record)) diff --git a/tests/control_plane_ts/shadow_drain.test.ts b/tests/control_plane_ts/shadow_drain.test.ts index 9a41f16f14..40618069cd 100644 --- a/tests/control_plane_ts/shadow_drain.test.ts +++ b/tests/control_plane_ts/shadow_drain.test.ts @@ -1,6 +1,5 @@ import assert from "node:assert/strict"; import test from "node:test"; -import {execFileSync} from "node:child_process"; import {readFile, readdir, writeFile, unlink} from "node:fs/promises"; import {join} from "node:path"; import type {JsonObject} from "../../loopx/control_plane/effect_program.ts"; @@ -8,8 +7,9 @@ import {drainShadowOutbox, SHADOW_DRAIN_SCHEMA} from "../../loopx/control_plane/ import {commitLocalAuthorityShadowEntry} from "../../loopx/control_plane/coordination/local_authority_shadow.ts"; import {fixture, pendingEntry, todo, type ShadowFixture} from "./shadow_file_fixture.ts"; import {productionScaleCoordinationFixture} from "./production_scale_coordination_fixture.ts"; +import {resolveTestPython} from "../../scripts/test-python.mjs"; -const python = execFileSync("python3", ["-c", "import sys; print(sys.executable)"], {encoding: "utf8"}).trim(); +const python = resolveTestPython(); function request(f: ShadowFixture, extra: JsonObject = {}): JsonObject { return {schema_version: SHADOW_DRAIN_SCHEMA, runtime_root: f.root, goal_id: "goal-a", python_executable: python, config_enabled: true, max_entries: 256, budget_seconds: 30, lock_timeout_seconds: 1, ...extra}; From 37c61125129e9a68e2dbadae87fd741c59beeb52 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 27 Sep 2026 23:47:40 +0800 Subject: [PATCH 4/5] test(authority): point the remaining stage2c seams at the batch owner Two Stage 2C lanes still addressed the retired Python per-entry commit. The bounded e2e built its witnessed selection through adapter._commit_entry_request, which the batch drain deleted, so the test raised AttributeError before exercising the rejection; it now builds the same request from durable entry bytes and keeps asserting that a native request cannot relabel a committed primary as abandoned. The mutation lane's replay_counted_as_delivery case edited the removed adapter counter and failed with mutation locator drift; the equivalent owner is the batch replay loop, so the case now mutates result.replayed++ in shadow_drain.ts and is still killed by the s2c2.sigkill_mid_drain ladder row. Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- examples/shared-goal-authority-e2e/mutants.py | 15 ++++++--------- .../test_runtime_shadow_bounded_e2e.py | 19 ++++++++++++++++++- 2 files changed, 24 insertions(+), 10 deletions(-) diff --git a/examples/shared-goal-authority-e2e/mutants.py b/examples/shared-goal-authority-e2e/mutants.py index d74a2b7719..93719f41fd 100644 --- a/examples/shared-goal-authority-e2e/mutants.py +++ b/examples/shared-goal-authority-e2e/mutants.py @@ -181,9 +181,9 @@ def command(self) -> list[str]: "const matched = localAuthorityShadowHeadDigest(request.projection) === localAuthorityShadowHeadDigest(lineage.head.head);", "const matched = true;")),), LADDER_ROW + "[s2c2.parity_divergent_detects_foreign_edit]"), - Case("replay_counted_as_delivery", ((COORDINATION + "local_authority_shadow_adapter.py", replacement( - " self._result.replayed += 1\n", - " self._result.delivered += 1\n")),), + Case("replay_counted_as_delivery", ((COORDINATION + "shadow_drain.ts", replacement( + "result.no_op += Number(noOp === true); result.entries.push(summary); result.replayed++; consumed++;", + "result.no_op += Number(noOp === true); result.entries.push(summary); result.delivered++; consumed++;")),), LADDER_ROW + "[s2c2.sigkill_mid_drain]"), ]) @@ -252,12 +252,9 @@ def apply(source: str) -> str: " resolved_source = state_file.resolve(strict=False)", " return # DELIBERATE MUTANT: allow another goal to bypass source authority.\n resolved_source = state_file.resolve(strict=False)")),), "tests/control_plane/test_shadow_writer_variant_e2e.py::test_other_goal_cannot_write_a_protected_goal_source_via_state_override[active_capture]"), - Case("cleanup_hides_verified_commit", ((COORDINATION + "local_authority_shadow_adapter.py", replacement( - " if self._result.cursor_before is None:\n" - " self._result.cursor_before = view.get(\"cursor\")\n" - " self._record_view(view)\n", - " if self._result.cursor_before is None:\n" - " self._result.cursor_before = view.get(\"cursor\")\n")),), + Case("cleanup_hides_verified_commit", ((COORDINATION + "shadow_drain.ts", replacement( + " if (plan.view !== null && typeof plan.view === \"object\") observe(plan.view as JsonObject);\n", + "")),), "tests/control_plane/test_shadow_drain_adversarial.py::test_cleanup_permission_failure_reports_verified_commit_and_recovers[before_commit]"), Case("native_update_maintenance", ((COORDINATION + "local_authority_runtime.ts", remove_native_update_maintenance),), diff --git a/tests/control_plane/test_runtime_shadow_bounded_e2e.py b/tests/control_plane/test_runtime_shadow_bounded_e2e.py index 976cfc36ac..3ba16109dc 100644 --- a/tests/control_plane/test_runtime_shadow_bounded_e2e.py +++ b/tests/control_plane/test_runtime_shadow_bounded_e2e.py @@ -16,6 +16,7 @@ from loopx.control_plane.effect_runtime import EffectRuntimeRejected from loopx.control_plane.coordination.coordination_state_contract_generated import ( + LOCAL_AUTHORITY_SHADOW_COMMIT_ENTRY_REQUEST_SCHEMA, TASK_LEASE_ACQUIRE_REQUEST_SCHEMA, ) from loopx.control_plane.coordination.runtime_shadow import ( @@ -513,7 +514,23 @@ def test_public_committed_primary_cannot_be_relabelled_abandoned_by_native_reque w.crash("before_commit", "todo", "add", "--role", "agent", "--text", "A committed primary is never abandoned") directory = outbox.partition_directory(w.runtime, w.goal, "todos") [entry] = outbox.list_entries(directory) - request = adapter._commit_entry_request(runtime_root=w.runtime, goal_id=w.goal, entry=entry) + # The batch drain owns public commits now, so the witnessed selection is + # built here from durable entry bytes instead of a retired private helper. + request = { + "schema_version": LOCAL_AUTHORITY_SHADOW_COMMIT_ENTRY_REQUEST_SCHEMA, + "runtime_root": str(w.runtime), + "goal_id": w.goal, + "entry_id": entry.entry_id, + "partition": entry.partition, + "seq": entry.seq, + "capture_lineage_id": entry.prepared.get("capture_lineage_id"), + "prepared_sha256": outbox.raw_bytes_digest(entry.prepared_path.read_bytes()), + "committed_sha256": ( + outbox.raw_bytes_digest(entry.committed_path.read_bytes()) + if entry.committed_path + else None + ), + } assert request["committed_sha256"] is not None request["resolution"] = "abandoned" before = {path.name: path.read_bytes() for path in directory.iterdir()} From aea02b9bb6be97737f84236a5d1f749e06273b47 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Mon, 28 Sep 2026 00:37:37 +0800 Subject: [PATCH 5/5] test(authority): hold the cross-runtime marker in the deferred-drain window The native batch drain takes the TypeScript mutation marker, so the stage2c2 deferred-drain window has to acquire the same cross-runtime lock the production readers take; a kernel-only flock no longer excludes the batch and let the rollback row drain early. Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../testing/authority_e2e_rows_stage2c2.py | 11 ++++++++--- 1 file changed, 8 insertions(+), 3 deletions(-) diff --git a/loopx/control_plane/testing/authority_e2e_rows_stage2c2.py b/loopx/control_plane/testing/authority_e2e_rows_stage2c2.py index c6cc12e364..ced7b1b6ad 100644 --- a/loopx/control_plane/testing/authority_e2e_rows_stage2c2.py +++ b/loopx/control_plane/testing/authority_e2e_rows_stage2c2.py @@ -27,7 +27,7 @@ from dataclasses import dataclass, field from pathlib import Path -from ...file_lock import exclusive_file_lock +from ...file_lock import exclusive_cross_runtime_file_lock from ..coordination import local_authority_shadow_outbox as shadow_outbox from .authority_e2e_fixtures import ( REPO_ROOT, @@ -406,9 +406,14 @@ def todo_count(workspace: GoalWorkspace) -> int: @contextmanager def hold_drain_lock(workspace: GoalWorkspace) -> Iterator[None]: - """Hold the stable maintenance lock so writers defer their post-commit drain.""" + """Hold the maintenance lock so writers defer their post-commit drain. - with exclusive_file_lock( + The native batch takes the TypeScript mutation marker, so the window must + hold the same cross-runtime lock the production readers take; a kernel-only + flock would no longer exclude it. + """ + + with exclusive_cross_runtime_file_lock( shadow_outbox.drain_lock_target(workspace.runtime_root, workspace.goal_id), operation="e2e_window", ):