diff --git a/docs/plans/2026-05-04-062-store-mainline-facade-boundary-cleanup-plan.md b/docs/plans/2026-05-04-062-store-mainline-facade-boundary-cleanup-plan.md new file mode 100644 index 0000000..cc16d01 --- /dev/null +++ b/docs/plans/2026-05-04-062-store-mainline-facade-boundary-cleanup-plan.md @@ -0,0 +1,106 @@ +--- +status: completed +created: 2026-05-04 +origin: user request +scope: store mainline CortexMemory facade boundary cleanup +--- + +# Store Mainline Facade Boundary Cleanup + +## Problem + +The normal `/api/v1/memory/store` path now has a clear staged flow in +`MemoryWriteService.add`, but several write-path helper services still reach +through `CortexMemory` for storage, filesystem, embedding, record projection, +signals, and URI helpers. That makes the store mainline look like it still +depends on the top-level compatibility facade rather than on the write domain +boundary. + +The goal is to let the store mainline depend on `MemoryWriteService` as its +local boundary. `CortexMemory` should keep compatibility wrappers, but the +normal store helpers should stop calling `self._write_service._orch` directly. + +## Scope + +In scope: + +- Add narrow dependency accessors on + `src/opencortex/services/memory_write_service.py` for store helpers: + storage, collection, filesystem, embedder, config, signal bus, entity index, + record/URI helpers, and derived projection sync. +- Update normal store-path helper services to call the write-service boundary + instead of `CortexMemory`: + - `src/opencortex/services/memory_write_context_builder.py` + - `src/opencortex/services/memory_write_derive_service.py` + - `src/opencortex/services/memory_write_embed_service.py` + - `src/opencortex/services/memory_write_dedup_service.py` + - `src/opencortex/services/memory_directory_record_service.py` + - `src/opencortex/services/memory_store_record_service.py` +- Remove `_orch` properties from those helper services when no longer needed. +- Preserve `CortexMemory` public and compatibility wrappers. +- Preserve `/api/v1/memory/store` behavior, dedup behavior, projection sync, + parent directory records, signals, and CortexFS fire-and-forget semantics. + +Out of scope: + +- Update/remove mutation cleanup in `MemoryMutationService`. +- Document ingest cleanup in `MemoryDocumentWriteService`. +- Batch/document benchmark paths. +- Deleting `CortexMemory` compatibility methods. +- Changing storage filter, TTL, projection, or entity-index semantics. + +## Implementation Units + +### 1. Add Write-Service Boundary Methods + +Add focused accessors/delegates to `MemoryWriteService` for exactly the +dependencies needed by the normal store helpers: + +- initialization and collection lookup +- storage and filesystem +- embedder and config +- memory signal bus and entity index +- URI/category helpers +- abstract/object payload helpers +- derive helpers used by write derive +- record loading for explicit URI/dedup +- TTL and anchor projection sync + +These delegates may still call `CortexMemory` internally. The cleanup target is +that store helper services depend on `MemoryWriteService`, not the top-level +facade. + +### 2. Move Store Helpers Off `_orch` + +Change the normal store helper services listed in scope to call the new +write-service boundary. Keep argument order and result payloads unchanged. + +### 3. Update Tests to Assert the New Boundary + +Adjust focused tests such as `tests/test_memory_store_record_service.py` so the +test double represents `MemoryWriteService` directly instead of a nested +`_orch` object. Existing behavioral assertions should remain the same. + +## Test Plan + +Focused tests: + +- `uv run --group dev pytest tests/test_memory_store_record_service.py tests/test_memory_write_context_builder.py tests/test_memory_write_derive_service.py tests/test_memory_write_embed_service.py tests/test_memory_write_dedup_service.py tests/test_memory_directory_record_service.py -q` +- `uv run --group dev pytest tests/test_write_dedup.py tests/test_http_server.py -q` + +Static checks: + +- `uv run --group dev ruff format --check src/opencortex/services/memory_write_service.py src/opencortex/services/memory_write_context_builder.py src/opencortex/services/memory_write_derive_service.py src/opencortex/services/memory_write_embed_service.py src/opencortex/services/memory_write_dedup_service.py src/opencortex/services/memory_directory_record_service.py src/opencortex/services/memory_store_record_service.py tests/test_memory_store_record_service.py` +- `uv run --group dev ruff check src/opencortex/services/memory_write_service.py src/opencortex/services/memory_write_context_builder.py src/opencortex/services/memory_write_derive_service.py src/opencortex/services/memory_write_embed_service.py src/opencortex/services/memory_write_dedup_service.py src/opencortex/services/memory_directory_record_service.py src/opencortex/services/memory_store_record_service.py tests/test_memory_store_record_service.py` + +## Risks + +- Some focused tests may intentionally construct helper services with + `SimpleNamespace(_orch=...)`; update those tests only where the helper's + boundary changes. +- Dedup and parent-directory writes share storage and embedding behavior with + normal store. Preserve collection names, filter shapes, and best-effort + filesystem behavior. +- This is a dependency-boundary cleanup, not a semantics refactor. Do not + change scoring, derive output, TTL values, projection payloads, or signal + payloads. diff --git a/src/opencortex/services/cortex_memory_services.py b/src/opencortex/services/cortex_memory_services.py index 415c67a..186cad2 100644 --- a/src/opencortex/services/cortex_memory_services.py +++ b/src/opencortex/services/cortex_memory_services.py @@ -6,9 +6,9 @@ from typing import TYPE_CHECKING, Callable, TypeVar if TYPE_CHECKING: + from opencortex.cortex_memory import CortexMemory from opencortex.lifecycle.background_tasks import BackgroundTaskManager from opencortex.lifecycle.bootstrapper import SubsystemBootstrapper - from opencortex.cortex_memory import CortexMemory from opencortex.services.derivation_service import DerivationService from opencortex.services.knowledge_service import KnowledgeService from opencortex.services.memory_admin_stats_service import ( @@ -45,11 +45,42 @@ def _cached(self, attr_name: str, factory: Callable[[], T]) -> T: def memory_service(self) -> "MemoryService": """Lazy-built MemoryService for delegated CRUD/query/scoring methods.""" from opencortex.services.memory_service import MemoryService + from opencortex.services.memory_write_service import MemoryWriteDependencies return self._cached( "_memory_service_instance", - lambda: MemoryService(self._orch), + lambda: self._build_memory_service( + MemoryService, + MemoryWriteDependencies, + ), + ) + + def _build_memory_service( + self, + memory_service_type: type["MemoryService"], + dependencies_type: type, + ) -> "MemoryService": + """Construct MemoryService and bind explicit write-path dependencies.""" + service = memory_service_type(self._orch) + if not hasattr(self._orch, "_config"): + return service + service.configure_write_dependencies( + dependencies_type( + config=self._orch._config, + storage=self._orch._storage, + fs=self._orch._fs, + embedder=self._orch._embedder, + memory_signal_bus=getattr(self._orch, "_memory_signal_bus", None), + entity_index=getattr(self._orch, "_entity_index", None), + memory_record_service=self._orch._memory_record_service, + derivation_service=self._orch._derivation_service, + session_lifecycle_service=self._orch._session_lifecycle_service, + ensure_init=self._orch._ensure_init, + get_collection=self._orch._get_collection, + feedback=service.feedback, + ) ) + return service @property def derivation_service(self) -> "DerivationService": diff --git a/src/opencortex/services/memory_directory_record_service.py b/src/opencortex/services/memory_directory_record_service.py index c0b2e34..b89f9d8 100644 --- a/src/opencortex/services/memory_directory_record_service.py +++ b/src/opencortex/services/memory_directory_record_service.py @@ -26,10 +26,6 @@ def __init__(self, write_service: "MemoryWriteService") -> None: """Bind the directory service to a write service facade.""" self._write_service = write_service - @property - def _orch(self) -> Any: - return self._write_service._orch - async def ensure_parent_records(self, parent_uri: str) -> None: """Ensure all ancestor directory records exist in the vector store.""" to_create = await self._collect_missing_ancestors(parent_uri) @@ -49,7 +45,6 @@ async def ensure_parent_records(self, parent_uri: str) -> None: async def _collect_missing_ancestors(self, parent_uri: str) -> List[str]: """Walk upward and return directory URIs missing from storage.""" - orch = self._orch uri = parent_uri to_create: List[str] = [] @@ -59,8 +54,8 @@ async def _collect_missing_ancestors(self, parent_uri: str) -> List[str]: except ValueError: break - existing = await orch._storage.filter( - orch._get_collection(), + existing = await self._write_service._storage.filter( + self._write_service._get_collection(), FilterExpr.eq("uri", uri).to_dict(), limit=1, ) @@ -84,10 +79,9 @@ async def _create_directory_record( effective_user: UserIdentifier, ) -> None: """Build and upsert one directory record.""" - orch = self._orch dir_ctx = Context( uri=dir_uri, - parent_uri=orch._derive_parent_uri(dir_uri), + parent_uri=self._write_service._derive_parent_uri(dir_uri), is_leaf=False, abstract="", user=effective_user, @@ -106,12 +100,14 @@ async def _create_directory_record( record["mergeable"] = False record["session_id"] = "" record["ttl_expires_at"] = "" - await orch._storage.upsert(orch._get_collection(), record) + await self._write_service._storage.upsert( + self._write_service._get_collection(), record + ) logger.debug("[MemoryService] Created directory record: %s", dir_uri) async def _embed_directory_name(self, *, dir_ctx: Context, uri: str) -> Any: """Embed the directory basename and attach its dense vector.""" - embedder = self._orch._embedder + embedder = self._write_service._embedder dir_name = uri.rstrip("/").rsplit("/", 1)[-1] if not embedder or not dir_name: return None diff --git a/src/opencortex/services/memory_service.py b/src/opencortex/services/memory_service.py index b62bd4f..5ea8e36 100644 --- a/src/opencortex/services/memory_service.py +++ b/src/opencortex/services/memory_service.py @@ -55,7 +55,10 @@ from opencortex.cortex_memory import CortexMemory from opencortex.services.memory_query_service import MemoryQueryService from opencortex.services.memory_scoring_service import MemoryScoringService - from opencortex.services.memory_write_service import MemoryWriteService + from opencortex.services.memory_write_service import ( + MemoryWriteDependencies, + MemoryWriteService, + ) _BATCH_ADD_CONCURRENCY = 8 _BATCH_ADD_TASK_CHUNK_SIZE = _BATCH_ADD_CONCURRENCY * 4 @@ -79,6 +82,14 @@ def __init__(self, orchestrator: "CortexMemory") -> None: at call time. Stored as ``self._orch``; not validated. """ self._orch = orchestrator + self._write_dependencies: "MemoryWriteDependencies | None" = None + + def configure_write_dependencies( + self, + dependencies: "MemoryWriteDependencies", + ) -> None: + """Bind explicit write-path dependencies from the service registry.""" + self._write_dependencies = dependencies @property def _memory_write_service(self) -> "MemoryWriteService": @@ -87,7 +98,9 @@ def _memory_write_service(self) -> "MemoryWriteService": cached = getattr(self, "_memory_write_service_instance", None) if cached is None: - cached = MemoryWriteService(self) + if self._write_dependencies is None: + raise RuntimeError("Memory write dependencies are not configured") + cached = MemoryWriteService(self, self._write_dependencies) self._memory_write_service_instance = cached return cached diff --git a/src/opencortex/services/memory_store_record_service.py b/src/opencortex/services/memory_store_record_service.py index 8a7d26e..ad80453 100644 --- a/src/opencortex/services/memory_store_record_service.py +++ b/src/opencortex/services/memory_store_record_service.py @@ -34,10 +34,6 @@ def __init__(self, write_service: "MemoryWriteService") -> None: """Bind the persistence service to a write service facade.""" self._write_service = write_service - @property - def _orch(self) -> Any: - return self._write_service._orch - async def persist_context_record( self, *, @@ -57,7 +53,6 @@ async def persist_context_record( is_leaf: bool, ) -> StoredRecordResult: """Assemble and persist a normal store record.""" - orch = self._orch record = ctx.to_dict() if ctx.vector: record["vector"] = ctx.vector @@ -85,9 +80,11 @@ async def persist_context_record( self._populate_flattened_source_fields(record, meta) upsert_started = asyncio.get_running_loop().time() - await orch._storage.upsert(orch._get_collection(), record) + await self._write_service._storage.upsert( + self._write_service._get_collection(), record + ) upsert_ms = int((asyncio.get_running_loop().time() - upsert_started) * 1000) - await orch._sync_anchor_projection_records( + await self._write_service._sync_anchor_projection_records( source_record=record, abstract_json=abstract_json, ) @@ -120,15 +117,18 @@ def _ttl_for_record( meta: Dict[str, Any], ) -> str: """Return the TTL string for short-lived record kinds.""" - orch = self._orch if context_type == "staging": - return orch._ttl_from_hours(orch._config.immediate_event_ttl_hours) + return self._write_service._ttl_from_hours( + self._write_service._config.immediate_event_ttl_hours + ) if ( (context_type or "memory") == "memory" and effective_category == "events" and meta.get("layer") == "merged" ): - return orch._ttl_from_hours(orch._config.merged_event_ttl_hours) + return self._write_service._ttl_from_hours( + self._write_service._config.merged_event_ttl_hours + ) return "" @staticmethod @@ -156,7 +156,7 @@ def _publish_memory_stored( effective_category: str, ) -> None: """Publish the post-store lifecycle signal when a bus exists.""" - signal_bus = getattr(self._orch, "_memory_signal_bus", None) + signal_bus = self._write_service._memory_signal_bus if signal_bus is None: return signal_bus.publish_nowait( @@ -179,9 +179,11 @@ def _sync_entity_index( entities: List[str], ) -> None: """Sync the entity index for entity-bearing records.""" - entity_index = getattr(self._orch, "_entity_index", None) + entity_index = self._write_service._entity_index if entity_index and entities: - entity_index.add(self._orch._get_collection(), str(record["id"]), entities) + entity_index.add( + self._write_service._get_collection(), str(record["id"]), entities + ) def _schedule_cortexfs_write( self, @@ -207,7 +209,7 @@ def _on_fs_done(task: asyncio.Task[Any]) -> None: ) fs_task = asyncio.create_task( - self._orch._fs.write_context( + self._write_service._fs.write_context( uri=uri, content=content, abstract=abstract, diff --git a/src/opencortex/services/memory_write_context_builder.py b/src/opencortex/services/memory_write_context_builder.py index 5e6cf12..33092fa 100644 --- a/src/opencortex/services/memory_write_context_builder.py +++ b/src/opencortex/services/memory_write_context_builder.py @@ -55,10 +55,6 @@ def __init__(self, write_service: "MemoryWriteService") -> None: """Bind the builder to a write service facade.""" self._write_service = write_service - @property - def _orch(self) -> Any: - return self._write_service._orch - async def resolve_target( self, *, @@ -70,24 +66,25 @@ async def resolve_target( uri: Optional[str], ) -> ResolvedWriteTarget: """Resolve URI, parent URI, existing record, and explicit metadata.""" - orch = self._orch resolved_meta = dict(meta or {}) explicit_entities = _merge_unique_strings(resolved_meta.get("entities")) explicit_topics = _merge_unique_strings(resolved_meta.get("topics")) if not uri: - resolved_uri = orch._auto_uri( + resolved_uri = self._write_service._auto_uri( context_type or "memory", category, abstract=abstract, ) - resolved_uri = await orch._resolve_unique_uri(resolved_uri) + resolved_uri = await self._write_service._resolve_unique_uri(resolved_uri) existing_record = None else: resolved_uri = uri - existing_record = await orch._get_record_by_uri(resolved_uri) + existing_record = await self._write_service._get_record_by_uri(resolved_uri) - resolved_parent_uri = parent_uri or orch._derive_parent_uri(resolved_uri) + resolved_parent_uri = parent_uri or self._write_service._derive_parent_uri( + resolved_uri + ) return ResolvedWriteTarget( uri=resolved_uri, parent_uri=resolved_parent_uri, @@ -113,7 +110,6 @@ def assemble_context( layers: Dict[str, Any], ) -> AssembledWriteContext: """Assemble the post-derive Context, metadata, and object payload.""" - orch = self._orch meta = target.meta derived_entities = layers.get("entities", []) if content and is_leaf else [] entities = _merge_unique_strings(derived_entities, target.explicit_entities) @@ -159,8 +155,10 @@ def assemble_context( elif embed_text: ctx.vectorize = Vectorize(embed_text) - effective_category = category or orch._extract_category_from_uri(target.uri) - abstract_json = orch._build_abstract_json( + effective_category = category or self._write_service._extract_category_from_uri( + target.uri + ) + abstract_json = self._write_service._build_abstract_json( uri=target.uri, context_type=context_type or "", category=effective_category, @@ -175,7 +173,9 @@ def assemble_context( ) if content and is_leaf: abstract_json["fact_points"] = layers.get("fact_points", []) - object_payload = orch._memory_object_payload(abstract_json, is_leaf=is_leaf) + object_payload = self._write_service._memory_object_payload( + abstract_json, is_leaf=is_leaf + ) return AssembledWriteContext( ctx=ctx, abstract=abstract, diff --git a/src/opencortex/services/memory_write_dedup_service.py b/src/opencortex/services/memory_write_dedup_service.py index f82d754..28ff840 100644 --- a/src/opencortex/services/memory_write_dedup_service.py +++ b/src/opencortex/services/memory_write_dedup_service.py @@ -4,12 +4,17 @@ from __future__ import annotations import asyncio +import json import logging from dataclasses import dataclass, field from typing import TYPE_CHECKING, Any, Dict, List, Optional, Tuple from opencortex.core.context import Context from opencortex.http.request_context import get_effective_project_id +from opencortex.services.derivation_service import ( + _merge_unique_strings, + _split_keyword_string, +) from opencortex.services.memory_filters import ( FilterExpr, and_filter, @@ -43,14 +48,6 @@ def __init__(self, write_service: "MemoryWriteService") -> None: """Bind the dedup service to a write service facade.""" self._write_service = write_service - @property - def _orch(self) -> Any: - return self._write_service._orch - - @property - def _service(self) -> Any: - return self._write_service._service - async def try_merge_duplicate( self, *, @@ -83,7 +80,7 @@ async def try_merge_duplicate( total_ms_at_match = int( (asyncio.get_running_loop().time() - add_started) * 1000 ) - existing_record = await self._orch._get_record_by_uri(existing_uri) + existing_record = await self._write_service._get_record_by_uri(existing_uri) persisted_owner_id = "" persisted_project_id = get_effective_project_id() if existing_record: @@ -125,7 +122,6 @@ async def check_duplicate( uid: str, ) -> Optional[Tuple[str, float]]: """Return duplicate ``(existing_uri, score)`` when one exists.""" - orch = self._orch try: dedup_filter = self._build_duplicate_filter( memory_kind=memory_kind, @@ -133,8 +129,8 @@ async def check_duplicate( tid=tid, uid=uid, ) - results = await orch._storage.search( - orch._get_collection(), + results = await self._write_service._storage.search( + self._write_service._get_collection(), query_vector=vector, filter=dedup_filter, limit=1, @@ -152,9 +148,8 @@ async def merge_into( self, existing_uri: str, new_abstract: str, new_content: str ) -> None: """Merge new content into an existing record and reinforce it.""" - orch = self._orch - records = await orch._storage.filter( - orch._get_collection(), + records = await self._write_service._storage.filter( + self._write_service._get_collection(), FilterExpr.eq("uri", existing_uri).to_dict(), limit=1, output_fields=["abstract", "overview"], @@ -162,7 +157,7 @@ async def merge_into( existing_content = "" if records: try: - existing_content = await orch._fs.read_file(existing_uri) + existing_content = await self._write_service._fs.read_file(existing_uri) except Exception: existing_content = "" @@ -171,12 +166,147 @@ async def merge_into( if new_content else existing_content ) - await self._service.update( - existing_uri, - abstract=new_abstract, - content=merged_content, + if records: + await self._update_merged_record( + existing_uri=existing_uri, + record=records[0], + abstract=new_abstract, + content=merged_content, + ) + await self._write_service.feedback(existing_uri, 0.5) + + async def _update_merged_record( + self, + *, + existing_uri: str, + record: Dict[str, Any], + abstract: str, + content: str, + ) -> None: + """Apply the minimal record update needed by dedup merge.""" + record_id = record.get("id", "") + if not record_id: + return + + next_meta = self._coerce_meta(record.get("meta", {})) + next_overview = str(record.get("overview", "") or "") + next_entities = _merge_unique_strings( + record.get("entities") or [], + next_meta.get("entities"), + ) + next_keywords_list = _merge_unique_strings( + next_meta.get("topics"), + _split_keyword_string(record.get("keywords", "")), ) - await self._service.feedback(existing_uri, 0.5) + derived_fact_points: Optional[List[str]] = None + if content: + derive_result = await self._write_service._derive_layers( + user_abstract=abstract, + content=content, + user_overview="", + ) + next_entities = _merge_unique_strings( + derive_result.get("entities", []), + next_entities, + ) + next_keywords_list = _merge_unique_strings( + next_keywords_list, + _split_keyword_string(derive_result.get("keywords", "")), + ) + next_anchor_handles = _merge_unique_strings( + next_meta.get("anchor_handles"), + derive_result.get("anchor_handles", []), + ) + if next_anchor_handles: + next_meta["anchor_handles"] = next_anchor_handles + raw_fps = derive_result.get("fact_points", []) + derived_fact_points = ( + [str(fp) for fp in raw_fps] if isinstance(raw_fps, list) else [] + ) + + update_data: Dict[str, Any] = {"abstract": abstract} + if next_keywords_list: + next_meta["topics"] = _merge_unique_strings( + next_meta.get("topics"), + next_keywords_list, + ) + update_data["keywords"] = ", ".join(next_keywords_list) + if next_entities: + update_data["entities"] = next_entities + if next_meta: + update_data["meta"] = next_meta + + embedder = self._write_service._embedder + if embedder: + loop = asyncio.get_running_loop() + embed_input = abstract + if next_keywords_list: + embed_input = f"{embed_input} {', '.join(next_keywords_list)}".strip() + result = await loop.run_in_executor(None, embedder.embed, embed_input) + update_data["vector"] = result.dense_vector + if result.sparse_vector: + update_data["sparse_vector"] = result.sparse_vector + + abstract_json = self._write_service._build_abstract_json( + uri=existing_uri, + context_type=str(record.get("context_type", "") or ""), + category=str(record.get("category", "") or ""), + abstract=abstract, + overview=next_overview, + content=content, + entities=next_entities, + meta=next_meta, + keywords=next_keywords_list, + parent_uri=str(record.get("parent_uri", "") or ""), + session_id=str(record.get("session_id", "") or ""), + ) + if derived_fact_points is not None: + abstract_json["fact_points"] = derived_fact_points + else: + prior_abstract_json = record.get("abstract_json") + if isinstance(prior_abstract_json, dict): + prior_fps = prior_abstract_json.get("fact_points") or [] + if isinstance(prior_fps, list): + abstract_json["fact_points"] = [str(fp) for fp in prior_fps] + update_data.update( + self._write_service._memory_object_payload( + abstract_json, + is_leaf=bool(record.get("is_leaf", False)), + ) + ) + update_data["abstract_json"] = abstract_json + + await self._write_service._storage.update( + self._write_service._get_collection(), + record_id, + update_data, + ) + updated_record = dict(record) + updated_record.update(update_data) + await self._write_service._sync_anchor_projection_records( + source_record=updated_record, + abstract_json=abstract_json, + ) + await self._write_service._fs.write_context( + uri=existing_uri, + content=content, + abstract=abstract, + overview=next_overview, + abstract_json=abstract_json, + ) + + @staticmethod + def _coerce_meta(raw_meta: Any) -> Dict[str, Any]: + """Return metadata as a dict, tolerating legacy encoded payloads.""" + if isinstance(raw_meta, str): + try: + decoded = json.loads(raw_meta) + except (json.JSONDecodeError, TypeError): + return {} + return decoded if isinstance(decoded, dict) else {} + if isinstance(raw_meta, dict): + return dict(raw_meta) + return {} @staticmethod def _build_duplicate_filter( @@ -214,7 +344,7 @@ def _publish_merge_signal( existing_record: Dict[str, Any], ) -> None: """Publish the dedup merge lifecycle signal when a bus exists.""" - signal_bus = getattr(self._orch, "_memory_signal_bus", None) + signal_bus = self._write_service._memory_signal_bus if signal_bus is None: return signal_bus.publish_nowait( diff --git a/src/opencortex/services/memory_write_derive_service.py b/src/opencortex/services/memory_write_derive_service.py index 7c8a48c..b2a9673 100644 --- a/src/opencortex/services/memory_write_derive_service.py +++ b/src/opencortex/services/memory_write_derive_service.py @@ -28,10 +28,6 @@ def __init__(self, write_service: "MemoryWriteService") -> None: """Bind the derive service to a write service facade.""" self._write_service = write_service - @property - def _orch(self) -> Any: - return self._write_service._orch - async def derive_for_write( self, *, @@ -42,10 +38,9 @@ async def derive_for_write( defer_derive: bool, ) -> MemoryWriteDeriveResult: """Derive or fallback-fill write summary fields for normal add().""" - orch = self._orch if content and is_leaf and not defer_derive: derive_started = asyncio.get_running_loop().time() - layers = await orch._derive_layers( + layers = await self._write_service._derive_layers( user_abstract=abstract, content=content, user_overview=overview, @@ -63,13 +58,13 @@ async def derive_for_write( if content and is_leaf and defer_derive: resolved_overview = overview if not resolved_overview: - resolved_overview = orch._fallback_overview_from_content( + resolved_overview = self._write_service._fallback_overview_from_content( user_overview=overview, content=content, ) resolved_abstract = abstract if not resolved_abstract: - resolved_abstract = orch._derive_abstract_from_overview( + resolved_abstract = self._write_service._derive_abstract_from_overview( user_abstract=abstract, overview=resolved_overview, content=content, diff --git a/src/opencortex/services/memory_write_embed_service.py b/src/opencortex/services/memory_write_embed_service.py index cc27840..9e3f3a4 100644 --- a/src/opencortex/services/memory_write_embed_service.py +++ b/src/opencortex/services/memory_write_embed_service.py @@ -28,13 +28,9 @@ def __init__(self, write_service: "MemoryWriteService") -> None: """Bind the embed service to a write service facade.""" self._write_service = write_service - @property - def _orch(self) -> Any: - return self._write_service._orch - async def embed_for_write(self, ctx: Context) -> MemoryWriteEmbedResult: """Embed a normal write context and attach its dense vector.""" - embedder = self._orch._embedder + embedder = self._write_service._embedder if not embedder: return MemoryWriteEmbedResult() diff --git a/src/opencortex/services/memory_write_service.py b/src/opencortex/services/memory_write_service.py index 280d98c..39029a8 100644 --- a/src/opencortex/services/memory_write_service.py +++ b/src/opencortex/services/memory_write_service.py @@ -9,6 +9,7 @@ import asyncio import logging +from dataclasses import dataclass from typing import TYPE_CHECKING, Any, Dict, List, Optional, Tuple from opencortex.core.context import Context @@ -35,15 +36,193 @@ logger = logging.getLogger(__name__) +@dataclass(frozen=True) +class MemoryWriteDependencies: + """Explicit subsystem bundle used by the normal store write path.""" + + config: Any + storage: Any + fs: Any + embedder: Any + memory_signal_bus: Any + entity_index: Any + memory_record_service: Any + derivation_service: Any + session_lifecycle_service: Any + ensure_init: Any + get_collection: Any + feedback: Any + + class MemoryWriteService: """Own memory write/mutation logic behind the MemoryService facade.""" - def __init__(self, memory_service: "MemoryService") -> None: + def __init__( + self, + memory_service: "MemoryService", + dependencies: MemoryWriteDependencies, + ) -> None: self._service = memory_service + self._deps = dependencies + + @property + def _config(self) -> Any: + """Cortex configuration for write-path helpers.""" + return self._deps.config + + @property + def _storage(self) -> Any: + """Vector storage owned by the memory facade.""" + return self._deps.storage + + @property + def _fs(self) -> Any: + """CortexFS instance owned by the memory facade.""" + return self._deps.fs + + @property + def _embedder(self) -> Any: + """Embedder used by normal write-path helpers.""" + return self._deps.embedder + + @property + def _memory_signal_bus(self) -> Any: + """Optional lifecycle signal bus for write-path notifications.""" + return self._deps.memory_signal_bus @property - def _orch(self) -> Any: - return self._service._orch + def _entity_index(self) -> Any: + """Optional entity index for write-path synchronization.""" + return self._deps.entity_index + + def _ensure_init(self) -> None: + """Require the parent memory facade to be initialized.""" + self._deps.ensure_init() + + def _get_collection(self) -> str: + """Return the active vector-store collection.""" + return self._deps.get_collection() + + def _auto_uri(self, context_type: str, category: str, abstract: str = "") -> str: + """Generate a memory URI through the record service boundary.""" + return self._deps.memory_record_service._auto_uri( + context_type=context_type, + category=category, + abstract=abstract, + ) + + async def _resolve_unique_uri(self, uri: str) -> str: + """Resolve one URI to a unique value.""" + return await self._deps.memory_record_service._resolve_unique_uri(uri) + + async def _get_record_by_uri(self, uri: str) -> Optional[Dict[str, Any]]: + """Load one record by URI through the session/record boundary.""" + return await self._deps.session_lifecycle_service._get_record_by_uri(uri) + + def _derive_parent_uri(self, uri: str) -> str: + """Derive the parent URI for a memory URI.""" + return self._deps.memory_record_service._derive_parent_uri(uri) + + def _extract_category_from_uri(self, uri: str) -> str: + """Extract the memory category from a URI.""" + return self._deps.memory_record_service._extract_category_from_uri(uri) + + def _build_abstract_json( + self, + *, + uri: str, + context_type: str, + category: str, + abstract: str, + overview: str, + content: str, + entities: List[str], + meta: Optional[Dict[str, Any]], + keywords: Optional[List[str]] = None, + parent_uri: str, + session_id: Optional[str], + ) -> Dict[str, Any]: + """Build the canonical abstract payload for a write record.""" + return self._deps.memory_record_service._build_abstract_json( + uri=uri, + context_type=context_type, + category=category, + abstract=abstract, + overview=overview, + content=content, + entities=entities, + meta=meta, + keywords=keywords, + parent_uri=parent_uri, + session_id=session_id or "", + ) + + def _memory_object_payload( + self, + abstract_json: Dict[str, Any], + *, + is_leaf: bool, + ) -> Dict[str, Any]: + """Project abstract payload into flat memory object fields.""" + return self._deps.memory_record_service._memory_object_payload( + abstract_json, is_leaf=is_leaf + ) + + async def _derive_layers( + self, + *, + user_abstract: str, + content: str, + user_overview: str, + ) -> Dict[str, Any]: + """Derive memory layers for write-path content.""" + return await self._deps.derivation_service._derive_layers( + user_abstract=user_abstract, + content=content, + user_overview=user_overview, + ) + + def _fallback_overview_from_content( + self, + *, + user_overview: str, + content: str, + ) -> str: + """Build a deterministic fallback overview for deferred derive.""" + return self._deps.derivation_service._fallback_overview_from_content( + user_overview=user_overview, + content=content, + ) + + def _derive_abstract_from_overview( + self, + *, + user_abstract: str, + overview: str, + content: str, + ) -> str: + """Build a deterministic fallback abstract for deferred derive.""" + return self._deps.derivation_service._derive_abstract_from_overview( + user_abstract=user_abstract, + overview=overview, + content=content, + ) + + def _ttl_from_hours(self, hours: int) -> str: + """Return the TTL string for a write-path record.""" + return self._deps.memory_record_service._ttl_from_hours(hours) + + async def _sync_anchor_projection_records( + self, + *, + source_record: Dict[str, Any], + abstract_json: Dict[str, Any], + ) -> None: + """Synchronize derived anchor/fact projection records.""" + await self._deps.memory_record_service._sync_anchor_projection_records( + source_record=source_record, + abstract_json=abstract_json, + ) # ========================================================================= # CRUD (U2 of plan 010) @@ -141,8 +320,7 @@ async def add( The created ``Context`` with ``meta["dedup_action"]`` set to ``"created"`` or ``"merged"``. """ - orch = self._orch - orch._ensure_init() + self._ensure_init() # Determine ingestion mode from opencortex.ingest.resolver import IngestModeResolver @@ -156,7 +334,7 @@ async def add( # Document mode: parse -> chunks -> write each with hierarchy if ingest_mode == "document" and content and is_leaf: - return await self._service._add_document( + return await self._document_write_service._add_document( content=content, abstract=abstract, overview=overview, @@ -265,7 +443,7 @@ async def add( # Ensure parent directory records exist in vector DB if is_leaf and parent_uri: - await self._service._ensure_parent_records(parent_uri) + await self._ensure_parent_records(parent_uri) store_result = await self._store_record_service.persist_context_record( ctx=ctx, @@ -413,6 +591,10 @@ async def _merge_into( new_content=new_content, ) + async def feedback(self, uri: str, reward: float) -> None: + """Apply scoring feedback for write-time merge reinforcement.""" + await self._deps.feedback(uri, reward) + async def _ensure_parent_records(self, parent_uri: str) -> None: """Ensure all ancestor directory records exist in the vector store.""" await self._directory_record_service.ensure_parent_records(parent_uri) diff --git a/tests/test_memory_directory_record_service.py b/tests/test_memory_directory_record_service.py index d6f271f..317f606 100644 --- a/tests/test_memory_directory_record_service.py +++ b/tests/test_memory_directory_record_service.py @@ -40,7 +40,7 @@ def _build_service( filter=AsyncMock(side_effect=filter_results), upsert=AsyncMock(), ) - orch = SimpleNamespace( + write_service = SimpleNamespace( _storage=storage, _embedder=embedder, _get_collection=MagicMock(return_value="context"), @@ -48,8 +48,7 @@ def _build_service( side_effect=lambda uri: uri.rsplit("/", 1)[0] if "/" in uri else None ), ) - write_service = SimpleNamespace(_orch=orch) - return MemoryDirectoryRecordService(write_service), orch + return MemoryDirectoryRecordService(write_service), write_service async def test_missing_ancestors_are_created_top_down(self) -> None: """Missing directory ancestors are upserted from root to leaf.""" diff --git a/tests/test_memory_store_record_service.py b/tests/test_memory_store_record_service.py index 47aaf55..77421ca 100644 --- a/tests/test_memory_store_record_service.py +++ b/tests/test_memory_store_record_service.py @@ -40,7 +40,7 @@ def _build_service(self) -> tuple[MemoryStoreRecordService, Any]: fs.write_context = AsyncMock() signal_bus = _SignalBus() entity_index = MagicMock() - orch = SimpleNamespace( + write_service = SimpleNamespace( _storage=storage, _fs=fs, _memory_signal_bus=signal_bus, @@ -53,8 +53,7 @@ def _build_service(self) -> tuple[MemoryStoreRecordService, Any]: _sync_anchor_projection_records=AsyncMock(), _ttl_from_hours=MagicMock(side_effect=lambda hours: f"ttl:{hours}"), ) - write_service = SimpleNamespace(_orch=orch) - return MemoryStoreRecordService(write_service), orch + return MemoryStoreRecordService(write_service), write_service async def test_persist_context_record_assembles_and_persists_record(self) -> None: """The service owns normal store record payload construction.""" diff --git a/tests/test_memory_write_context_builder.py b/tests/test_memory_write_context_builder.py index d99b61c..5c1dafb 100644 --- a/tests/test_memory_write_context_builder.py +++ b/tests/test_memory_write_context_builder.py @@ -30,7 +30,7 @@ def memory_object_payload( "mergeable": is_leaf, } - orch = SimpleNamespace( + write_service = SimpleNamespace( _auto_uri=MagicMock( return_value="opencortex://tenant/user/memories/preferences/generated" ), @@ -45,8 +45,7 @@ def memory_object_payload( _build_abstract_json=MagicMock(side_effect=build_abstract_json), _memory_object_payload=MagicMock(side_effect=memory_object_payload), ) - write_service = SimpleNamespace(_orch=orch) - return MemoryWriteContextBuilder(write_service), orch + return MemoryWriteContextBuilder(write_service), write_service async def test_resolve_target_auto_uri_and_explicit_metadata(self) -> None: """Auto URI resolution copies meta and extracts explicit fields.""" diff --git a/tests/test_memory_write_dedup_service.py b/tests/test_memory_write_dedup_service.py index bfa528a..aafd701 100644 --- a/tests/test_memory_write_dedup_service.py +++ b/tests/test_memory_write_dedup_service.py @@ -39,12 +39,33 @@ def _build_service( ) -> tuple[MemoryWriteDedupService, Any, Any]: storage = MagicMock() storage.search = AsyncMock(return_value=search_results or []) - storage.filter = AsyncMock(return_value=[{"uri": "target-uri"}]) + storage.filter = AsyncMock( + return_value=[ + { + "id": "target-id", + "uri": "target-uri", + "abstract": "old abstract", + "overview": "old overview", + "content": "old content", + "context_type": "memory", + "category": "preferences", + "is_leaf": True, + "parent_uri": "parent-uri", + "session_id": "session-1", + "entities": ["Alice"], + "keywords": "auth", + "abstract_json": {"fact_points": ["old fact"]}, + } + ] + ) + storage.update = AsyncMock() fs = MagicMock() fs.read_file = AsyncMock(return_value="old content") - orch = SimpleNamespace( + fs.write_context = AsyncMock() + write_service = SimpleNamespace( _storage=storage, _fs=fs, + _embedder=None, _memory_signal_bus=_SignalBus(), _get_collection=MagicMock(return_value="context"), _get_record_by_uri=AsyncMock( @@ -58,13 +79,36 @@ def _build_service( "category": "preferences", } ), + _build_abstract_json=MagicMock( + side_effect=lambda **kwargs: { + **kwargs, + "memory_kind": "preference", + } + ), + _memory_object_payload=MagicMock( + return_value={ + "memory_kind": "preference", + "merge_signature": "sig", + "mergeable": True, + } + ), + _derive_layers=AsyncMock( + return_value={ + "entities": ["Bob"], + "keywords": "merged", + "fact_points": ["new fact"], + "anchor_handles": ["Bob"], + } + ), + _sync_anchor_projection_records=AsyncMock(), + feedback=AsyncMock(), ) memory_service = SimpleNamespace( update=AsyncMock(), feedback=AsyncMock(), ) - write_service = SimpleNamespace(_orch=orch, _service=memory_service) - return MemoryWriteDedupService(write_service), orch, memory_service + write_service._service = memory_service + return MemoryWriteDedupService(write_service), write_service, memory_service async def test_check_duplicate_builds_scope_and_project_filter(self) -> None: """Duplicate search applies tenant, scope, kind, signature, and project.""" @@ -222,12 +266,15 @@ async def test_try_merge_duplicate_merges_and_publishes_signal(self) -> None: self.assertEqual(ctx.uri, "target-uri") self.assertEqual(ctx.meta["dedup_action"], "merged") self.assertEqual(ctx.meta["dedup_score"], 0.93) - memory_service.update.assert_awaited_once_with( - "target-uri", - abstract="new abstract", - content="old content\n---\nnew content", - ) - memory_service.feedback.assert_awaited_once_with("target-uri", 0.5) + memory_service.update.assert_not_awaited() + orch.feedback.assert_awaited_once_with("target-uri", 0.5) + orch._storage.update.assert_awaited_once() + update_data = orch._storage.update.await_args.args[2] + self.assertEqual(update_data["abstract"], "new abstract") + self.assertEqual(update_data["abstract_json"]["fact_points"], ["new fact"]) + orch._derive_layers.assert_awaited_once() + orch._sync_anchor_projection_records.assert_awaited_once() + orch._fs.write_context.assert_awaited_once() self.assertEqual(len(orch._memory_signal_bus.signals), 1) signal = orch._memory_signal_bus.signals[0] @@ -246,7 +293,7 @@ async def test_merge_into_without_new_content_preserves_existing_content( self, ) -> None: """Empty incoming content keeps the existing filesystem content.""" - service, _orch, memory_service = self._build_service() + service, orch, memory_service = self._build_service() await service.merge_into( "target-uri", @@ -254,9 +301,9 @@ async def test_merge_into_without_new_content_preserves_existing_content( new_content="", ) - memory_service.update.assert_awaited_once_with( - "target-uri", - abstract="new abstract", - content="old content", + memory_service.update.assert_not_awaited() + orch.feedback.assert_awaited_once_with("target-uri", 0.5) + orch._fs.write_context.assert_awaited_once() + self.assertEqual( + orch._fs.write_context.await_args.kwargs["content"], "old content" ) - memory_service.feedback.assert_awaited_once_with("target-uri", 0.5) diff --git a/tests/test_memory_write_derive_service.py b/tests/test_memory_write_derive_service.py index cdcb134..13398ef 100644 --- a/tests/test_memory_write_derive_service.py +++ b/tests/test_memory_write_derive_service.py @@ -14,7 +14,7 @@ class TestMemoryWriteDeriveService(unittest.IsolatedAsyncioTestCase): """Verify derive/fallback behavior extracted from MemoryWriteService.add.""" def _build_service(self) -> tuple[MemoryWriteDeriveService, SimpleNamespace]: - orch = SimpleNamespace( + write_service = SimpleNamespace( _derive_layers=AsyncMock( return_value={ "abstract": "derived abstract", @@ -26,8 +26,7 @@ def _build_service(self) -> tuple[MemoryWriteDeriveService, SimpleNamespace]: _fallback_overview_from_content=MagicMock(return_value="fallback overview"), _derive_abstract_from_overview=MagicMock(return_value="fallback abstract"), ) - write_service = SimpleNamespace(_orch=orch) - return MemoryWriteDeriveService(write_service), orch + return MemoryWriteDeriveService(write_service), write_service async def test_content_leaf_derives_missing_abstract_and_overview(self) -> None: """Non-deferred content leaf writes call _derive_layers and fill blanks.""" diff --git a/tests/test_memory_write_embed_service.py b/tests/test_memory_write_embed_service.py index 69264c8..1ff4a7a 100644 --- a/tests/test_memory_write_embed_service.py +++ b/tests/test_memory_write_embed_service.py @@ -31,8 +31,7 @@ def _build_service( self, embedder: _SpyEmbedder | None, ) -> MemoryWriteEmbedService: - orch = SimpleNamespace(_embedder=embedder) - write_service = SimpleNamespace(_orch=orch) + write_service = SimpleNamespace(_embedder=embedder) return MemoryWriteEmbedService(write_service) async def test_no_embedder_returns_empty_result(self) -> None: diff --git a/tests/test_orchestrator_services.py b/tests/test_orchestrator_services.py index 74eda5e..cf09772 100644 --- a/tests/test_orchestrator_services.py +++ b/tests/test_orchestrator_services.py @@ -5,8 +5,8 @@ from opencortex.cortex_memory import CortexMemory from opencortex.orchestrator import MemoryOrchestrator -from opencortex.services.memory_service import MemoryService from opencortex.services.cortex_memory_services import CortexMemoryServices +from opencortex.services.memory_service import MemoryService from opencortex.services.orchestrator_services import MemoryOrchestratorServices