From 77221351b2ef1e9bc4bd43bdb5bf7345e06c2aa9 Mon Sep 17 00:00:00 2001 From: huan <140176752+HUAN2022A@users.noreply.github.com> Date: Sun, 13 Sep 2026 18:11:52 +0800 Subject: [PATCH] Serialize decision ledger appends through the store lock Concurrent appends each read the ledger, appended one decision, and wrote the whole file back, so the last writer silently dropped every decision appended in between; a sleep and a decide racing could forget an operator's refusal. DecisionLedger now takes the store layout and appends under the store lock, staging the new content and replacing the file in one step. --- packages/core/src/agent_memory/core/ledger.py | 22 +++++++++---- packages/core/src/agent_memory/core/manage.py | 4 +-- tests/system/test_concurrency.py | 31 +++++++++++++++++++ 3 files changed, 49 insertions(+), 8 deletions(-) diff --git a/packages/core/src/agent_memory/core/ledger.py b/packages/core/src/agent_memory/core/ledger.py index f7d1c37b..fbafba3c 100644 --- a/packages/core/src/agent_memory/core/ledger.py +++ b/packages/core/src/agent_memory/core/ledger.py @@ -8,9 +8,11 @@ from __future__ import annotations import dataclasses -import pathlib import re +from .locking import store_lock +from .paths import StoreLayout + LEDGER_FILENAME = "decisions.md" VERDICT_ACCEPTED = "accepted" VERDICT_REJECTED = "rejected" @@ -41,8 +43,9 @@ def render(self) -> str: class DecisionLedger: - def __init__(self, path: pathlib.Path): - self._path = path + def __init__(self, layout: StoreLayout): + self._layout = layout + self._path = layout.dream_reports / LEDGER_FILENAME def decided(self) -> dict[str, Decision]: if not self._path.exists(): @@ -61,7 +64,14 @@ def decided(self) -> dict[str, Decision]: return found def append(self, decision: Decision) -> Decision: - self._path.parent.mkdir(parents=True, exist_ok=True) - head = self._path.read_text(encoding="utf-8") if self._path.exists() else HEADING + "\n\n" - self._path.write_text(head + decision.render() + "\n", encoding="utf-8") + with store_lock(self._layout): + self._path.parent.mkdir(parents=True, exist_ok=True) + head = ( + self._path.read_text(encoding="utf-8") + if self._path.exists() + else HEADING + "\n\n" + ) + staged = self._path.with_name(self._path.name + ".pending") + staged.write_text(head + decision.render() + "\n", encoding="utf-8") + staged.replace(self._path) return decision diff --git a/packages/core/src/agent_memory/core/manage.py b/packages/core/src/agent_memory/core/manage.py index 0f444ebd..2684cc43 100644 --- a/packages/core/src/agent_memory/core/manage.py +++ b/packages/core/src/agent_memory/core/manage.py @@ -23,7 +23,7 @@ from .clock import Clock from .database import Database from .errors import FieldError, MemoryStoreError, NotFoundError, ValidationError -from .ledger import LEDGER_FILENAME, VERDICT_ACCEPTED, VERDICT_REJECTED, Decision, DecisionLedger +from .ledger import VERDICT_ACCEPTED, VERDICT_REJECTED, Decision, DecisionLedger from .pending import Pending from .record import DATE_FIELDS, MemoryRecord from .sessions import Pointer, parse_pointer @@ -351,7 +351,7 @@ def _entry(self, name: str) -> MemoryRecord: return record def _ledger(self) -> DecisionLedger: - return DecisionLedger(self._store.layout.dream_reports / LEDGER_FILENAME) + return DecisionLedger(self._store.layout) def _usage(self) -> tuple[dict[str, int], dict[str, int], dict[str, str]]: """Lifetime counts decide what was never useful; only new reads earn weight, so one diff --git a/tests/system/test_concurrency.py b/tests/system/test_concurrency.py index a2db2e26..a2698d79 100644 --- a/tests/system/test_concurrency.py +++ b/tests/system/test_concurrency.py @@ -2,8 +2,11 @@ import multiprocessing as mp +from agent_memory.core.ledger import Decision, DecisionLedger from agent_memory.core.store import Store +DECISIONS_PER_WRITER = 6 + def _write(root: str, name: str) -> None: store = Store(root, agent=name) @@ -34,3 +37,31 @@ def test_two_processes_recording_at_once_lose_nothing(tmp_path): assert {"writer-alpha", "writer-beta"} <= names index_lines = store.layout.memory_index.read_text(encoding="utf-8") assert "writer-alpha" in index_lines and "writer-beta" in index_lines + + +def _append(root: str, name: str) -> None: + ledger = DecisionLedger(Store(root).layout) + for step in range(DECISIONS_PER_WRITER): + ledger.append( + Decision(proposal_id=f"{name}-{step}", verdict="rejected", at="2026-01-01T00:00:00Z") + ) + + +def test_two_processes_appending_decisions_lose_nothing(tmp_path): + root = tmp_path / "store" + Store(root).init() + context = mp.get_context("spawn") + writers = ("writer-alpha", "writer-beta", "writer-gamma") + workers = [ + context.Process(target=_append, args=(str(root), name)) + for name in writers + ] + for worker in workers: + worker.start() + for worker in workers: + worker.join() + assert [worker.exitcode for worker in workers] == [0, 0, 0] + + decided = DecisionLedger(Store(root).layout).decided() + expected = {f"{name}-{step}" for name in writers for step in range(DECISIONS_PER_WRITER)} + assert expected <= set(decided)