diff --git a/src/basic_memory/alembic/versions/s2p3e4c5w6k7_add_project_partition_position.py b/src/basic_memory/alembic/versions/s2p3e4c5w6k7_add_project_partition_position.py new file mode 100644 index 000000000..ae7dc1a15 --- /dev/null +++ b/src/basic_memory/alembic/versions/s2p3e4c5w6k7_add_project_partition_position.py @@ -0,0 +1,88 @@ +"""Add the strict accepted-change partition head to projects. + +Revision ID: s2p3e4c5w6k7 +Revises: d2e3f4a5b6c7 +Create Date: 2026-08-29 20:15:00.000000 + +""" + +from typing import Sequence, Union + +from alembic import op +import sqlalchemy as sa + + +revision: str = "s2p3e4c5w6k7" +down_revision: Union[str, None] = "d2e3f4a5b6c7" +branch_labels: Union[str, Sequence[str], None] = None +depends_on: Union[str, Sequence[str], None] = None + + +def upgrade() -> None: + """Add the project partition head and its durable accepted evidence.""" + with op.batch_alter_table("project", schema=None) as batch_op: + batch_op.add_column( + sa.Column( + "partition_position", + sa.Integer(), + server_default=sa.text("0"), + nullable=False, + ) + ) + + op.create_table( + "accepted_project_note_change", + sa.Column("id", sa.Integer(), nullable=False), + sa.Column("project_id", sa.Integer(), nullable=False), + sa.Column("project_external_id", sa.String(), nullable=False), + sa.Column("partition_position", sa.Integer(), nullable=False), + sa.Column("entity_id", sa.Integer(), nullable=False), + sa.Column("note_external_id", sa.String(), nullable=False), + sa.Column("permalink", sa.Text(), nullable=False), + sa.Column("title", sa.Text(), nullable=False), + sa.Column("operation", sa.String(), nullable=False), + sa.Column("file_path", sa.Text(), nullable=False), + sa.Column("previous_file_path", sa.Text(), nullable=True), + sa.Column("accepted_at", sa.DateTime(timezone=True), nullable=False), + sa.Column("source", sa.String(), nullable=False), + sa.Column("db_version", sa.Integer(), nullable=True), + sa.Column("db_checksum", sa.String(), nullable=True), + sa.Column("actor_user_profile_id", sa.String(), nullable=True), + sa.Column("actor_kind", sa.String(), nullable=True), + sa.Column("actor_name", sa.String(), nullable=True), + sa.Column("materialized_at", sa.DateTime(timezone=True), nullable=True), + sa.ForeignKeyConstraint(["project_id"], ["project.id"], ondelete="CASCADE"), + sa.PrimaryKeyConstraint("id"), + sa.UniqueConstraint( + "project_id", + "partition_position", + name="uq_accepted_project_note_change_project_position", + ), + ) + op.create_index( + "ix_accepted_project_note_change_project_materialized", + "accepted_project_note_change", + ["project_id", "materialized_at"], + unique=False, + ) + op.create_index( + "ix_accepted_project_note_change_note_external_id", + "accepted_project_note_change", + ["note_external_id"], + unique=False, + ) + + +def downgrade() -> None: + """Remove accepted evidence and the project partition head.""" + op.drop_index( + "ix_accepted_project_note_change_note_external_id", + table_name="accepted_project_note_change", + ) + op.drop_index( + "ix_accepted_project_note_change_project_materialized", + table_name="accepted_project_note_change", + ) + op.drop_table("accepted_project_note_change") + with op.batch_alter_table("project", schema=None) as batch_op: + batch_op.drop_column("partition_position") diff --git a/src/basic_memory/indexing/accepted_note_mutation_runner.py b/src/basic_memory/indexing/accepted_note_mutation_runner.py index ffa3b9b25..6320ec6ab 100644 --- a/src/basic_memory/indexing/accepted_note_mutation_runner.py +++ b/src/basic_memory/indexing/accepted_note_mutation_runner.py @@ -2,7 +2,7 @@ from __future__ import annotations -from dataclasses import dataclass +from dataclasses import dataclass, replace from datetime import UTC, datetime from enum import StrEnum from typing import NoReturn, Protocol @@ -46,6 +46,10 @@ ) from basic_memory.runtime.note_move import normalize_note_move_destination_path from basic_memory.runtime.note_object_metadata import NOTE_SOURCE_COLLABORATION_RELAY +from basic_memory.runtime.project_partition import ( + RuntimeAcceptedProjectNoteChange, + RuntimeProjectNoteOperation, +) from basic_memory.runtime.storage import ( NoteExternalId, ProjectExternalId, @@ -64,6 +68,7 @@ type AcceptedNoteMutationChange = RuntimeAcceptedNoteChange[RuntimeNoteContentResponsePayload] type AcceptedNoteMutationUserProfileId = UUID +ACCEPTED_NOTE_DELETE_SOURCE: RuntimeNoteChangeSource = "delete_note" class AcceptedNoteMutationRejectKind(StrEnum): @@ -142,6 +147,7 @@ class AcceptedNoteCreateMutation: data: EntitySchema actor: AcceptedNoteMutationActor source: RuntimeNoteChangeSource + publish_graph_facts: bool = True @dataclass(frozen=True, slots=True) @@ -155,6 +161,7 @@ class AcceptedNoteUpdateMutation: source: RuntimeNoteChangeSource # db_checksum the caller last synced; None means no precondition (issue #1445). base_checksum: str | None = None + publish_graph_facts: bool = True @dataclass(frozen=True, slots=True) @@ -214,6 +221,18 @@ async def get_by_external_id( external_id: ProjectExternalId, ) -> Project | None: ... + async def advance_partition_position( + self, + session: AsyncSession, + project_id: ProjectId, + ) -> int: ... + + async def record_accepted_note_change( + self, + session: AsyncSession, + change: RuntimeAcceptedProjectNoteChange, + ) -> None: ... + class AcceptedNoteMutationEntityRepository(Protocol): """Entity lookup capability for accepted-note mutations.""" @@ -321,6 +340,92 @@ class AcceptedNoteMutationResult: relation_publication: RelationGenerationPublication | None = None +async def record_accepted_project_note_change( + session: AsyncSession, + *, + project: Project, + entity: Entity, + operation: RuntimeProjectNoteOperation, + accepted_at: datetime, + source: RuntimeNoteChangeSource, + previous_file_path: RuntimeFilePath | None, + note_content: NoteContent | None, + actor: AcceptedNoteMutationActor | None, + dependencies: AcceptedNoteMutationDependencies, +) -> RuntimeAcceptedProjectNoteChange: + """Claim and describe one accepted change in the project's strict partition.""" + if entity.permalink is None: + raise RuntimeError( + f"Accepted note is missing permalink for entity_id={entity.id}" + ) + position = await dependencies.project_repository.advance_partition_position( + session, + project.id, + ) + change = RuntimeAcceptedProjectNoteChange( + project_id=project.id, + project_external_id=project.external_id, + partition_position=position, + entity_id=entity.id, + note_external_id=entity.external_id, + permalink=entity.permalink, + title=entity.title, + operation=operation, + file_path=entity.file_path, + previous_file_path=previous_file_path, + accepted_at=accepted_at, + source=source, + db_version=note_content.db_version if note_content is not None else None, + db_checksum=note_content.db_checksum if note_content is not None else None, + actor_user_profile_id=actor.user_profile_id if actor is not None else None, + actor_kind=actor.kind if actor is not None else None, + actor_name=actor.name if actor is not None else None, + ) + await dependencies.project_repository.record_accepted_note_change(session, change) + return change + + +def attach_accepted_project_note_change( + change: AcceptedNoteMutationChange, + project_change: RuntimeAcceptedProjectNoteChange, +) -> AcceptedNoteMutationChange: + """Carry accepted partition evidence through existing runtime follow-up work.""" + materialization = ( + replace(change.materialization, project_change=project_change) + if change.materialization is not None + else None + ) + file_delete = ( + replace(change.file_delete, project_change=project_change) + if change.file_delete is not None + else None + ) + return replace( + change, + project_change=project_change, + materialization=materialization, + file_delete=file_delete, + ) + + +def apply_accepted_note_graph_policy( + prepared_write: AcceptedPreparedNoteWrite, + *, + publish_graph_facts: bool, +) -> AcceptedPreparedNoteWrite: + """Keep canonical Markdown while suppressing graph facts for derived documents.""" + if publish_graph_facts: + return prepared_write + prepared = prepared_write.prepared + graph_silent_markdown = prepared.entity_markdown.model_copy( + update={"observations": [], "relations": []} + ) + return replace( + prepared_write, + prepared=replace(prepared, entity_markdown=graph_silent_markdown), + ) + + def accepted_note_integrity_rejection(error: IntegrityError) -> AcceptedNoteMutationRejection: """Map repository integrity errors into portable accepted-note rejections.""" conflict_kind = classify_accepted_note_write_conflict(str(error.orig or error)) @@ -523,6 +628,31 @@ async def run_accepted_note_delete( ) ) + # The entity lookup above is intentionally unlocked so an already-missing + # delete stays idempotent. Once the note exists, claim its mutation lock and + # reload it before deleting: another delete may have removed the row while we + # waited, while an update or move may have changed its accepted evidence. + await lock_accepted_note_content_for_entity_mutation( + session, + project_id=project.id, + entity_id=entity.id, + ) + entity = await entity_repository.get_by_external_id( + session, + request.entity_external_id, + load_relations=False, + ) + if entity is None: + return AcceptedNoteMutationResult( + change=await delete_accepted_note( + session, + project_id=project.id, + entity=None, + repositories=dependencies.write_repositories, + ) + ) + + await session.refresh(entity) note_content = await load_accepted_note_content( session, project_id=project.id, @@ -530,14 +660,34 @@ async def run_accepted_note_delete( dependencies=dependencies, missing_kind=None, ) + if note_content is not None: + await session.refresh(note_content) + change = await delete_accepted_note( + session, + project_id=project.id, + entity=entity, + note_content=note_content, + repositories=dependencies.write_repositories, + ) + # The legacy entity-delete route also removes binary resources. Keep that + # behavior, but do not claim an accepted-note partition position for data + # that cannot participate in Markdown indexing or Wiki projection. + if not runtime_content_type_is_markdown(entity): + return AcceptedNoteMutationResult(change=change) + project_change = await record_accepted_project_note_change( + session, + project=project, + entity=entity, + operation=RuntimeProjectNoteOperation.deleted, + accepted_at=accepted_note_mutation_utc_now(), + source=ACCEPTED_NOTE_DELETE_SOURCE, + previous_file_path=None, + note_content=note_content, + actor=None, + dependencies=dependencies, + ) return AcceptedNoteMutationResult( - change=await delete_accepted_note( - session, - project_id=project.id, - entity=entity, - note_content=note_content, - repositories=dependencies.write_repositories, - ) + change=attach_accepted_project_note_change(change, project_change) ) @@ -583,7 +733,10 @@ async def _run_accepted_note_create( check_storage_exists=dependencies.verify_storage_absent_on_create, session=session, ) - + prepared_write = apply_accepted_note_graph_policy( + prepared_write, + publish_graph_facts=request.publish_graph_facts, + ) prepared = prepared_write.prepared entity = await create_accepted_pending_entity( session, @@ -602,15 +755,30 @@ async def _run_accepted_note_create( self_relation_resolver=preparer, repositories=dependencies.write_repositories, ) + project_change = await record_accepted_project_note_change( + session, + project=project, + entity=entity, + operation=RuntimeProjectNoteOperation.created, + accepted_at=now, + source=request.source, + previous_file_path=None, + note_content=persisted.note_content, + actor=request.actor, + dependencies=dependencies, + ) return AcceptedNoteMutationResult( - change=plan_accepted_note_write_change( - status_code=201, - entity=entity, - note_content=persisted.note_content, - actor_user_profile_id=request.actor.user_profile_id, - actor_kind=request.actor.kind, - actor_name=request.actor.name, - fallback_source=request.source, + change=attach_accepted_project_note_change( + plan_accepted_note_write_change( + status_code=201, + entity=entity, + note_content=persisted.note_content, + actor_user_profile_id=request.actor.user_profile_id, + actor_kind=request.actor.kind, + actor_name=request.actor.name, + fallback_source=request.source, + ), + project_change, ), relation_publication=persisted.relation_publication, ) @@ -640,6 +808,32 @@ async def _run_accepted_note_update( load_relations=False, ) created = entity is None + current_note_content: NoteContent | None = None + + if entity is not None: + # The first lookup resolves the addressed identity, but it is not a + # stable source-path snapshot. Claim the note lock, then refresh the + # entity before any path-sensitive preparation: a concurrent move may + # have changed the accepted predecessor while this PUT was waiting. + if not runtime_content_type_is_markdown(entity): + reject_accepted_note_mutation( + AcceptedNoteMutationRejectKind.unsupported_media_type, + "Only markdown note mutations are supported by the note-content path.", + ) + current_note_content = await load_required_accepted_note_content( + session, + project_id=project.id, + entity_id=entity.id, + dependencies=dependencies, + missing_kind=AcceptedNoteMutationRejectKind.conflict, + ) + await session.refresh(entity) + if not runtime_content_type_is_markdown(entity): + reject_accepted_note_mutation( + AcceptedNoteMutationRejectKind.unsupported_media_type, + "Only markdown note mutations are supported by the note-content path.", + ) + existing_file_path = entity.file_path if entity is not None else None vacated_source: tuple[RuntimeFilePath, RuntimeFileChecksum | None] | None = None @@ -694,14 +888,6 @@ async def _run_accepted_note_update( ) current_note_content = None else: - # A PUT replacement can only target a markdown note. A watcher-indexed binary - # entity has no markdown note_content to replace, so reject with 415 like the - # edit/move paths instead of a permanent-looking 409 content-backfill retry. - if not runtime_content_type_is_markdown(entity): - reject_accepted_note_mutation( - AcceptedNoteMutationRejectKind.unsupported_media_type, - "Only markdown note mutations are supported by the note-content path.", - ) # Local source-of-truth guard: a PUT that renames onto a destination file that # exists on disk but is not yet indexed would overwrite/lose that unindexed # write. Mirror the create/move storage check before committing DB/search to @@ -714,13 +900,7 @@ async def _run_accepted_note_update( ) except EntityAlreadyExistsError as error: reject_accepted_note_mutation(AcceptedNoteMutationRejectKind.conflict, str(error)) - current_note_content = await load_required_accepted_note_content( - session, - project_id=project.id, - entity_id=entity.id, - dependencies=dependencies, - missing_kind=AcceptedNoteMutationRejectKind.conflict, - ) + assert current_note_content is not None # A PUT replacement may also rename the note. Capture the exact source bytes before # persistence mutates the entity and note_content to the destination version so delayed # cleanup cannot let a later project index recreate the old path as a ghost. @@ -780,6 +960,10 @@ async def _run_accepted_note_update( except (ParseError, ValueError) as error: reject_accepted_note_mutation(AcceptedNoteMutationRejectKind.bad_request, str(error)) + prepared_write = apply_accepted_note_graph_policy( + prepared_write, + publish_graph_facts=request.publish_graph_facts, + ) prepared = prepared_write.prepared persisted = await persist_accepted_note_snapshot( session, @@ -806,16 +990,43 @@ async def _run_accepted_note_update( file_path=vacated_source[0], file_checksum=vacated_source[1], ) + operation = ( + RuntimeProjectNoteOperation.created + if created + else ( + RuntimeProjectNoteOperation.moved + if existing_file_path != entity.file_path + else RuntimeProjectNoteOperation.updated + ) + ) + previous_file_path = ( + existing_file_path if operation == RuntimeProjectNoteOperation.moved else None + ) + project_change = await record_accepted_project_note_change( + session, + project=project, + entity=entity, + operation=operation, + accepted_at=now, + source=request.source, + previous_file_path=previous_file_path, + note_content=persisted.note_content, + actor=request.actor, + dependencies=dependencies, + ) return AcceptedNoteMutationResult( - change=plan_accepted_note_write_change( - status_code=201 if created else 200, - entity=entity, - note_content=persisted.note_content, - actor_user_profile_id=request.actor.user_profile_id, - actor_kind=request.actor.kind, - actor_name=request.actor.name, - cleanup_after_write=persisted.previous_file_delete, - fallback_source=request.source, + change=attach_accepted_project_note_change( + plan_accepted_note_write_change( + status_code=201 if created else 200, + entity=entity, + note_content=persisted.note_content, + actor_user_profile_id=request.actor.user_profile_id, + actor_kind=request.actor.kind, + actor_name=request.actor.name, + cleanup_after_write=persisted.previous_file_delete, + fallback_source=request.source, + ), + project_change, ), relation_publication=persisted.relation_publication, ) @@ -869,15 +1080,30 @@ async def _run_accepted_note_edit( self_relation_resolver=preparer, repositories=dependencies.write_repositories, ) + project_change = await record_accepted_project_note_change( + session, + project=project, + entity=entity, + operation=RuntimeProjectNoteOperation.updated, + accepted_at=now, + source=request.source, + previous_file_path=None, + note_content=persisted.note_content, + actor=request.actor, + dependencies=dependencies, + ) return AcceptedNoteMutationResult( - change=plan_accepted_note_write_change( - status_code=200, - entity=entity, - note_content=persisted.note_content, - actor_user_profile_id=request.actor.user_profile_id, - actor_kind=request.actor.kind, - actor_name=request.actor.name, - fallback_source=request.source, + change=attach_accepted_project_note_change( + plan_accepted_note_write_change( + status_code=200, + entity=entity, + note_content=persisted.note_content, + actor_user_profile_id=request.actor.user_profile_id, + actor_kind=request.actor.kind, + actor_name=request.actor.name, + fallback_source=request.source, + ), + project_change, ), relation_publication=persisted.relation_publication, ) @@ -904,6 +1130,10 @@ async def _run_accepted_note_move( entity_external_id=request.entity_external_id, dependencies=dependencies, ) + # The identity lookup precedes the NoteContent lock. Refresh after the lock + # so an overlapping move records the committed source path it actually + # replaces, rather than the path observed while waiting. + await session.refresh(entity) existing_file_path = entity.file_path # The destination filename keeps its requested casing; only the parent # directory resolves against existing folders (issue #1326). The path is @@ -994,17 +1224,32 @@ async def _run_accepted_note_move( file_path=existing_file_path, file_checksum=vacated_source_checksum, ) + project_change = await record_accepted_project_note_change( + session, + project=project, + entity=entity, + operation=RuntimeProjectNoteOperation.moved, + accepted_at=now, + source=request.source, + previous_file_path=existing_file_path, + note_content=persisted.note_content, + actor=request.actor, + dependencies=dependencies, + ) return AcceptedNoteMutationResult( - change=plan_accepted_note_write_change( - status_code=200, - entity=entity, - note_content=persisted.note_content, - actor_user_profile_id=request.actor.user_profile_id, - actor_kind=request.actor.kind, - actor_name=request.actor.name, - previous_file_path=existing_file_path, - cleanup_after_write=persisted.previous_file_delete, - fallback_source=request.source, + change=attach_accepted_project_note_change( + plan_accepted_note_write_change( + status_code=200, + entity=entity, + note_content=persisted.note_content, + actor_user_profile_id=request.actor.user_profile_id, + actor_kind=request.actor.kind, + actor_name=request.actor.name, + previous_file_path=existing_file_path, + cleanup_after_write=persisted.previous_file_delete, + fallback_source=request.source, + ), + project_change, ), relation_publication=persisted.relation_publication, ) diff --git a/src/basic_memory/indexing/wiki_projector.py b/src/basic_memory/indexing/wiki_projector.py index 267911210..861581aca 100644 --- a/src/basic_memory/indexing/wiki_projector.py +++ b/src/basic_memory/indexing/wiki_projector.py @@ -625,7 +625,11 @@ def _render_document( body: str, include_okf_version: bool, ) -> bytes: - frontmatter = ["---", f"type: {note_type}"] + frontmatter = [ + "---", + f"type: {note_type}", + "bm_parse_semantics: false", + ] if include_okf_version: frontmatter.append(f'okf_version: "{OKF_VERSION}"') frontmatter.extend( diff --git a/src/basic_memory/models/__init__.py b/src/basic_memory/models/__init__.py index 186d45b5e..cc5285b78 100644 --- a/src/basic_memory/models/__init__.py +++ b/src/basic_memory/models/__init__.py @@ -9,11 +9,12 @@ Observation, Relation, ) -from basic_memory.models.project import Project +from basic_memory.models.project import AcceptedProjectNoteChange, Project from basic_memory.models.relation_search_refresh import RelationSearchRefresh __all__ = [ "Base", + "AcceptedProjectNoteChange", "Entity", "NoteContent", "NoteFileVacate", diff --git a/src/basic_memory/models/project.py b/src/basic_memory/models/project.py index a39ddaacf..fa388eefb 100644 --- a/src/basic_memory/models/project.py +++ b/src/basic_memory/models/project.py @@ -11,13 +11,16 @@ Boolean, DateTime, Float, + ForeignKey, Index, + UniqueConstraint, + text, event, ) from sqlalchemy.orm import Mapped, mapped_column, relationship from basic_memory.models.base import Base -from basic_memory.utils import generate_permalink +from basic_memory.utils import ensure_timezone_aware, generate_permalink class Project(Base): @@ -70,9 +73,26 @@ class Project(Base): last_scan_timestamp: Mapped[Optional[float]] = mapped_column(Float, nullable=True) last_file_count: Mapped[Optional[int]] = mapped_column(Integer, nullable=True) + # Strict, project-local ordering for accepted durable changes. This is the + # generic partition head that future event-journal projectors can reuse; it + # is deliberately not derived from timestamps, webhooks, or scan activity. + partition_position: Mapped[int] = mapped_column( + Integer, + nullable=False, + default=0, + server_default=text("0"), + ) + # Define relationships to entities, observations, and relations # These relationships will be established once we add project_id to those models entities = relationship("Entity", back_populates="project", cascade="all, delete-orphan") + accepted_note_changes = relationship( + "AcceptedProjectNoteChange", + back_populates="project", + cascade="all, delete-orphan", + order_by="AcceptedProjectNoteChange.partition_position", + passive_deletes=True, + ) @override def __repr__(self) -> str: # pragma: no cover @@ -90,3 +110,67 @@ def set_project_permalink(mapper, connection, project): # If the name changed or permalink is empty, regenerate permalink if not project.permalink or project.permalink != generate_permalink(project.name): project.permalink = generate_permalink(project.name) + + +class AcceptedProjectNoteChange(Base): + """Durable evidence for one accepted note mutation in project order. + + The row intentionally retains note identity and path values instead of an + entity foreign key: delete evidence must survive after the entity is gone. + """ + + __tablename__ = "accepted_project_note_change" + __table_args__ = ( + UniqueConstraint( + "project_id", + "partition_position", + name="uq_accepted_project_note_change_project_position", + ), + Index( + "ix_accepted_project_note_change_project_materialized", + "project_id", + "materialized_at", + ), + Index( + "ix_accepted_project_note_change_note_external_id", + "note_external_id", + ), + ) + + id: Mapped[int] = mapped_column(Integer, primary_key=True) + project_id: Mapped[int] = mapped_column( + ForeignKey("project.id", ondelete="CASCADE"), + nullable=False, + ) + project_external_id: Mapped[str] = mapped_column(String, nullable=False) + partition_position: Mapped[int] = mapped_column(Integer, nullable=False) + entity_id: Mapped[int] = mapped_column(Integer, nullable=False) + note_external_id: Mapped[str] = mapped_column(String, nullable=False) + permalink: Mapped[str] = mapped_column(Text, nullable=False) + title: Mapped[str] = mapped_column(Text, nullable=False) + operation: Mapped[str] = mapped_column(String, nullable=False) + file_path: Mapped[str] = mapped_column(Text, nullable=False) + previous_file_path: Mapped[Optional[str]] = mapped_column(Text, nullable=True) + accepted_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False) + source: Mapped[str] = mapped_column(String, nullable=False) + db_version: Mapped[Optional[int]] = mapped_column(Integer, nullable=True) + db_checksum: Mapped[Optional[str]] = mapped_column(String, nullable=True) + actor_user_profile_id: Mapped[Optional[str]] = mapped_column(String, nullable=True) + actor_kind: Mapped[Optional[str]] = mapped_column(String, nullable=True) + actor_name: Mapped[Optional[str]] = mapped_column(String, nullable=True) + materialized_at: Mapped[Optional[datetime]] = mapped_column( + DateTime(timezone=True), + nullable=True, + ) + + project = relationship("Project", back_populates="accepted_note_changes") + + @override + def __getattribute__(self, name: str): + value = super().__getattribute__(name) + # SQLite drops timezone information when persisting DateTime columns. + # Normalize project-journal timestamps at the model boundary so the + # projector receives the same aware instants on every supported backend. + if name in {"accepted_at", "materialized_at"} and isinstance(value, datetime): + return ensure_timezone_aware(value) + return value diff --git a/src/basic_memory/repository/project_repository.py b/src/basic_memory/repository/project_repository.py index e5d64dbb3..2a4b08045 100644 --- a/src/basic_memory/repository/project_repository.py +++ b/src/basic_memory/repository/project_repository.py @@ -1,16 +1,18 @@ """Repository for managing projects in Basic Memory.""" +from datetime import datetime from pathlib import Path from typing import Any, override, Optional, Sequence, Union from loguru import logger -from sqlalchemy import Executable, inspect as sa_inspect, select, text +from sqlalchemy import Executable, inspect as sa_inspect, select, text, update from sqlalchemy.exc import NoResultFound, OperationalError from sqlalchemy.ext.asyncio import AsyncSession -from basic_memory.models.project import Project +from basic_memory.models.project import AcceptedProjectNoteChange, Project from basic_memory.repository.repository import Repository +from basic_memory.runtime.project_partition import RuntimeAcceptedProjectNoteChange async def _load_sqlite_vec_on_session(session) -> bool: @@ -153,6 +155,115 @@ async def get_default_project(self, session: AsyncSession) -> Optional[Project]: query = self.select().where(Project.is_default.is_(True)) return await self.find_one(session, query) + async def advance_partition_position( + self, + session: AsyncSession, + project_id: int, + ) -> int: + """Atomically claim the next strict position in one project partition.""" + statement = ( + update(Project) + .where(Project.id == project_id) + .values(partition_position=Project.partition_position + 1) + .returning(Project.partition_position) + .execution_options(synchronize_session=False) + ) + position = (await session.execute(statement)).scalar_one_or_none() + if position is None: + raise RuntimeError(f"Project partition is missing for project_id={project_id}") + return int(position) + + async def record_accepted_note_change( + self, + session: AsyncSession, + change: RuntimeAcceptedProjectNoteChange, + ) -> None: + """Persist replayable accepted evidence in the mutation transaction.""" + persisted = AcceptedProjectNoteChange( + project_id=change.project_id, + project_external_id=change.project_external_id, + partition_position=change.partition_position, + entity_id=change.entity_id, + note_external_id=change.note_external_id, + permalink=change.permalink, + title=change.title, + operation=change.operation.value, + file_path=change.file_path, + previous_file_path=change.previous_file_path, + accepted_at=change.accepted_at, + source=change.source, + db_version=change.db_version, + db_checksum=change.db_checksum, + actor_user_profile_id=( + str(change.actor_user_profile_id) + if change.actor_user_profile_id is not None + else None + ), + actor_kind=change.actor_kind, + actor_name=change.actor_name, + ) + session.add(persisted) + await session.flush() + + async def list_accepted_note_changes( + self, + session: AsyncSession, + project_id: int, + *, + after_position: int = 0, + through_position: int | None = None, + ) -> Sequence[AcceptedProjectNoteChange]: + """List one contiguous project evidence range in strict order.""" + statement = select(AcceptedProjectNoteChange).where( + AcceptedProjectNoteChange.project_id == project_id, + AcceptedProjectNoteChange.partition_position > after_position, + ) + if through_position is not None: + statement = statement.where( + AcceptedProjectNoteChange.partition_position <= through_position + ) + result = await session.execute( + statement.order_by(AcceptedProjectNoteChange.partition_position) + ) + return tuple(result.scalars().all()) + + async def mark_accepted_note_change_materialized( + self, + session: AsyncSession, + project_id: int, + partition_position: int, + *, + materialized_at: datetime, + ) -> bool: + """Record when this note's accepted evidence reached canonical storage.""" + target_note_external_id = ( + await session.execute( + select(AcceptedProjectNoteChange.note_external_id).where( + AcceptedProjectNoteChange.project_id == project_id, + AcceptedProjectNoteChange.partition_position == partition_position, + ) + ) + ).scalar_one_or_none() + if target_note_external_id is None: + return False + + # A newer materialized generation is canonical evidence that every older + # accepted generation of this note has been superseded. Retire those rows + # together so an obsolete queued write cannot block downstream projectors. + statement = ( + update(AcceptedProjectNoteChange) + .where( + AcceptedProjectNoteChange.project_id == project_id, + AcceptedProjectNoteChange.note_external_id == target_note_external_id, + AcceptedProjectNoteChange.partition_position <= partition_position, + AcceptedProjectNoteChange.materialized_at.is_(None), + ) + .values(materialized_at=materialized_at) + .returning(AcceptedProjectNoteChange.id) + .execution_options(synchronize_session=False) + ) + return bool((await session.execute(statement)).scalars().all()) + async def get_active_projects(self, session: AsyncSession) -> Sequence[Project]: """Get all active projects.""" query = self.select().where(Project.is_active == True) # noqa: E712 diff --git a/src/basic_memory/runtime/accepted_note_changes.py b/src/basic_memory/runtime/accepted_note_changes.py index 34b12ba3f..ca7fd4053 100644 --- a/src/basic_memory/runtime/accepted_note_changes.py +++ b/src/basic_memory/runtime/accepted_note_changes.py @@ -29,6 +29,7 @@ RuntimePendingNoteMaterializationSource, plan_pending_note_materialization, ) +from basic_memory.runtime.project_partition import RuntimeAcceptedProjectNoteChange from basic_memory.runtime.storage import ( NoteExternalId, ProjectId, @@ -213,6 +214,7 @@ class RuntimeAcceptedNoteChange(Generic[_PayloadT_co]): status_code: int payload: _PayloadT_co + project_change: RuntimeAcceptedProjectNoteChange | None = None materialization: RuntimePendingNoteMaterialization | None = None file_delete: RuntimePendingNoteFileDelete | None = None # Surviving notes whose relations pointed at a deleted target. Their search diff --git a/src/basic_memory/runtime/cleanup.py b/src/basic_memory/runtime/cleanup.py index 8a195460a..03b9d1d20 100644 --- a/src/basic_memory/runtime/cleanup.py +++ b/src/basic_memory/runtime/cleanup.py @@ -21,6 +21,7 @@ RuntimeFileChecksum, RuntimeFilePath, ) +from basic_memory.runtime.project_partition import RuntimeAcceptedProjectNoteChange RUNTIME_FILE_SNAPSHOT_TIMESTAMP_MATCH_EPSILON_SECONDS = 0.001 @@ -196,6 +197,7 @@ class RuntimeNoteFileDeleteJobRequest: entity_id: RuntimeEntityId file_path: RuntimeFilePath file_checksum: RuntimeFileChecksum | None = None + project_change: RuntimeAcceptedProjectNoteChange | None = None # Live note path after the move that scheduled this cleanup; a local adapter # skips the delete when it shares a physical file with file_path. Not part of # dedupe_key: it does not change the logical identity of the delete. @@ -230,6 +232,7 @@ def plan_note_file_delete_job_request( entity_id=file_delete.entity_id, file_path=file_delete.file_path, file_checksum=file_delete.file_checksum, + project_change=file_delete.project_change, live_file_path=file_delete.live_file_path, ) diff --git a/src/basic_memory/runtime/job_payloads.py b/src/basic_memory/runtime/job_payloads.py index 73519f56d..7b54e65bd 100644 --- a/src/basic_memory/runtime/job_payloads.py +++ b/src/basic_memory/runtime/job_payloads.py @@ -18,6 +18,7 @@ VALID_NOTE_OBJECT_SOURCES, normalize_actor_name, ) +from basic_memory.runtime.project_partition import RuntimeAcceptedProjectNoteChange DELETE_NOTE_FILE_ENTRYPOINT: JobEntrypoint = "delete_note_file" @@ -31,6 +32,7 @@ class RuntimeNoteFileDeleteJobPayload(BaseModel): entity_id: int file_path: str file_checksum: str | None = None + project_change: RuntimeAcceptedProjectNoteChange | None = None @classmethod def from_runtime_request(cls, request: RuntimeNoteFileDeleteJobRequest) -> Self: @@ -40,6 +42,7 @@ def from_runtime_request(cls, request: RuntimeNoteFileDeleteJobRequest) -> Self: entity_id=request.entity_id, file_path=request.file_path, file_checksum=request.file_checksum, + project_change=request.project_change, ) def to_runtime_request(self) -> RuntimeNoteFileDeleteJobRequest: @@ -49,6 +52,7 @@ def to_runtime_request(self) -> RuntimeNoteFileDeleteJobRequest: entity_id=self.entity_id, file_path=self.file_path, file_checksum=self.file_checksum, + project_change=self.project_change, ) def runtime_job_request( @@ -72,6 +76,7 @@ class RuntimeNoteMaterializationJobPayload(BaseModel): entity_id: int db_version: int db_checksum: str + project_change: RuntimeAcceptedProjectNoteChange | None = None actor_user_profile_id: UUID | None = None actor_kind: str | None = None actor_name: str | None = None @@ -122,6 +127,7 @@ def from_runtime_request(cls, request: RuntimeNoteMaterializationJobRequest) -> entity_id=request.entity_id, db_version=request.db_version, db_checksum=request.db_checksum, + project_change=request.project_change, actor_user_profile_id=request.actor_user_profile_id, actor_kind=request.actor_kind, actor_name=request.actor_name, @@ -138,6 +144,7 @@ def to_runtime_request(self) -> RuntimeNoteMaterializationJobRequest: entity_id=self.entity_id, db_version=self.db_version, db_checksum=self.db_checksum, + project_change=self.project_change, actor_user_profile_id=self.actor_user_profile_id, actor_kind=self.actor_kind, actor_name=self.actor_name, diff --git a/src/basic_memory/runtime/note_content_deletes.py b/src/basic_memory/runtime/note_content_deletes.py index 2f4337156..5d8bd6126 100644 --- a/src/basic_memory/runtime/note_content_deletes.py +++ b/src/basic_memory/runtime/note_content_deletes.py @@ -14,6 +14,7 @@ RuntimeFilePath, runtime_content_type_is_markdown, ) +from basic_memory.runtime.project_partition import RuntimeAcceptedProjectNoteChange class RuntimeDeletedNoteEntitySource(RuntimeContentTypeSource, Protocol): @@ -222,6 +223,7 @@ class RuntimePendingNoteFileDelete: entity_id: RuntimeEntityId file_path: RuntimeFilePath file_checksum: RuntimeFileChecksum | None = None + project_change: RuntimeAcceptedProjectNoteChange | None = None # The note's live path after the move that scheduled this cleanup. Object # storage treats case-different keys as distinct; a local adapter re-checks # it against the physical filesystem before deleting because a case-only diff --git a/src/basic_memory/runtime/note_materialization_planning.py b/src/basic_memory/runtime/note_materialization_planning.py index efe768fc6..5935f14b1 100644 --- a/src/basic_memory/runtime/note_materialization_planning.py +++ b/src/basic_memory/runtime/note_materialization_planning.py @@ -9,6 +9,7 @@ from uuid import UUID from basic_memory.runtime.note_content_deletes import RuntimePendingNoteFileDelete +from basic_memory.runtime.project_partition import RuntimeAcceptedProjectNoteChange from basic_memory.runtime.storage import ( ProjectId, RuntimeEntityId, @@ -93,6 +94,7 @@ class RuntimePendingNoteMaterialization: entity_id: RuntimeEntityId db_version: RuntimeNoteContentVersion db_checksum: RuntimeNoteContentChecksum + project_change: RuntimeAcceptedProjectNoteChange | None = None actor_user_profile_id: UUID | None = None actor_kind: RuntimeNoteActorKind | None = None actor_name: RuntimeNoteActorName | None = None @@ -137,6 +139,7 @@ class RuntimeNoteMaterializationJobRequest: entity_id: RuntimeEntityId db_version: RuntimeNoteContentVersion db_checksum: RuntimeNoteContentChecksum + project_change: RuntimeAcceptedProjectNoteChange | None = None actor_user_profile_id: UUID | None = None actor_kind: RuntimeNoteActorKind | None = None actor_name: RuntimeNoteActorName | None = None @@ -194,6 +197,7 @@ def plan_note_materialization_job_request( entity_id=materialization.entity_id, db_version=materialization.db_version, db_checksum=materialization.db_checksum, + project_change=materialization.project_change, actor_user_profile_id=materialization.actor_user_profile_id, actor_kind=materialization.actor_kind, actor_name=materialization.actor_name, diff --git a/src/basic_memory/runtime/note_object_metadata.py b/src/basic_memory/runtime/note_object_metadata.py index 350cd7947..ca77324ef 100644 --- a/src/basic_memory/runtime/note_object_metadata.py +++ b/src/basic_memory/runtime/note_object_metadata.py @@ -23,6 +23,7 @@ NOTE_OBJECT_ACTOR_KIND_METADATA = "bm-actor-kind" NOTE_OBJECT_ACTOR_KIND_MCP_CLIENT = "mcp_client" +NOTE_OBJECT_ACTOR_KIND_SYSTEM = "system" NOTE_OBJECT_ACTOR_NAME_METADATA = "bm-actor-name" NOTE_OBJECT_ACTOR_USER_PROFILE_ID_METADATA = "bm-actor-user-profile-id" NOTE_OBJECT_DB_CHECKSUM_METADATA = "bm-db-checksum" @@ -32,7 +33,7 @@ NOTE_OBJECT_FILE_VERSION_METADATA = "bm-file-version" NOTE_OBJECT_SOURCE_METADATA = "bm-note-source" VALID_NOTE_OBJECT_ACTOR_KINDS: frozenset[RuntimeNoteActorKind] = frozenset( - {NOTE_OBJECT_ACTOR_KIND_MCP_CLIENT} + {NOTE_OBJECT_ACTOR_KIND_MCP_CLIENT, NOTE_OBJECT_ACTOR_KIND_SYSTEM} ) # web_v2 = a note write originating from the web-v2 UI. Distinguishing it from # `api` lets clients tell a genuine web-UI edit apart from api/materialization @@ -42,8 +43,18 @@ # note.updated event echoes this source as the write's actor origin. # document_ingestion = a hosted worker accepting a deterministic document # extraction or ingestion-run note through the canonical note mutation path. +# wiki_projector = the deterministic OKF projector accepting generated index +# and log notes through that same path before materializing them as Markdown. VALID_NOTE_OBJECT_SOURCES: frozenset[RuntimeNoteChangeSource] = frozenset( - {"api", "collaboration_relay", "document_ingestion", "mcp", "s3_webhook", "web_v2"} + { + "api", + "collaboration_relay", + "document_ingestion", + "mcp", + "s3_webhook", + "web_v2", + "wiki_projector", + } ) # Named because the accepted-note write path special-cases relay writes: the # relay superseding its own prior write is never a real conflict (#1589). diff --git a/src/basic_memory/runtime/project_partition.py b/src/basic_memory/runtime/project_partition.py new file mode 100644 index 000000000..660de208a --- /dev/null +++ b/src/basic_memory/runtime/project_partition.py @@ -0,0 +1,83 @@ +"""Portable evidence for one accepted change in a strict project partition.""" + +from __future__ import annotations + +from dataclasses import dataclass +from datetime import datetime +from enum import StrEnum +from uuid import UUID + +from basic_memory.runtime.storage import ( + NoteExternalId, + ProjectExternalId, + ProjectId, + RuntimeEntityId, + RuntimeFilePath, + RuntimeNoteActorKind, + RuntimeNoteActorName, + RuntimeNoteChangeSource, + RuntimeNoteContentChecksum, + RuntimeNoteContentVersion, +) + +type ProjectPartitionPosition = int + + +class RuntimeProjectNoteOperation(StrEnum): + """Accepted note operation recorded for project-wide consumers.""" + + created = "created" + updated = "updated" + moved = "moved" + deleted = "deleted" + + +@dataclass(frozen=True, slots=True) +class RuntimeAcceptedProjectNoteChange: + """Replay-complete accepted note evidence carried to runtime follow-ups.""" + + project_id: ProjectId + project_external_id: ProjectExternalId + partition_position: ProjectPartitionPosition + entity_id: RuntimeEntityId + note_external_id: NoteExternalId + permalink: str + title: str + operation: RuntimeProjectNoteOperation + file_path: RuntimeFilePath + accepted_at: datetime + source: RuntimeNoteChangeSource + previous_file_path: RuntimeFilePath | None = None + db_version: RuntimeNoteContentVersion | None = None + db_checksum: RuntimeNoteContentChecksum | None = None + actor_user_profile_id: UUID | None = None + actor_kind: RuntimeNoteActorKind | None = None + actor_name: RuntimeNoteActorName | None = None + + def __post_init__(self) -> None: + if self.partition_position <= 0: + raise ValueError("Accepted project change position must be positive") + if not self.project_external_id.strip(): + raise ValueError("Accepted project change requires project_external_id") + if not self.note_external_id.strip(): + raise ValueError("Accepted project change requires note_external_id") + if not self.permalink.strip(): + raise ValueError("Accepted project change requires permalink") + if not self.title.strip(): + raise ValueError("Accepted project change requires title") + if not self.file_path.strip(): + raise ValueError("Accepted project change requires file_path") + if self.accepted_at.tzinfo is None: + raise ValueError("Accepted project change accepted_at must be timezone-aware") + if not self.source.strip(): + raise ValueError("Accepted project change requires source") + if (self.db_version is None) != (self.db_checksum is None): + raise ValueError( + "Accepted project change revision requires both db_version and db_checksum" + ) + if self.db_version is not None and self.db_version <= 0: + raise ValueError("Accepted project change db_version must be positive") + if self.db_checksum is not None and not self.db_checksum.strip(): + raise ValueError("Accepted project change db_checksum must not be empty") + if self.operation == RuntimeProjectNoteOperation.moved and not self.previous_file_path: + raise ValueError("Moved accepted project change requires previous_file_path") diff --git a/tests/api/v2/test_accepted_note_atomicity.py b/tests/api/v2/test_accepted_note_atomicity.py index 339f3a409..b57f47d75 100644 --- a/tests/api/v2/test_accepted_note_atomicity.py +++ b/tests/api/v2/test_accepted_note_atomicity.py @@ -11,7 +11,14 @@ from basic_memory import db from basic_memory.deps.services import get_note_content_materialization_provider -from basic_memory.models import Entity, NoteContent, Observation, Project, Relation +from basic_memory.models import ( + AcceptedProjectNoteChange, + Entity, + NoteContent, + Observation, + Project, + Relation, +) from basic_memory.runtime.note_content import ( RuntimeAcceptedNoteChange, RuntimeNoteContentResponsePayload, @@ -24,6 +31,8 @@ class PersistedAcceptedSnapshot: entity: Entity note_content: NoteContent + project_partition_position: int + accepted_project_changes: tuple[AcceptedProjectNoteChange, ...] observations: tuple[Observation, ...] relations: tuple[Relation, ...] search_content: str @@ -36,6 +45,7 @@ async def _load_persisted_snapshot( entity_id: int, ) -> PersistedAcceptedSnapshot: async with db.scoped_session(session_maker) as session: + project = await session.get(Project, project_id) entity = await session.get(Entity, entity_id) note_content = await session.get(NoteContent, entity_id) observations = tuple( @@ -62,6 +72,15 @@ async def _load_persisted_snapshot( ) ).all() ) + accepted_project_changes = tuple( + ( + await session.scalars( + select(AcceptedProjectNoteChange) + .where(AcceptedProjectNoteChange.project_id == project_id) + .order_by(AcceptedProjectNoteChange.partition_position) + ) + ).all() + ) search_content = ( await session.execute( text(""" @@ -77,9 +96,12 @@ async def _load_persisted_snapshot( assert entity is not None assert note_content is not None + assert project is not None return PersistedAcceptedSnapshot( entity=entity, note_content=note_content, + project_partition_position=project.partition_position, + accepted_project_changes=accepted_project_changes, observations=observations, relations=relations, search_content=str(search_content), @@ -164,6 +186,11 @@ async def test_create_and_update_persist_complete_snapshot_at_materialization_bo assert created_snapshot.note_content.markdown_content == created.content assert created_snapshot.note_content.db_version == 1 assert created_snapshot.note_content.file_write_status == "pending" + assert created_snapshot.project_partition_position == 1 + assert [change.partition_position for change in created_snapshot.accepted_project_changes] == [ + 1 + ] + assert created_snapshot.accepted_project_changes[0].materialized_at is None assert [observation.content for observation in created_snapshot.observations] == [ "Create snapshot observation" ] @@ -199,6 +226,11 @@ async def test_create_and_update_persist_complete_snapshot_at_materialization_bo assert updated_snapshot.note_content.markdown_content == updated.content assert updated_snapshot.note_content.db_version == 2 assert updated_snapshot.note_content.file_write_status == "pending" + assert updated_snapshot.project_partition_position == 2 + assert [change.partition_position for change in updated_snapshot.accepted_project_changes] == [ + 1, + 2, + ] assert [observation.content for observation in updated_snapshot.observations] == [ "Replacing update observation" ] @@ -208,3 +240,12 @@ async def test_create_and_update_persist_complete_snapshot_at_materialization_bo assert "Replacing update observation" in updated_snapshot.search_content assert "Create snapshot observation" not in updated_snapshot.search_content assert len(materializer.accepted_changes) == 2 + created_change, updated_change = materializer.accepted_changes + assert created_change.project_change is not None + assert created_change.project_change.partition_position == 1 + assert created_change.materialization is not None + assert created_change.materialization.project_change is created_change.project_change + assert updated_change.project_change is not None + assert updated_change.project_change.partition_position == 2 + assert updated_change.materialization is not None + assert updated_change.materialization.project_change is updated_change.project_change diff --git a/tests/fixtures/wiki_projector/basic_projection.json b/tests/fixtures/wiki_projector/basic_projection.json index 0ef9090d1..da6de9e63 100644 --- a/tests/fixtures/wiki_projector/basic_projection.json +++ b/tests/fixtures/wiki_projector/basic_projection.json @@ -3,9 +3,9 @@ "project_id": "project-88", "through_partition_position": 3, "expected_sha256": { - "guides/index.md": "c6481bd17c663a3c2595cb35cd22a03da46c6d55465a274482a91363b3af244a", - "guides/log.md": "a3bf317e661d40a156a7481478938b8a652d4a69184b06cfa8c1b8b5355307fb", - "index.md": "33a4af87c4c1a3d4e128f780e67a8a7084aa1796118fb7a5334bbec4a4e0728d", - "log.md": "8e21ce176e930556f6a39f30412af7c488f9724b4118d2946879b72d0eb6c2f3" + "guides/index.md": "e0b43158cd8459fc441fb3d4f81100b0de8cf29b4ab324cc01a7b7c77e2c5b63", + "guides/log.md": "4f9a35fe972a1bcd7c38495e710b5bd9c712ec8e136033706f7436ecaee57031", + "index.md": "95771dff8c870604fdd4d4fce2d9a6dad8be1ea13f9ec56996f4b13218c378d2", + "log.md": "058f0f21ec2171d2e922c816ae97419cd2d942f3595db175d0b74a5b1cec9a97" } } diff --git a/tests/indexing/test_accepted_note_mutation_runner.py b/tests/indexing/test_accepted_note_mutation_runner.py index ce1f1683c..d92f5f3b3 100644 --- a/tests/indexing/test_accepted_note_mutation_runner.py +++ b/tests/indexing/test_accepted_note_mutation_runner.py @@ -2,7 +2,7 @@ from __future__ import annotations -from collections.abc import Iterator, Sequence +from collections.abc import Callable, Iterator, Sequence from dataclasses import dataclass from datetime import UTC, datetime from pathlib import Path @@ -50,6 +50,10 @@ from basic_memory.repository.observation_repository import ObservationGenerationWriteResult from basic_memory.repository.relation_repository import RelationGenerationWriteResult from basic_memory.runtime.note_content import RuntimeAcceptedNoteResponse +from basic_memory.runtime.project_partition import ( + RuntimeAcceptedProjectNoteChange, + RuntimeProjectNoteOperation, +) from basic_memory.schemas.base import Entity as EntitySchema from basic_memory.schemas.request import EditEntityRequest from basic_memory.services.exceptions import EntityAlreadyExistsError @@ -148,8 +152,10 @@ class _MutationSession: def __init__(self) -> None: self.deleted: list[object] = [] self.added: list[object] = [] + self.refreshed: list[object] = [] self.flush_count = 0 self.scalar_count = 0 + self.refresh_effect: Callable[[object], None] | None = None self.bind = SimpleNamespace(dialect=SimpleNamespace(name="sqlite")) async def delete(self, value: object) -> None: @@ -170,6 +176,11 @@ async def scalar(self, statement: object) -> int: async def flush(self) -> None: self.flush_count += 1 + async def refresh(self, value: object) -> None: + self.refreshed.append(value) + if self.refresh_effect is not None: + self.refresh_effect(value) + class _CreatePreparer: def __init__( @@ -331,9 +342,12 @@ async def load_current_file_checksum(self, project: Project, file_path: str) -> class _ProjectRepository: - def __init__(self, project: Project | None) -> None: + def __init__(self, project: Project | None, *, next_partition_position: int = 1) -> None: self.project = project self.calls: list[tuple[AsyncSession, str]] = [] + self.partition_calls: list[tuple[AsyncSession, int]] = [] + self.recorded_changes: list[RuntimeAcceptedProjectNoteChange] = [] + self.next_partition_position = next_partition_position async def get_by_external_id( self, @@ -343,6 +357,24 @@ async def get_by_external_id( self.calls.append((session, external_id)) return self.project + async def advance_partition_position( + self, + session: AsyncSession, + project_id: int, + ) -> int: + self.partition_calls.append((session, project_id)) + position = self.next_partition_position + self.next_partition_position += 1 + return position + + async def record_accepted_note_change( + self, + session: AsyncSession, + change: RuntimeAcceptedProjectNoteChange, + ) -> None: + _ = session + self.recorded_changes.append(change) + class _EntityLookupRepository: def __init__( @@ -772,6 +804,25 @@ async def test_run_accepted_note_create_persists_prepared_markdown( assert change.materialization.actor_kind == "user" assert change.materialization.actor_name == "Ada" assert change.materialization.previous_file_path is None + assert project_repository.partition_calls == [(session, project.id)] + assert change.project_change is not None + project_change = change.project_change + assert project_change.partition_position == 1 + assert project_change.operation is RuntimeProjectNoteOperation.created + assert project_change.project_external_id == "project-123" + assert project_change.note_external_id == "note-123" + assert project_change.permalink == "accepted" + assert project_change.file_path == "notes/accepted.md" + assert project_change.previous_file_path is None + assert project_change.accepted_at == _NOW + assert project_change.source == "api" + assert project_change.db_version == 1 + assert project_change.db_checksum == note_content.db_checksum + assert project_change.actor_user_profile_id == _ACTOR_ID + assert project_change.actor_kind == "user" + assert project_change.actor_name == "Ada" + assert project_repository.recorded_changes == [project_change] + assert change.materialization.project_change is project_change assert result.relation_publication is not None assert result.relation_publication.generation == 1 assert persistence_calls[0].await_count == 1 @@ -910,12 +961,61 @@ async def test_run_accepted_note_update_replaces_existing_note_content( assert change.materialization is not None assert change.materialization.db_version == 2 assert change.materialization.previous_file_path is None + assert project_repository.partition_calls == [(cast(AsyncSession, session), project.id)] + assert change.project_change is not None + assert change.project_change.operation is RuntimeProjectNoteOperation.moved + assert change.project_change.previous_file_path == "notes/accepted.md" + assert change.project_change.file_path == "notes/replacement.md" + assert change.project_change.db_version == 2 + assert change.materialization.project_change is change.project_change assert result.relation_publication is not None assert result.relation_publication.generation == 2 assert persistence_calls[0].await_count == 1 assert persistence_calls[1].await_count == 0 +@pytest.mark.asyncio +async def test_run_accepted_note_update_refreshes_source_path_after_lock() -> None: + session = _MutationSession() + schema = _schema() + project = _project() + entity = _entity(file_path="notes/accepted.md") + note_content = _note_content(entity) + + def refresh_after_concurrent_move(value: object) -> None: + if value is entity: + entity.file_path = "archive/accepted.md" + + session.refresh_effect = refresh_after_concurrent_move + project_repository = _ProjectRepository(project) + preparer = _CreatePreparer(_prepared()) + + result = await run_accepted_note_update( + cast(AsyncSession, session), + request=AcceptedNoteUpdateMutation( + project_external_id="project-123", + entity_external_id="note-123", + data=schema, + actor=AcceptedNoteMutationActor(user_profile_id=_ACTOR_ID), + source="api", + ), + dependencies=_dependencies( + project_repository=project_repository, + entity_lookup_repository=_EntityLookupRepository(by_external_id=entity), + note_content_lookup_repository=_NoteContentLookupRepository(note_content), + preparer_factory=_PreparerFactory(preparer), + pending_entity_repository=_PendingEntityRepository(entity), + note_content_accept_repository=_NoteContentAcceptRepository(note_content), + search_repository=_SearchRepository(), + ), + ) + + assert result.change.project_change is not None + assert result.change.project_change.operation is RuntimeProjectNoteOperation.moved + assert result.change.project_change.previous_file_path == "archive/accepted.md" + assert preparer.replace_calls[0][0].file_path == "notes/accepted.md" + + @pytest.mark.asyncio async def test_run_accepted_note_update_accepts_matching_base_checksum() -> None: # The caller's synced base ("old-checksum" in the fixture) still matches the @@ -1486,6 +1586,12 @@ async def test_run_accepted_note_edit_applies_patch_against_db_content( assert change.status_code == 200 assert change.materialization is not None assert change.materialization.source == "mcp" + assert project_repository.partition_calls == [(cast(AsyncSession, session), project.id)] + assert change.project_change is not None + assert change.project_change.operation is RuntimeProjectNoteOperation.updated + assert change.project_change.previous_file_path is None + assert change.project_change.actor_user_profile_id is None + assert change.materialization.project_change is change.project_change assert persistence_calls[0].await_count == 1 assert persistence_calls[1].await_count == 0 @@ -1646,6 +1752,15 @@ async def test_run_accepted_note_move_carries_previous_path_and_materialized_cle assert change.status_code == 200 assert change.materialization is not None assert change.materialization.previous_file_path == "notes/accepted.md" + assert project_repository.partition_calls == [(cast(AsyncSession, session), project.id)] + assert change.project_change is not None + assert change.project_change.operation is RuntimeProjectNoteOperation.moved + assert change.project_change.previous_file_path == "notes/accepted.md" + assert change.project_change.file_path == "archive/accepted.md" + assert change.project_change.actor_user_profile_id == _ACTOR_ID + assert change.project_change.actor_kind == "mcp" + assert change.project_change.actor_name == "Claude" + assert change.materialization.project_change is change.project_change cleanup = change.materialization.cleanup_after_write if expected_source_checksum is None: assert cleanup is None @@ -1701,6 +1816,55 @@ async def test_run_accepted_note_move_rejects_same_file_path() -> None: assert exc_info.value.rejection.detail == "Source and destination paths are the same." +@pytest.mark.asyncio +async def test_run_accepted_note_move_refreshes_source_path_after_lock( + persistence_calls: tuple[AsyncMock, AsyncMock], +) -> None: + session = _MutationSession() + project = _project() + prepared = _prepared_replacement() + prepared_move = _prepared_move() + entity = _entity(file_path="notes/original.md") + note_content = _note_content(entity) + project_repository = _ProjectRepository(project) + entity_lookup_repository = _EntityLookupRepository(by_external_id=entity) + note_content_lookup_repository = _NoteContentLookupRepository(note_content) + preparer = _CreatePreparer(prepared, prepared_move=prepared_move) + + def move_while_waiting(value: object) -> None: + if value is entity: + entity.file_path = "notes/intermediate.md" + + session.refresh_effect = move_while_waiting + + result = await run_accepted_note_move( + cast(AsyncSession, session), + request=AcceptedNoteMoveMutation( + project_external_id="project-123", + entity_external_id="note-123", + destination_path="archive/accepted.md", + actor=AcceptedNoteMutationActor(user_profile_id=None), + source="mcp", + ), + dependencies=_dependencies( + project_repository=project_repository, + entity_lookup_repository=entity_lookup_repository, + note_content_lookup_repository=note_content_lookup_repository, + preparer_factory=_PreparerFactory(preparer), + pending_entity_repository=_PendingEntityRepository(entity), + note_content_accept_repository=_NoteContentAcceptRepository(note_content), + search_repository=_SearchRepository(), + ), + ) + + assert session.refreshed == [entity] + assert result.change.project_change is not None + assert result.change.project_change.previous_file_path == "notes/intermediate.md" + assert result.change.materialization is not None + assert result.change.materialization.previous_file_path == "notes/intermediate.md" + assert persistence_calls[1].await_count == 1 + + @pytest.mark.asyncio async def test_run_accepted_note_create_resolves_directory_casing() -> None: """A unique case-insensitive folder match redirects the create (#1326).""" @@ -1987,15 +2151,60 @@ async def test_run_accepted_note_delete_removes_entity_and_returns_cleanup() -> assert session.deleted == [entity] assert search_repository.deleted_entity_ids == [entity.id] assert search_repository.deleted_vector_entity_ids == [entity.id] - assert session.scalar_count == 1 + assert session.scalar_count == 2 + assert session.refreshed == [entity, note_content] assert change.status_code == 200 assert change.file_delete is not None assert change.file_delete.file_path == "notes/accepted.md" assert change.file_delete.file_checksum == "file-checksum" + assert project_repository.partition_calls == [(cast(AsyncSession, session), project.id)] + assert change.project_change is not None + assert change.project_change.operation is RuntimeProjectNoteOperation.deleted + assert change.project_change.file_path == "notes/accepted.md" + assert change.project_change.source == "delete_note" + assert change.project_change.db_version == note_content.db_version + assert change.project_change.db_checksum == note_content.db_checksum + assert change.project_change.actor_user_profile_id is None + assert change.file_delete.project_change is change.project_change assert change.relation_cleanup_entity_ids == frozenset() assert result.relation_publication is None +@pytest.mark.asyncio +async def test_run_accepted_note_delete_stays_idempotent_after_concurrent_delete() -> None: + session = _MutationSession() + project = _project() + entity = _entity(file_path="notes/accepted.md") + entity_lookup_repository = _EntityLookupRepository(by_external_id=entity) + entity_lookup_repository.get_by_external_id = AsyncMock(side_effect=[entity, None]) + search_repository = _SearchRepository() + + result = await run_accepted_note_delete( + cast(AsyncSession, session), + request=AcceptedNoteDeleteMutation( + project_external_id="project-123", + entity_external_id="note-123", + ), + dependencies=_dependencies( + project_repository=_ProjectRepository(project), + entity_lookup_repository=entity_lookup_repository, + note_content_lookup_repository=_NoteContentLookupRepository(), + preparer_factory=_PreparerFactory(_CreatePreparer(_prepared())), + pending_entity_repository=_PendingEntityRepository(entity), + note_content_accept_repository=_NoteContentAcceptRepository(_note_content(entity)), + search_repository=search_repository, + ), + ) + + assert result.change.status_code == 200 + assert result.change.payload == {"deleted": False} + assert session.scalar_count == 1 + assert session.refreshed == [] + assert session.deleted == [] + assert search_repository.deleted_entity_ids == [] + assert search_repository.deleted_vector_entity_ids == [] + + def _prepared_with_graph( *, observations: Sequence[AcceptedObservationWrite], @@ -2073,6 +2282,57 @@ async def test_run_accepted_note_create_returns_graph_publication() -> None: assert result.relation_publication.relations[0].target_name == "XSYS Target" +@pytest.mark.asyncio +async def test_run_accepted_note_create_can_suppress_derived_graph_facts() -> None: + """Derived documents keep their Markdown without recursively expanding the graph.""" + session = cast(AsyncSession, object()) + prepared = _prepared_with_graph( + observations=[ + AcceptedObservationWrite( + content="Generated list item", + category="note", + context=None, + tags=None, + ) + ], + relations=[ + AcceptedRelationWrite( + relation_type="links_to", + target_name="Source Note", + context=None, + ) + ], + ) + entity = _entity() + note_content = _note_content(entity) + + result = await run_accepted_note_create( + session, + request=AcceptedNoteCreateMutation( + project_external_id="project-123", + data=_schema(), + actor=AcceptedNoteMutationActor(user_profile_id=None, kind="system"), + source="wiki_projector", + publish_graph_facts=False, + ), + dependencies=_dependencies( + project_repository=_ProjectRepository(_project()), + entity_lookup_repository=_EntityLookupRepository(), + note_content_lookup_repository=_NoteContentLookupRepository(), + preparer_factory=_PreparerFactory(_CreatePreparer(prepared)), + pending_entity_repository=_PendingEntityRepository(entity), + note_content_accept_repository=_NoteContentAcceptRepository(note_content), + search_repository=_SearchRepository(), + ), + ) + + assert isinstance(result.change.payload, RuntimeAcceptedNoteResponse) + assert result.change.payload.markdown_content == "# Accepted\n" + assert result.relation_publication is not None + assert result.relation_publication.observations == () + assert result.relation_publication.relations == () + + @pytest.mark.asyncio async def test_run_accepted_note_create_pre_resolves_only_unambiguous_self_links() -> None: """Safe self aliases resolve inline while ambiguous title aliases stay deferred.""" @@ -2177,6 +2437,58 @@ async def test_run_accepted_note_update_returns_replacement_graph() -> None: assert result.relation_publication.relations[0].target_name == "Other" +@pytest.mark.asyncio +async def test_run_accepted_note_update_can_clear_derived_graph_facts() -> None: + """A graph-silent replacement publishes empty sets so earlier facts are removed.""" + session = _MutationSession() + prepared = _prepared_with_graph( + observations=[ + AcceptedObservationWrite( + content="Generated list item", + category="note", + context=None, + tags=None, + ) + ], + relations=[ + AcceptedRelationWrite( + relation_type="links_to", + target_name="Source Note", + context=None, + ) + ], + ) + entity = _entity(file_path="notes/accepted.md") + note_content = _note_content(entity) + + result = await run_accepted_note_update( + cast(AsyncSession, session), + request=AcceptedNoteUpdateMutation( + project_external_id="project-123", + entity_external_id="note-123", + data=_schema(), + actor=AcceptedNoteMutationActor(user_profile_id=None, kind="system"), + source="wiki_projector", + publish_graph_facts=False, + ), + dependencies=_dependencies( + project_repository=_ProjectRepository(_project()), + entity_lookup_repository=_EntityLookupRepository(by_external_id=entity), + note_content_lookup_repository=_NoteContentLookupRepository(note_content), + preparer_factory=_PreparerFactory(_CreatePreparer(prepared)), + pending_entity_repository=_PendingEntityRepository(entity), + note_content_accept_repository=_NoteContentAcceptRepository(note_content), + search_repository=_SearchRepository(), + ), + ) + + assert isinstance(result.change.payload, RuntimeAcceptedNoteResponse) + assert result.change.payload.markdown_content == "# Accepted\n" + assert result.relation_publication is not None + assert result.relation_publication.observations == () + assert result.relation_publication.relations == () + + @pytest.mark.asyncio async def test_run_accepted_note_edit_returns_empty_replacement_graph() -> None: """An edit that drops the graph returns empty sets for fenced cleanup.""" diff --git a/tests/indexing/test_wiki_projector.py b/tests/indexing/test_wiki_projector.py index 3c23419be..fd6fe1b64 100644 --- a/tests/indexing/test_wiki_projector.py +++ b/tests/indexing/test_wiki_projector.py @@ -388,6 +388,8 @@ def test_projection_renders_root_and_affected_directory_indexes_and_logs() -> No assert "[[guides/index|Guides]]" in rendered["index.md"] assert "[[overview|Overview]]" in rendered["index.md"] assert "Updated [[guides/setup|Setup]]" in rendered["guides/log.md"] + assert "bm_parse_semantics: false" in rendered["index.md"] + assert "bm_parse_semantics: false" in rendered["log.md"] assert plan.result.source_watermark == 3 assert plan.result.output_watermark == 3 assert plan.result.created == 4 diff --git a/tests/repository/test_project_partition_repository.py b/tests/repository/test_project_partition_repository.py new file mode 100644 index 000000000..aec850900 --- /dev/null +++ b/tests/repository/test_project_partition_repository.py @@ -0,0 +1,233 @@ +"""Database regressions for strict project partition positions.""" + +from __future__ import annotations + +from datetime import UTC, datetime + +import pytest +from sqlalchemy import inspect, select +from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker + +from basic_memory import db +from basic_memory.models import AcceptedProjectNoteChange, Project +from basic_memory.repository.project_repository import ProjectRepository +from basic_memory.runtime.project_partition import ( + RuntimeAcceptedProjectNoteChange, + RuntimeProjectNoteOperation, +) + + +_ACCEPTED_AT = datetime(2026, 8, 29, 20, 15, tzinfo=UTC) + + +def test_project_delete_delegates_partition_journal_cleanup_to_database() -> None: + relationship = inspect(Project).relationships["accepted_note_changes"] + + assert relationship.passive_deletes is True + + +def _accepted_change( + project: Project, + *, + partition_position: int, + entity_id: int = 42, + note_external_id: str | None = None, + db_version: int = 3, +) -> RuntimeAcceptedProjectNoteChange: + return RuntimeAcceptedProjectNoteChange( + project_id=project.id, + project_external_id=project.external_id, + partition_position=partition_position, + entity_id=entity_id, + note_external_id=note_external_id or f"note-{entity_id}", + permalink=f"accepted-evidence-{entity_id}", + title="Accepted evidence", + operation=RuntimeProjectNoteOperation.updated, + file_path="notes/accepted-evidence.md", + accepted_at=_ACCEPTED_AT, + source="api", + db_version=db_version, + db_checksum=f"accepted-checksum-{db_version}", + ) + + +@pytest.mark.asyncio +async def test_project_partition_positions_advance_in_transaction_order( + session_maker: async_sessionmaker[AsyncSession], + test_project: Project, +) -> None: + repository = ProjectRepository() + + async with db.scoped_session(session_maker) as session: + first = await repository.advance_partition_position(session, test_project.id) + second = await repository.advance_partition_position(session, test_project.id) + + assert (first, second) == (1, 2) + async with db.scoped_session(session_maker) as session: + persisted = await session.get(Project, test_project.id) + assert persisted is not None + assert persisted.partition_position == 2 + + +@pytest.mark.asyncio +async def test_accepted_project_note_change_is_replayable_and_materialization_aware( + session_maker: async_sessionmaker[AsyncSession], + test_project: Project, +) -> None: + repository = ProjectRepository() + + async with db.scoped_session(session_maker) as session: + position = await repository.advance_partition_position(session, test_project.id) + await repository.record_accepted_note_change( + session, + _accepted_change(test_project, partition_position=position), + ) + + async with db.scoped_session(session_maker) as session: + changes = await repository.list_accepted_note_changes( + session, + test_project.id, + after_position=0, + through_position=1, + ) + assert len(changes) == 1 + assert changes[0].operation == RuntimeProjectNoteOperation.updated.value + assert changes[0].permalink == "accepted-evidence-42" + assert changes[0].db_checksum == "accepted-checksum-3" + assert changes[0].accepted_at.tzinfo is not None + assert changes[0].materialized_at is None + assert await repository.mark_accepted_note_change_materialized( + session, + test_project.id, + 1, + materialized_at=_ACCEPTED_AT, + ) + + async with db.scoped_session(session_maker) as session: + [materialized] = await repository.list_accepted_note_changes( + session, + test_project.id, + ) + assert materialized.materialized_at is not None + assert materialized.materialized_at.tzinfo is not None + assert not await repository.mark_accepted_note_change_materialized( + session, + test_project.id, + 1, + materialized_at=_ACCEPTED_AT, + ) + + +@pytest.mark.asyncio +async def test_materializing_newer_note_change_satisfies_superseded_positions( + session_maker: async_sessionmaker[AsyncSession], + test_project: Project, +) -> None: + repository = ProjectRepository() + + async with db.scoped_session(session_maker) as session: + for position, entity_id in ((1, 42), (2, 99), (3, 42)): + claimed_position = await repository.advance_partition_position( + session, + test_project.id, + ) + assert claimed_position == position + await repository.record_accepted_note_change( + session, + _accepted_change( + test_project, + partition_position=position, + entity_id=entity_id, + db_version=position, + ), + ) + + assert await repository.mark_accepted_note_change_materialized( + session, + test_project.id, + 3, + materialized_at=_ACCEPTED_AT, + ) + + async with db.scoped_session(session_maker) as session: + changes = await repository.list_accepted_note_changes(session, test_project.id) + assert [change.materialized_at is not None for change in changes] == [True, False, True] + + +@pytest.mark.asyncio +async def test_materializing_note_change_does_not_match_reused_entity_id( + session_maker: async_sessionmaker[AsyncSession], + test_project: Project, +) -> None: + repository = ProjectRepository() + + async with db.scoped_session(session_maker) as session: + for position, note_external_id in ((1, "deleted-note"), (2, "replacement-note")): + claimed_position = await repository.advance_partition_position( + session, + test_project.id, + ) + assert claimed_position == position + await repository.record_accepted_note_change( + session, + _accepted_change( + test_project, + partition_position=position, + entity_id=42, + note_external_id=note_external_id, + db_version=position, + ), + ) + + assert await repository.mark_accepted_note_change_materialized( + session, + test_project.id, + 2, + materialized_at=_ACCEPTED_AT, + ) + + async with db.scoped_session(session_maker) as session: + changes = await repository.list_accepted_note_changes(session, test_project.id) + assert [change.materialized_at is not None for change in changes] == [False, True] + + +@pytest.mark.asyncio +async def test_project_partition_position_rolls_back_with_rejected_transaction( + session_maker: async_sessionmaker[AsyncSession], + test_project: Project, +) -> None: + repository = ProjectRepository() + + with pytest.raises(RuntimeError, match="reject accepted mutation"): + async with db.scoped_session(session_maker) as session: + position = await repository.advance_partition_position(session, test_project.id) + await repository.record_accepted_note_change( + session, + _accepted_change(test_project, partition_position=position), + ) + raise RuntimeError("reject accepted mutation") + + async with db.scoped_session(session_maker) as session: + persisted = await session.get(Project, test_project.id) + assert persisted is not None + assert persisted.partition_position == 0 + rolled_back_changes = ( + await session.scalars( + select(AcceptedProjectNoteChange).where( + AcceptedProjectNoteChange.project_id == test_project.id + ) + ) + ).all() + assert rolled_back_changes == [] + assert await repository.advance_partition_position(session, test_project.id) == 1 + + +@pytest.mark.asyncio +async def test_project_partition_advance_rejects_missing_project( + session_maker: async_sessionmaker[AsyncSession], +) -> None: + repository = ProjectRepository() + + with pytest.raises(RuntimeError, match="project_id=999999"): + async with db.scoped_session(session_maker) as session: + await repository.advance_partition_position(session, 999999) diff --git a/tests/runtime/test_project_partition.py b/tests/runtime/test_project_partition.py new file mode 100644 index 000000000..301c205cc --- /dev/null +++ b/tests/runtime/test_project_partition.py @@ -0,0 +1,111 @@ +"""Portable project partition evidence and propagation tests.""" + +from __future__ import annotations + +from datetime import UTC, datetime +from uuid import UUID + +import pytest + +from basic_memory.runtime.cleanup import plan_note_file_delete_job_request +from basic_memory.runtime.note_content_deletes import RuntimePendingNoteFileDelete +from basic_memory.runtime.note_materialization_planning import ( + RuntimePendingNoteMaterialization, + plan_note_materialization_job_request, +) +from basic_memory.runtime.project_partition import ( + RuntimeAcceptedProjectNoteChange, + RuntimeProjectNoteOperation, +) + + +def _project_change( + operation: RuntimeProjectNoteOperation = RuntimeProjectNoteOperation.updated, +) -> RuntimeAcceptedProjectNoteChange: + return RuntimeAcceptedProjectNoteChange( + project_id=7, + project_external_id="project-123", + partition_position=4, + entity_id=42, + note_external_id="note-123", + permalink="accepted", + title="Accepted", + operation=operation, + file_path="notes/accepted.md", + accepted_at=datetime(2026, 8, 29, 12, tzinfo=UTC), + source="api", + previous_file_path=( + "notes/old.md" if operation is RuntimeProjectNoteOperation.moved else None + ), + db_version=3, + db_checksum="db-checksum", + actor_user_profile_id=UUID("11111111-1111-4111-8111-111111111111"), + actor_kind="user", + actor_name="Ada", + ) + + +def test_project_change_requires_complete_revision_identity() -> None: + with pytest.raises(ValueError, match="both db_version and db_checksum"): + RuntimeAcceptedProjectNoteChange( + project_id=7, + project_external_id="project-123", + partition_position=1, + entity_id=42, + note_external_id="note-123", + permalink="accepted", + title="Accepted", + operation=RuntimeProjectNoteOperation.updated, + file_path="notes/accepted.md", + accepted_at=datetime(2026, 8, 29, 12, tzinfo=UTC), + source="api", + db_version=3, + ) + + +def test_moved_project_change_requires_previous_path() -> None: + with pytest.raises(ValueError, match="requires previous_file_path"): + RuntimeAcceptedProjectNoteChange( + project_id=7, + project_external_id="project-123", + partition_position=1, + entity_id=42, + note_external_id="note-123", + permalink="accepted", + title="Accepted", + operation=RuntimeProjectNoteOperation.moved, + file_path="notes/accepted.md", + accepted_at=datetime(2026, 8, 29, 12, tzinfo=UTC), + source="api", + ) + + +def test_project_change_survives_materialization_job_flattening() -> None: + project_change = _project_change() + request = plan_note_materialization_job_request( + RuntimePendingNoteMaterialization( + project_id=7, + entity_id=42, + db_version=3, + db_checksum="db-checksum", + project_change=project_change, + source="api", + ) + ) + + assert request.project_change is project_change + + +def test_project_change_survives_file_delete_job_flattening() -> None: + project_change = _project_change(RuntimeProjectNoteOperation.deleted) + request = plan_note_file_delete_job_request( + RuntimePendingNoteFileDelete( + project_id=7, + entity_id=42, + file_path="notes/accepted.md", + file_checksum="file-checksum", + project_change=project_change, + ) + ) + + assert request.project_change is project_change diff --git a/tests/runtime/test_runtime_job_payloads.py b/tests/runtime/test_runtime_job_payloads.py index 0b5918b5c..96df54232 100644 --- a/tests/runtime/test_runtime_job_payloads.py +++ b/tests/runtime/test_runtime_job_payloads.py @@ -1,9 +1,11 @@ """Tests for portable runtime worker payload boundaries.""" +from datetime import UTC, datetime from uuid import UUID import pytest +from basic_memory.indexing.wiki_projector import WIKI_PROJECTOR_SOURCE from basic_memory.runtime.cleanup import RuntimeNoteFileDeleteJobRequest from basic_memory.runtime.job_payloads import ( DELETE_NOTE_FILE_ENTRYPOINT, @@ -13,7 +15,35 @@ ) from basic_memory.runtime.jobs import RuntimeJobRequest from basic_memory.runtime.note_content import RuntimeNoteMaterializationJobRequest -from basic_memory.runtime.note_object_metadata import NOTE_OBJECT_ACTOR_KIND_MCP_CLIENT +from basic_memory.runtime.note_object_metadata import ( + NOTE_OBJECT_ACTOR_KIND_MCP_CLIENT, + NOTE_OBJECT_ACTOR_KIND_SYSTEM, +) +from basic_memory.runtime.project_partition import ( + RuntimeAcceptedProjectNoteChange, + RuntimeProjectNoteOperation, +) + + +def _project_change() -> RuntimeAcceptedProjectNoteChange: + return RuntimeAcceptedProjectNoteChange( + project_id=101, + project_external_id="project-123", + partition_position=7, + entity_id=42, + note_external_id="note-123", + permalink="a", + title="A", + operation=RuntimeProjectNoteOperation.updated, + file_path="notes/a.md", + accepted_at=datetime(2026, 8, 29, 12, tzinfo=UTC), + source="mcp", + db_version=4, + db_checksum="db-sum", + actor_user_profile_id=UUID("33333333-3333-4333-8333-333333333333"), + actor_kind=NOTE_OBJECT_ACTOR_KIND_MCP_CLIENT, + actor_name="Claude Code", + ) def test_runtime_note_file_delete_job_payload_round_trips_runtime_request() -> None: @@ -23,6 +53,7 @@ def test_runtime_note_file_delete_job_payload_round_trips_runtime_request() -> N entity_id=42, file_path="notes/a.md", file_checksum="file-sum", + project_change=_project_change(), ) payload = RuntimeNoteFileDeleteJobPayload.from_runtime_request(runtime_request) @@ -62,6 +93,7 @@ def test_runtime_note_materialization_job_payload_round_trips_runtime_request() entity_id=42, db_version=4, db_checksum="db-sum", + project_change=_project_change(), actor_user_profile_id=UUID("33333333-3333-3333-3333-333333333333"), actor_kind=NOTE_OBJECT_ACTOR_KIND_MCP_CLIENT, actor_name="Claude Code", @@ -119,6 +151,23 @@ def test_runtime_note_materialization_job_payload_normalizes_origin_fields() -> assert payload.source == "mcp" +def test_runtime_note_materialization_job_payload_accepts_wiki_projector_source() -> None: + """Generated OKF notes preserve their projector source through materialization.""" + payload = RuntimeNoteMaterializationJobPayload( + project_id=101, + entity_id=42, + db_version=4, + db_checksum="db-sum", + actor_kind=NOTE_OBJECT_ACTOR_KIND_SYSTEM, + actor_name="Basic Memory Wiki Projector", + source=WIKI_PROJECTOR_SOURCE, + ) + + assert payload.actor_kind == NOTE_OBJECT_ACTOR_KIND_SYSTEM + assert payload.actor_name == "Basic Memory Wiki Projector" + assert payload.source == WIKI_PROJECTOR_SOURCE + + def test_runtime_note_materialization_job_payload_rejects_unknown_origin_fields() -> None: """Bad queued origins should fail before they become materialized file metadata.""" with pytest.raises(ValueError, match="unsupported note materialization actor kind"):