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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
40 changes: 27 additions & 13 deletions src/basic_memory/indexing/file_index_checking.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
FileIndexTarget,
IndexedChecksums,
build_file_index_plan,
current_file_index_decision,
move_orphan_file_index_decision,
plan_file_index_target_from_current,
plan_file_index_target_from_observed,
Expand Down Expand Up @@ -46,7 +47,7 @@ async def load_current_file_checksum(


class IndexedFileChecksumRow(Protocol):
"""Tuple-like row of file path, indexed checksum and sync checksum."""
"""Tuple-like row of file path, indexed checksum, sync checksum and accepted checksum."""

def __getitem__(self, index: int, /) -> object:
"""Return a row field by positional index."""
Expand All @@ -62,7 +63,7 @@ async def get_by_file_paths(
*,
content_types: Mapping[str, str | None] | None = None,
) -> Sequence[IndexedFileChecksumRow]:
"""Return rows whose first three fields are file path, checksum and sync checksum."""
"""Return rows of file path, checksum, sync checksum and accepted content checksum."""


def indexed_checksums_by_path(
Expand All @@ -73,6 +74,7 @@ def indexed_checksums_by_path(
str(row[0]): IndexedChecksums(
checksum=None if row[1] is None else str(row[1]),
sync_checksum=None if row[2] is None else str(row[2]),
accepted_content_checksum=None if row[3] is None else str(row[3]),
)
for row in rows
}
Expand Down Expand Up @@ -293,14 +295,12 @@ class FileIndexChecker:
# ghost from a legitimate byte-identical copy (which has no marker and stays a new file).
moved_entity_source: MovedEntitySource | None = None
move_vacate_source: MoveVacateSource | None = None
# The move-vacate marker stores the moved note's *content* checksum. When storage freshness
# keys on a different checksum domain than note content (e.g. cloud indexes by S3 ETag but
# note content is SHA-256), the gate must compare the marker against the current object's
# content checksum, not the freshness checksum, or the marker never matches and every orphan
# is either ignored or wrongly retired (basic-memory-cloud#1601). Left unset when the two
# domains coincide (local uses one SHA-256 checksum everywhere), in which case the gate falls
# back to the freshness `current_checksum` and behavior is unchanged.
move_orphan_checksum_source: CurrentFileChecksumSource | None = None
# The current object's *content* checksum, for checks keyed on note content: the
# move-vacate marker (#1601) and an accepted note's own storage echo. When storage freshness
# keys on a different checksum domain than note content (cloud indexes by S3 ETag, note
# content is SHA-256) this reads the content checksum. Left unset when the two domains
# coincide (local uses one SHA-256 checksum everywhere): the freshness checksum serves both.
content_checksum_source: CurrentFileChecksumSource | None = None

async def detect(self, targets: Sequence[FileIndexTarget]) -> FileIndexPlan:
"""Return the file paths whose current storage object still needs indexing."""
Expand Down Expand Up @@ -382,6 +382,20 @@ async def inspect_target(
indexed=indexed,
current_checksum=current_checksum,
)
# A DB-first write's own storage echo can arrive before its storage checksum is
# recorded. If the object holds exactly the accepted content, there is nothing to read.
if (
decision.status == FileIndexDecisionStatus.read
and indexed is not None
and indexed.accepted_content_checksum is not None
):
content_checksum = (
await self.content_checksum_source.load_current_file_checksum(target.path)
if self.content_checksum_source is not None
else current_checksum
)
if indexed.holds_accepted_content(content_checksum):
decision = current_file_index_decision(target.path)
return decision, current_checksum

async def _apply_move_orphan_gate(
Expand Down Expand Up @@ -423,7 +437,7 @@ async def _apply_move_orphan_gate(
)
)
# The marker records the moved note's content checksum. Compare it against the object's
# content checksum (from move_orphan_checksum_source when the freshness domain differs),
# content checksum (from content_checksum_source when the freshness domain differs),
# falling back to the freshness checksum when the two domains coincide. The moved-entity
# checksum below stays in the freshness/entity domain, so gap-(a) keeps using current.
gated_candidates: list[
Expand All @@ -434,8 +448,8 @@ async def _apply_move_orphan_gate(
if marker is None:
continue
content_checksum = (
await self.move_orphan_checksum_source.load_current_file_checksum(item.target.path)
if self.move_orphan_checksum_source is not None
await self.content_checksum_source.load_current_file_checksum(item.target.path)
if self.content_checksum_source is not None
else item.current_checksum
)
if marker.checksum is not None and marker.checksum != content_checksum:
Expand Down
16 changes: 16 additions & 0 deletions src/basic_memory/indexing/file_index_planning.py
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,8 @@ class IndexedChecksums:

checksum: FileIndexChecksum | None
sync_checksum: FileIndexChecksum | None = None
# The content checksum of the note's accepted Markdown, for notes written DB-first.
accepted_content_checksum: FileIndexChecksum | None = None

def recognizes(self, storage_checksum: FileIndexChecksum | None) -> bool:
"""Whether a stored file with `storage_checksum` is already indexed.
Expand All @@ -38,6 +40,20 @@ def recognizes(self, storage_checksum: FileIndexChecksum | None) -> bool:
return False
return storage_checksum in (self.checksum, self.sync_checksum)

def holds_accepted_content(self, content_checksum: FileIndexChecksum | None) -> bool:
"""Whether a stored file whose content hashes to `content_checksum` is the accepted note.

A DB-first write indexes its accepted content before the file is written, and its
storage checksum is recorded only after the write, so the file's own storage
notification can arrive first, including for a brand-new note whose storage
checksum is still unset. Matching the content closes that gap. The indexed-checksum
lookup withholds the accepted checksum from an incomplete row (pending graph
publication), which is therefore still read.
"""
if content_checksum is None:
return False
return content_checksum == self.accepted_content_checksum


@dataclass(frozen=True, slots=True)
class FileIndexTarget:
Expand Down
28 changes: 24 additions & 4 deletions src/basic_memory/repository/entity_repository.py
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@
from sqlalchemy.orm.interfaces import LoaderOption
from sqlalchemy.engine import Row

from basic_memory.models.knowledge import Entity, Observation, Relation
from basic_memory.models.knowledge import Entity, NoteContent, Observation, Relation
from basic_memory.models.relation_search_refresh import RelationSearchRefresh
from basic_memory.repository.repository import Repository
from basic_memory.runtime.storage import (
Expand Down Expand Up @@ -448,10 +448,30 @@ async def get_by_file_paths(
(or_(*incomplete_projection), None),
else_=Entity.checksum,
).label("checksum")
# The same mask hides the accepted checksum, so an incomplete projection is
# read and repaired even when its content matches. A first write that has not
# recorded a storage checksum yet is not masked and can still match.
accepted_content_checksum = case(
(or_(*incomplete_projection), None),
else_=NoteContent.db_checksum,
).label("accepted_content_checksum")
# The sync checksum is returned as stored: an incomplete row's masked checksum
# already forces a read, whatever the sync checksum says.
path_query = select(Entity.file_path, indexed_checksum, Entity.sync_checksum).where(
Entity.file_path.in_(paths)
# already forces a read, whatever the sync checksum says. The accepted content
# checksum lets a storage echo of the accepted note be recognized by its content
# when its storage checksum has not been recorded yet.
path_query = (
select(
Entity.file_path,
indexed_checksum,
Entity.sync_checksum,
accepted_content_checksum,
)
.outerjoin(
NoteContent,
(NoteContent.entity_id == Entity.id)
& (NoteContent.file_path == Entity.file_path),
)
.where(Entity.file_path.in_(paths))
)
queries.append(self._add_project_filter(path_query))
query = union_all(*queries)
Expand Down
114 changes: 114 additions & 0 deletions test-int/test_accepted_content_echo_gate.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,114 @@
"""A DB-first note's storage echo is current when the object holds the accepted content."""

from hashlib import sha256
from pathlib import Path
from typing import Literal

import pytest
from httpx import AsyncClient
from sqlalchemy import select, update
from sqlalchemy.ext.asyncio import AsyncEngine, AsyncSession, async_sessionmaker

from basic_memory import db
from basic_memory.index.local_runtime import LocalStorageFileMetadataSource
from basic_memory.indexing.file_index_checking import (
FileIndexChecker,
RepositoryIndexedFileChecksumSource,
)
from basic_memory.indexing.file_index_planning import FileIndexDecisionStatus, FileIndexTarget
from basic_memory.markdown import EntityParser
from basic_memory.markdown.markdown_processor import MarkdownProcessor
from basic_memory.models import Entity, NoteContent, Project, RelationSearchRefresh
from basic_memory.repository import EntityRepository
from basic_memory.services import FileService

type EchoState = Literal["accepted", "first-write", "publication-pending", "other-bytes"]


@pytest.mark.parametrize(
("state", "expected"),
[
pytest.param("accepted", FileIndexDecisionStatus.current, id="accepted-content"),
# A brand-new note has no storage checksum until its first materialization settles.
pytest.param("first-write", FileIndexDecisionStatus.current, id="first-write-echo"),
pytest.param("publication-pending", FileIndexDecisionStatus.read, id="repair-still-reads"),
pytest.param("other-bytes", FileIndexDecisionStatus.read, id="external-change-reads"),
],
)
async def test_storage_echo_before_the_storage_checksum_is_recorded(
client: AsyncClient,
test_project: Project,
engine_factory: tuple[AsyncEngine, async_sessionmaker[AsyncSession]],
state: EchoState,
expected: FileIndexDecisionStatus,
) -> None:
"""The window between a materializer's file write and its recording the new checksum."""
_, session_maker = engine_factory
created = await client.post(
f"/v2/projects/{test_project.external_id}/knowledge/write",
json={"note": {"title": "Echo", "directory": "notes", "content": "Revision one"}},
)
assert created.json()["kind"] == "created", created.text
file_path = "notes/Echo.md"
home = Path(test_project.path)

# Revision two is accepted and its bytes are in storage, but the entity still records
# revision one's checksum: the materializer has written the file and not yet settled.
async with db.scoped_session(session_maker) as session:
note_content = await session.scalar(
select(NoteContent).where(NoteContent.file_path == file_path)
)
assert note_content is not None
revision_two = f"{note_content.markdown_content}\nRevision two\n"
await session.execute(
update(NoteContent)
.where(NoteContent.entity_id == note_content.entity_id)
.values(
markdown_content=revision_two,
db_checksum=sha256(revision_two.encode()).hexdigest(),
db_version=note_content.db_version + 1,
)
)
if state == "publication-pending":
session.add(
RelationSearchRefresh(
project_id=test_project.id,
entity_id=note_content.entity_id,
publication_generation=note_content.db_version + 1,
)
)
if state == "first-write":
await session.execute(
update(Entity).where(Entity.id == note_content.entity_id).values(checksum=None)
)
entity_checksum = await session.scalar(
select(Entity.checksum).where(Entity.id == note_content.entity_id)
)
stored = revision_two if state != "other-bytes" else "Someone else's edit\n"
(home / file_path).write_bytes(stored.encode())
assert entity_checksum != sha256(stored.encode()).hexdigest()

entity_repository = EntityRepository(project_id=test_project.id)
file_service = FileService(home, MarkdownProcessor(EntityParser(home)))
checker = FileIndexChecker(
indexed_checksum_source=RepositoryIndexedFileChecksumSource(
session_maker=session_maker,
entity_repository=entity_repository,
),
current_checksum_source=LocalStorageFileMetadataSource(file_service),
)
plan = await checker.detect(
(
FileIndexTarget(
path=file_path,
observed_checksum=sha256(stored.encode()).hexdigest(),
),
)
)

status = (
FileIndexDecisionStatus.read
if file_path in plan.paths_to_read
else next(decision.status for decision in plan.decisions if decision.path == file_path)
)
assert status is expected
12 changes: 6 additions & 6 deletions tests/indexing/test_change_planning.py
Original file line number Diff line number Diff line change
Expand Up @@ -221,12 +221,12 @@ async def get_by_file_paths(
self,
session: object,
paths: tuple[str, ...],
) -> list[tuple[str, str | None, str | None]]:
) -> list[tuple[str, str | None, str | None, str | None]]:
self.loaded_checksum_paths = paths
return [
("unchanged.md", "same-checksum", None),
("modified.md", "old-checksum", None),
("null-checksum.md", None, None),
("unchanged.md", "same-checksum", None, None),
("modified.md", "old-checksum", None, None),
("null-checksum.md", None, None, None),
]

async def find_by_checksums(
Expand Down Expand Up @@ -353,11 +353,11 @@ async def get_by_file_paths(
self,
session: object,
paths: tuple[str, ...],
) -> list[tuple[str, str | None, str | None]]:
) -> list[tuple[str, str | None, str | None, str | None]]:
self.path_batch_sizes.append(len(paths))
# Echo each requested path back as an indexed row so the merged result
# can be checked for completeness across batches.
return [(path, f"checksum-{path}", None) for path in paths]
return [(path, f"checksum-{path}", None, None) for path in paths]

async def find_by_checksums(
self,
Expand Down
14 changes: 7 additions & 7 deletions tests/indexing/test_file_index_checking.py
Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,7 @@ def begin(self) -> FakeSessionContext:

@dataclass(slots=True)
class RecordingChecksumRepository:
rows: list[tuple[object, object | None, object | None]]
rows: list[tuple[object, object | None, object | None, object | None]]
calls: list[tuple[object, tuple[str, ...]]] = field(default_factory=list)

async def get_by_file_paths(
Expand All @@ -55,7 +55,7 @@ async def get_by_file_paths(
file_paths: Sequence[str],
*,
content_types: Mapping[str, str | None] | None = None,
) -> list[tuple[object, object | None, object | None]]:
) -> list[tuple[object, object | None, object | None, object | None]]:
self.calls.append((session, tuple(file_paths)))
return self.rows

Expand Down Expand Up @@ -260,8 +260,8 @@ async def test_repository_indexed_file_checksum_source_maps_repository_rows() ->
session = object()
repository = RecordingChecksumRepository(
rows=[
("notes/a.md", "etag-a", None),
("notes/b.md", None, None),
("notes/a.md", "etag-a", None, "sha-a"),
("notes/b.md", None, None, None),
]
)
source = RepositoryIndexedFileChecksumSource(
Expand All @@ -272,7 +272,7 @@ async def test_repository_indexed_file_checksum_source_maps_repository_rows() ->
checksums = await source.load_indexed_file_checksums(["notes/a.md", "notes/b.md"])

assert checksums == {
"notes/a.md": IndexedChecksums("etag-a"),
"notes/a.md": IndexedChecksums("etag-a", accepted_content_checksum="sha-a"),
"notes/b.md": IndexedChecksums(None),
}
assert repository.calls == [(session, ("notes/a.md", "notes/b.md"))]
Expand Down Expand Up @@ -370,7 +370,7 @@ async def test_checker_gate_uses_content_checksum_source_when_domains_differ() -
"""The gate compares the marker against a distinct content-checksum source when set (#1601).

Cloud freshness keys on the S3 ETag while note content (and the marker) use a SHA-256 content
checksum. With move_orphan_checksum_source wired, the gate must compare the marker against the
checksum. With content_checksum_source wired, the gate must compare the marker against the
content checksum, not the ETag freshness checksum — otherwise the marker never matches and the
leftover source is wrongly retired and re-indexed as a ghost.
"""
Expand All @@ -387,7 +387,7 @@ async def test_checker_gate_uses_content_checksum_source_when_domains_differ() -
current_checksum_source=StubCurrentChecksumSource({"koncept/note.md": "etag-freshness"}),
moved_entity_source=moved,
move_vacate_source=vacate,
move_orphan_checksum_source=content_source,
content_checksum_source=content_source,
)

plan = await checker.detect(
Expand Down
4 changes: 2 additions & 2 deletions tests/indexing/test_relation_persistence.py
Original file line number Diff line number Diff line change
Expand Up @@ -638,7 +638,7 @@ async def cleanup_relation_generations(
relations=[IndexedRelation("links_to", "Target", None)],
)
assert await change_detector.load_indexed_file_checksums((sample_entity.file_path,)) == {
sample_entity.file_path: IndexedChecksums(checksum)
sample_entity.file_path: IndexedChecksums(checksum, accepted_content_checksum=checksum)
}
async with db.scoped_session(session_maker) as session:
refreshes = await relation_repository.list_pending_search_refreshes(
Expand Down Expand Up @@ -710,7 +710,7 @@ async def test_generation_zero_relation_forces_generation_publication(
)

assert await change_detector.load_indexed_file_checksums((sample_entity.file_path,)) == {
sample_entity.file_path: IndexedChecksums(checksum)
sample_entity.file_path: IndexedChecksums(checksum, accepted_content_checksum=checksum)
}
async with db.scoped_session(session_maker) as session:
relations = await relation_repository.find_by_type(session, "links_to")
Expand Down
Loading