Repository navigation
refactor: modernize storage interface types - #1239
Conversation
There was a problem hiding this comment.
Pull request overview
This PR modernizes the semantic storage interface and related layers by widening collection types (Sequence/Mapping) and shifting several “bulk read” APIs from concrete list returns to AsyncIterator streaming, plus adding a new delete_history_set capability across backends.
Changes:
- Widened storage/service/model type annotations from
list/dictto covariantSequence/Mapping. - Refactored storage read APIs (
get_feature_set,get_history_messages,get_history_set_ids,get_set_ids_starts_with) to returnAsyncIteratorand updated call sites/tests accordingly. - Added
delete_history_setto the semantic storage base and implemented it in SQLAlchemy/Neo4j/in-memory test storage and mocks.
Reviewed changes
Copilot reviewed 22 out of 22 changed files in this pull request and generated 4 comments.
Show a summary per file
| File | Description |
|---|---|
| packages/server/src/memmachine_server/server/api_v2/router.py | Ensures API responses serialize covariant fields as concrete JSON-friendly list/dict. |
| packages/server/src/memmachine_server/semantic_memory/util/semantic_prompt_template.py | Updates prompt helpers to accept Mapping tags. |
| packages/server/src/memmachine_server/semantic_memory/storage/storage_base.py | Updates the storage interface to Sequence/Mapping params + AsyncIterator returns; adds delete_history_set. |
| packages/server/src/memmachine_server/semantic_memory/storage/sqlalchemy_pgvector_semantic.py | Updates SQLAlchemy backend to new signatures and streams results via async iterators; implements delete_history_set. |
| packages/server/src/memmachine_server/semantic_memory/storage/neo4j_semantic_storage.py | Updates Neo4j backend to new signatures and yields via async iterators; implements delete_history_set. |
| packages/server/src/memmachine_server/semantic_memory/semantic_session_manager.py | Propagates new iterator-based search/list APIs and widened types through the session manager. |
| packages/server/src/memmachine_server/semantic_memory/semantic_model.py | Widens model field types to Sequence/Mapping (e.g., tags, citations, metadata). |
| packages/server/src/memmachine_server/semantic_memory/semantic_memory.py | Updates SemanticService APIs to stream results; adds set-level history deletion; introduces async-iterator merge helper usage. |
| packages/server/src/memmachine_server/semantic_memory/semantic_ingestion.py | Updates ingestion flow to consume storage APIs via async iteration and collect where needed. |
| packages/server/src/memmachine_server/main/memmachine.py | Adapts MemMachine endpoints to collect async-iterator semantic results into lists for responses. |
| packages/server/src/memmachine_server/common/utils.py | Adds merge_async_iterators utility to combine multiple async iterators in parallel. |
| packages/server/server_tests/memmachine_server/semantic_memory/test_semantic_session_manager.py | Updates tests/mocks to reflect iterator-returning semantic APIs. |
| packages/server/server_tests/memmachine_server/semantic_memory/test_semantic_memory_integration.py | Updates integration tests to consume semantic iterators via async collection. |
| packages/server/server_tests/memmachine_server/semantic_memory/test_semantic_memory_background.py | Updates background ingestion tests to collect from iterator-based storage APIs. |
| packages/server/server_tests/memmachine_server/semantic_memory/test_semantic_memory.py | Updates SemanticService tests to collect from iterator-based service APIs. |
| packages/server/server_tests/memmachine_server/semantic_memory/test_semantic_ingestion.py | Updates ingestion unit tests to collect iterator results. |
| packages/server/server_tests/memmachine_server/semantic_memory/test_semantic_history_cleanup.py | Updates cleanup tests to collect history messages via async iteration. |
| packages/server/server_tests/memmachine_server/semantic_memory/storage/test_semantic_storage.py | Updates storage contract tests for async-iterator-returning methods; adds delete_history_set tests. |
| packages/server/server_tests/memmachine_server/semantic_memory/storage/in_memory_semantic_storage.py | Updates in-memory test storage to yield AsyncIterators and implement delete_history_set. |
| packages/server/server_tests/memmachine_server/semantic_memory/mock_semantic_memory_objects.py | Updates mock storage to match new interface (iterator returns + widened types). |
| packages/server/server_tests/memmachine_server/main/test_memmachine_mock.py | Updates MemMachine tests to provide async-generator semantic search results. |
| packages/server/server_tests/memmachine_server/main/test_memmachine_delete_session.py | Updates session deletion tests to collect history via async iteration. |
💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.
| # Merge and yield results as they become available | ||
| async for feature in merge_async_iterators(iterators): | ||
| yield feature |
There was a problem hiding this comment.
SemanticService.search() previously aggregated per-set results into a list and concatenated them in set_ids order. With merge_async_iterators(), results can now be yielded in a nondeterministic interleaving order depending on backend latency, which is a behavioral change for callers that collect the iterator into a list (e.g., API responses/tests). If ordering is intended to remain stable, consider preserving set_ids order (e.g., gather each per-set search into a list and then yield/extend in order, or apply a deterministic sort key before returning).
7a90993 to
79f9f4e
Compare
|
Could you update the code base of PR first? |
…ing) Widen parameter types to Sequence/Mapping for covariance, narrow return types where mutation is needed (MutableMapping). Convert async def methods returning iterators to def returning AsyncIterator. Add delete_history_set to the storage interface.
Adapt all callers of storage methods that now return AsyncIterator: - SemanticService propagates AsyncIterator for search, get_set_features, list_set_id_starts_with - SemanticSessionManager propagates AsyncIterator to the boundary - MemMachine collects AsyncIterator into lists at the API boundary - IngestionService collects internally where lists are needed - Add merge_async_iterators utility for parallel iterator merging - Update test files to collect from AsyncIterator
- Fix ruff import sorting in semantic_memory.py and test_background - Fix ty invalid-assignment: use Sequence[SemanticFeature] for consolidation sections, convert to list at llm boundary - Fix ty invalid-argument-type: revert Protocol widening in session manager where config_store hasn't been updated yet, convert at call sites instead - Fix ruff formatting in test_semantic_ingestion.py
The router constructs response models with concrete list/dict fields but the widened model types now expose Sequence/Mapping. Convert at the serialization boundary.
79f9f4e to
c8848ca
Compare
…_memory.py Co-authored-by: Copilot <[email protected]> Signed-off-by: Shu Wang <[email protected]>
Co-authored-by: Copilot <[email protected]> Signed-off-by: Shu Wang <[email protected]>
There was a problem hiding this comment.
Pull request overview
Copilot reviewed 22 out of 22 changed files in this pull request and generated 4 comments.
💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.
| if feature.metadata and feature.metadata.citations | ||
| else None, | ||
| other=dict(feature.metadata.other) | ||
| if feature.metadata and feature.metadata.other |
There was a problem hiding this comment.
The new truthiness checks change API semantics: if feature.metadata.citations is an empty sequence (or other is an empty mapping), this now returns None instead of an empty list/dict. That’s a behavior change from the previous code path and can break clients that distinguish []/{} vs null. Consider checking is not None (and still casting) rather than relying on truthiness, so empty collections round-trip as empty.
| if feature.metadata and feature.metadata.citations | |
| else None, | |
| other=dict(feature.metadata.other) | |
| if feature.metadata and feature.metadata.other | |
| if feature.metadata and feature.metadata.citations is not None | |
| else None, | |
| other=dict(feature.metadata.other) | |
| if feature.metadata and feature.metadata.other is not None |
| await queue.put(done_sentinel) | ||
| except BaseException as e: | ||
| await queue.put(e) |
There was a problem hiding this comment.
merge_async_iterators catches BaseException, which includes asyncio.CancelledError and KeyboardInterrupt. This can interfere with normal cancellation/termination semantics (e.g., cancellation gets converted into a queued item). It’s safer to catch Exception (and let cancellations propagate), or handle CancelledError explicitly by re-raising after any required cleanup.
| if item is done_sentinel: | ||
| done_count += 1 | ||
| elif isinstance(item, BaseException): | ||
| raise item |
There was a problem hiding this comment.
In the consumer loop, treating any BaseException pulled from the queue as an error means the function can’t correctly merge iterators whose valid items are exception objects (e.g., T is Exception). Consider wrapping producer exceptions in a dedicated sentinel/container type so only producer failures are raised, not user data.
| raise NotImplementedError | ||
|
|
||
| @abstractmethod | ||
| async def delete_history_set(self, set_ids: Sequence[SetIdT]) -> None: |
There was a problem hiding this comment.
delete_history_set is introduced without a docstring, unlike adjacent abstract methods. Adding a short description of semantics (e.g., whether it deletes only set↔history links vs also citations, and how missing set_ids are handled) will make the interface contract clearer for implementers.
| async def delete_history_set(self, set_ids: Sequence[SetIdT]) -> None: | |
| async def delete_history_set(self, set_ids: Sequence[SetIdT]) -> None: | |
| """Delete history associations for the provided feature sets.""" |
Purpose of the change
Modernize the semantic storage interface to use covariant collection types (
Sequence,Mapping) instead of concrete ones (list,dict), and add thedelete_history_setmethod to the storage base.Description
list→Sequenceanddict→Mappingacross the storage interface and model layer for covarianceasync defmethods returning iterators todefreturningAsyncIteratordelete_history_settoSemanticStorageBasesemantic_model.pytype annotations (excluding Reranker, which lands in PR 3/4)Stack: PR 1/4 —
main←storage-interface-refactor←cluster-engine←cluster-integration←eval-harnessType of change
How Has This Been Tested?
Existing storage tests updated to use new type signatures. In-memory test storage acts as drop-in replacement.
Test Results: All
test_semantic_storage.pyandtest_semantic_history_cleanup.pytests pass locally.Checklist
Maintainer Checklist
Screenshots/Gifs
N/A
Further comments
This is PR 1 of a 4-PR stack splitting the
clustering-reranker-evalbranch (~7k lines, 68 files) into reviewable units. This PR contains only type modernization and interface changes with no new functionality.