Removed legacy storage management cruft. Tidied documents.

This commit is contained in:
Cyber MacGeddon 2026-01-05 11:20:08 +00:00
parent 25563bae3c
commit 96ae0e1c60
4 changed files with 87 additions and 236 deletions

View file

@ -233,9 +233,13 @@ When a user initiates collection deletion through the librarian service:
#### Collection Management Interface #### Collection Management Interface
All store writers implement a standardized collection management interface with a common schema: **⚠️ LEGACY APPROACH - REPLACED BY CONFIG-BASED PATTERN**
**Message Schema (`StorageManagementRequest`):** The queue-based architecture described below has been replaced with a config-based approach using `CollectionConfigHandler`. All storage backends now receive collection updates via config push messages instead of dedicated management queues.
~~All store writers implement a standardized collection management interface with a common schema:~~
~~**Message Schema (`StorageManagementRequest`):**~~
```json ```json
{ {
"operation": "create-collection" | "delete-collection", "operation": "create-collection" | "delete-collection",
@ -244,24 +248,26 @@ All store writers implement a standardized collection management interface with
} }
``` ```
**Queue Architecture:** ~~**Queue Architecture:**~~
- **Vector Store Management Queue** (`vector-storage-management`): Vector/embedding stores - ~~**Vector Store Management Queue** (`vector-storage-management`): Vector/embedding stores~~
- **Object Store Management Queue** (`object-storage-management`): Object/document stores - ~~**Object Store Management Queue** (`object-storage-management`): Object/document stores~~
- **Triple Store Management Queue** (`triples-storage-management`): Graph/RDF stores - ~~**Triple Store Management Queue** (`triples-storage-management`): Graph/RDF stores~~
- **Storage Response Queue** (`storage-management-response`): All responses sent here - ~~**Storage Response Queue** (`storage-management-response`): All responses sent here~~
Each store writer implements: **Current Implementation:**
- **Collection Management Handler**: Processes `StorageManagementRequest` messages
- **Create Collection Operation**: Establishes collection in storage backend
- **Delete Collection Operation**: Removes all data associated with collection
- **Collection State Tracking**: Maintains knowledge of which collections exist
- **Message Processing**: Consumes from dedicated management queue
- **Status Reporting**: Returns success/failure via `StorageManagementResponse`
- **Idempotent Operations**: Safe to call create/delete multiple times
**Supported Operations:** All storage backends now use `CollectionConfigHandler`:
- `create-collection`: Create collection in storage backend - **Config Push Integration**: Storage services register for config push notifications
- `delete-collection`: Remove all collection data from storage backend - **Automatic Synchronization**: Collections created/deleted based on config changes
- **Declarative Model**: Collections defined in config service, backends sync to match
- **No Request/Response**: Eliminates coordination overhead and response tracking
- **Collection State Tracking**: Maintained via `known_collections` cache
- **Idempotent Operations**: Safe to process same config multiple times
Each storage backend implements:
- `create_collection(user: str, collection: str, metadata: dict)` - Create collection structures
- `delete_collection(user: str, collection: str)` - Remove all collection data
- `collection_exists(user: str, collection: str) -> bool` - Validate before writes
#### Cassandra Triple Store Refactor #### Cassandra Triple Store Refactor
@ -365,62 +371,33 @@ Comprehensive testing will cover:
- `triples_collection` table for SPO queries and deletion tracking - `triples_collection` table for SPO queries and deletion tracking
- Collection deletion implemented with read-then-delete pattern - Collection deletion implemented with read-then-delete pattern
### 🔄 In Progress Components ### ✅ Migration to Config-Based Pattern - COMPLETED
1. **Collection Creation Broadcast** (`trustgraph-flow/trustgraph/librarian/collection_manager.py`) **All storage backends have been migrated from the queue-based pattern to the config-based `CollectionConfigHandler` pattern.**
- Update `update_collection()` to send "create-collection" to storage backends
- Wait for confirmations from all storage processors
- Handle creation failures appropriately
2. **Document Submission Handler** (`trustgraph-flow/trustgraph/librarian/service.py` or similar) Completed migrations:
- Check if collection exists when document submitted - ✅ `trustgraph-flow/trustgraph/storage/triples/cassandra/write.py`
- If not exists: Create collection with defaults before processing document - ✅ `trustgraph-flow/trustgraph/storage/triples/neo4j/write.py`
- Trigger same "create-collection" broadcast as `tg-set-collection` - ✅ `trustgraph-flow/trustgraph/storage/triples/memgraph/write.py`
- Ensure collection established before document flows to storage processors - ✅ `trustgraph-flow/trustgraph/storage/triples/falkordb/write.py`
- ✅ `trustgraph-flow/trustgraph/storage/doc_embeddings/qdrant/write.py`
- ✅ `trustgraph-flow/trustgraph/storage/graph_embeddings/qdrant/write.py`
- ✅ `trustgraph-flow/trustgraph/storage/doc_embeddings/milvus/write.py`
- ✅ `trustgraph-flow/trustgraph/storage/graph_embeddings/milvus/write.py`
- ✅ `trustgraph-flow/trustgraph/storage/doc_embeddings/pinecone/write.py`
- ✅ `trustgraph-flow/trustgraph/storage/graph_embeddings/pinecone/write.py`
- ✅ `trustgraph-flow/trustgraph/storage/objects/cassandra/write.py`
### ❌ Pending Components All backends now:
- Inherit from `CollectionConfigHandler`
- Register for config push notifications via `self.register_config_handler(self.on_collection_config)`
- Implement `create_collection(user, collection, metadata)` and `delete_collection(user, collection)`
- Use `collection_exists(user, collection)` to validate before writes
- Automatically sync with config service changes
1. **Collection State Tracking** - Need to implement in each storage backend: Legacy queue-based infrastructure removed:
- **Cassandra Triples**: Use `triples_collection` table with marker triples - ✅ Removed `StorageManagementRequest` and `StorageManagementResponse` schemas
- **Neo4j/Memgraph/FalkorDB**: Create `:CollectionMetadata` nodes - ✅ Removed storage management queue topic definitions
- **Qdrant/Milvus/Pinecone**: Use native collection APIs - ✅ Removed storage management consumer/producer from all backends
- **Cassandra Objects**: Add collection metadata tracking - ✅ Removed `on_storage_management` handlers from all backends
2. **Storage Management Handlers** - Need "create-collection" support in 12 files:
- `trustgraph-flow/trustgraph/storage/triples/cassandra/write.py`
- `trustgraph-flow/trustgraph/storage/triples/neo4j/write.py`
- `trustgraph-flow/trustgraph/storage/triples/memgraph/write.py`
- `trustgraph-flow/trustgraph/storage/triples/falkordb/write.py`
- `trustgraph-flow/trustgraph/storage/doc_embeddings/qdrant/write.py`
- `trustgraph-flow/trustgraph/storage/graph_embeddings/qdrant/write.py`
- `trustgraph-flow/trustgraph/storage/doc_embeddings/milvus/write.py`
- `trustgraph-flow/trustgraph/storage/graph_embeddings/milvus/write.py`
- `trustgraph-flow/trustgraph/storage/doc_embeddings/pinecone/write.py`
- `trustgraph-flow/trustgraph/storage/graph_embeddings/pinecone/write.py`
- `trustgraph-flow/trustgraph/storage/objects/cassandra/write.py`
- Plus any other storage implementations
3. **Write Operation Validation** - Add collection existence checks to all `store_*` methods
4. **Query Operation Handling** - Update queries to return empty for non-existent collections
### Next Implementation Steps
**Phase 1: Core Infrastructure (2-3 days)**
1. Add collection state tracking methods to all storage backends
2. Implement `collection_exists()` and `create_collection()` methods
**Phase 2: Storage Handlers (1 week)**
3. Add "create-collection" handlers to all storage processors
4. Add write validation to reject non-existent collections
5. Update query handling for non-existent collections
**Phase 3: Collection Manager (2-3 days)**
6. Update collection_manager to broadcast creates
7. Implement response tracking and error handling
**Phase 4: Testing (3-5 days)**
8. End-to-end testing of explicit creation workflow
9. Test all storage backends
10. Validate error handling and edge cases

View file

@ -62,19 +62,20 @@ When tenant A starts flow `tenant-a-prod` and tenant B starts flow `tenant-b-pro
- **Impact:** Config, cores, and librarian services - **Impact:** Config, cores, and librarian services
- **Blocks:** Multiple tenants cannot use separate Cassandra keyspaces - **Blocks:** Multiple tenants cannot use separate Cassandra keyspaces
### Issue #4: Collection Management Architecture ### Issue #4: Collection Management Architecture ✅ COMPLETED
- **Current:** Collections stored in Cassandra librarian keyspace via separate collections table - **Previous:** Collections stored in Cassandra librarian keyspace via separate collections table
- **Current:** Librarian uses 4 hardcoded storage management topics to coordinate collection create/delete: - **Previous:** Librarian used 4 hardcoded storage management topics to coordinate collection create/delete:
- `vector_storage_management_topic` - `vector_storage_management_topic`
- `object_storage_management_topic` - `object_storage_management_topic`
- `triples_storage_management_topic` - `triples_storage_management_topic`
- `storage_management_response_topic` - `storage_management_response_topic`
- **Problems:** - **Problems (Resolved):**
- Hardcoded topics cannot be customized for multi-tenant deployments - Hardcoded topics could not be customized for multi-tenant deployments
- Complex async coordination between librarian and 4+ storage services - Complex async coordination between librarian and 4+ storage services
- Separate Cassandra table and management infrastructure - Separate Cassandra table and management infrastructure
- Non-persistent request/response queues for critical operations - Non-persistent request/response queues for critical operations
- **Solution:** Migrate collections to config service storage, use config push for distribution - **Solution Implemented:** Migrated collections to config service storage, use config push for distribution
- **Status:** All storage backends migrated to `CollectionConfigHandler` pattern
## Solution ## Solution
@ -448,7 +449,9 @@ async def delete_collection(self, user, collection):
- Breaking change acceptable - no data migration needed - Breaking change acceptable - no data migration needed
- Simplifies librarian service significantly - Simplifies librarian service significantly
#### Change 10: Storage Services - Config-Based Collection Management #### Change 10: Storage Services - Config-Based Collection Management ✅ COMPLETED
**Status:** All 11 storage backends have been migrated to use `CollectionConfigHandler`.
**Affected Services (11 total):** **Affected Services (11 total):**
- Document embeddings: milvus, pinecone, qdrant - Document embeddings: milvus, pinecone, qdrant
@ -708,9 +711,9 @@ All services inheriting from AsyncProcessor or FlowProcessor:
- tables/library.py (collections table removal) - tables/library.py (collections table removal)
- schema/services/collection.py (timestamp removal) - schema/services/collection.py (timestamp removal)
**Deferred Changes (Change 10):** **Completed Changes (Change 10):** ✅
- All storage services (11 total) - will subscribe to config push for collection updates - All storage services (11 total) - migrated to config push for collection updates via `CollectionConfigHandler`
- Storage management schema (potentially removable if unused elsewhere) - Storage management schema removed from `storage.py`
## Future Considerations ## Future Considerations
@ -749,10 +752,11 @@ Some services use **per-user keyspaces** dynamically, where each user gets their
- Update collection schema (remove timestamps) - Update collection schema (remove timestamps)
- **Outcome:** Eliminates hardcoded storage management topics, simplifies librarian - **Outcome:** Eliminates hardcoded storage management topics, simplifies librarian
### Phase 3: Storage Service Updates (Change 10) - Deferred ### Phase 3: Storage Service Updates (Change 10) ✅ COMPLETED
- Update all storage services to use config push for collections - Updated all storage services to use config push for collections via `CollectionConfigHandler`
- Remove storage management request/response infrastructure - Removed storage management request/response infrastructure
- **Outcome:** Complete config-based collection management - Removed legacy schema definitions
- **Outcome:** Complete config-based collection management achieved
## References ## References
- GitHub Issue: https://github.com/trustgraph-ai/trustgraph/issues/582 - GitHub Issue: https://github.com/trustgraph-ai/trustgraph/issues/582

View file

@ -1,45 +1,8 @@
from dataclasses import dataclass # This file previously contained legacy storage management queue definitions
# (StorageManagementRequest, StorageManagementResponse, and related topics).
from ..core.primitives import Error #
from ..core.topic import topic # These have been removed as collection management now uses a config-based
# approach via CollectionConfigHandler instead of request/response queues.
############################################################################ #
# This file is kept for potential future storage-related schema definitions.
# Storage management operations
@dataclass
class StorageManagementRequest:
"""Request for storage management operations sent to store processors"""
operation: str = "" # e.g., "delete-collection"
user: str = ""
collection: str = ""
@dataclass
class StorageManagementResponse:
"""Response from storage processors for management operations"""
error: Error | None = None # Only populated if there's an error, if null success
############################################################################
# Storage management topics
# Topics for sending collection management requests to different storage types
vector_storage_management_topic = topic(
'vector-storage-management', qos='q0', namespace='request'
)
object_storage_management_topic = topic(
'object-storage-management', qos='q0', namespace='request'
)
triples_storage_management_topic = topic(
'triples-storage-management', qos='q0', namespace='request'
)
# Topic for receiving responses from storage processors
storage_management_response_topic = topic(
'storage-management', qos='q0', namespace='response'
)
############################################################################

View file

@ -13,9 +13,8 @@ from cassandra import ConsistencyLevel
from .... schema import ExtractedObject from .... schema import ExtractedObject
from .... schema import RowSchema, Field from .... schema import RowSchema, Field
from .... schema import StorageManagementRequest, StorageManagementResponse
from .... schema import object_storage_management_topic, storage_management_response_topic
from .... base import FlowProcessor, ConsumerSpec, ProducerSpec from .... base import FlowProcessor, ConsumerSpec, ProducerSpec
from .... base import CollectionConfigHandler
from .... base.cassandra_config import add_cassandra_args, resolve_cassandra_config from .... base.cassandra_config import add_cassandra_args, resolve_cassandra_config
# Module logger # Module logger
@ -23,7 +22,7 @@ logger = logging.getLogger(__name__)
default_ident = "objects-write" default_ident = "objects-write"
class Processor(FlowProcessor): class Processor(CollectionConfigHandler, FlowProcessor):
def __init__(self, **params): def __init__(self, **params):
@ -64,39 +63,9 @@ class Processor(FlowProcessor):
) )
) )
# Set up storage management consumer and producer directly # Register config handlers
# (FlowProcessor doesn't support topic-based specs outside of flows)
from .... base import Consumer, Producer, ConsumerMetrics, ProducerMetrics
storage_request_metrics = ConsumerMetrics(
processor=self.id, flow=None, name="storage-request"
)
storage_response_metrics = ProducerMetrics(
processor=self.id, flow=None, name="storage-response"
)
# Create storage management consumer
self.storage_request_consumer = Consumer(
taskgroup=self.taskgroup,
backend=self.pubsub,
flow=None,
topic=object_storage_management_topic,
subscriber=f"{id}-storage",
schema=StorageManagementRequest,
handler=self.on_storage_management,
metrics=storage_request_metrics,
)
# Create storage management response producer
self.storage_response_producer = Producer(
backend=self.pubsub,
topic=storage_management_response_topic,
schema=StorageManagementResponse,
metrics=storage_response_metrics,
)
# Register config handler for schema updates
self.register_config_handler(self.on_schema_config) self.register_config_handler(self.on_schema_config)
self.register_config_handler(self.on_collection_config)
# Cache of known keyspaces/tables # Cache of known keyspaces/tables
self.known_keyspaces: Set[str] = set() self.known_keyspaces: Set[str] = set()
@ -347,28 +316,14 @@ class Processor(FlowProcessor):
obj = msg.value() obj = msg.value()
logger.info(f"Storing {len(obj.values)} objects for schema {obj.schema_name} from {obj.metadata.id}") logger.info(f"Storing {len(obj.values)} objects for schema {obj.schema_name} from {obj.metadata.id}")
# Validate collection/keyspace exists before accepting writes # Validate collection exists before accepting writes
safe_keyspace = self.sanitize_name(obj.metadata.user) if not self.collection_exists(obj.metadata.user, obj.metadata.collection):
if safe_keyspace not in self.known_keyspaces: error_msg = (
# Check if keyspace actually exists in Cassandra f"Collection {obj.metadata.collection} does not exist. "
self.connect_cassandra() f"Create it first via collection management API."
check_keyspace_cql = """ )
SELECT keyspace_name FROM system_schema.keyspaces logger.error(error_msg)
WHERE keyspace_name = %s raise ValueError(error_msg)
"""
result = self.session.execute(check_keyspace_cql, (safe_keyspace,))
# Check if result is None (mock case) or has no rows
if result is None or not result.one():
error_msg = (
f"Collection {obj.metadata.collection} does not exist. "
f"Create it first via collection management API."
)
logger.error(error_msg)
raise ValueError(error_msg)
# Cache it if it exists
self.known_keyspaces.add(safe_keyspace)
if safe_keyspace not in self.known_tables:
self.known_tables[safe_keyspace] = set()
# Get schema definition # Get schema definition
schema = self.schemas.get(obj.schema_name) schema = self.schemas.get(obj.schema_name)
@ -447,55 +402,7 @@ class Processor(FlowProcessor):
logger.error(f"Failed to insert object {obj_index}: {e}", exc_info=True) logger.error(f"Failed to insert object {obj_index}: {e}", exc_info=True)
raise raise
async def on_storage_management(self, msg, consumer, flow): async def create_collection(self, user: str, collection: str, metadata: dict):
"""Handle storage management requests for collection operations"""
request = msg.value()
logger.info(f"Received storage management request: {request.operation} for {request.user}/{request.collection}")
try:
if request.operation == "create-collection":
await self.create_collection(request.user, request.collection)
# Send success response
response = StorageManagementResponse(
error=None # No error means success
)
await self.storage_response_producer.send(response)
logger.info(f"Successfully created collection {request.user}/{request.collection}")
elif request.operation == "delete-collection":
await self.delete_collection(request.user, request.collection)
# Send success response
response = StorageManagementResponse(
error=None # No error means success
)
await self.storage_response_producer.send(response)
logger.info(f"Successfully deleted collection {request.user}/{request.collection}")
else:
logger.warning(f"Unknown storage management operation: {request.operation}")
# Send error response
from .... schema import Error
response = StorageManagementResponse(
error=Error(
type="unknown_operation",
message=f"Unknown operation: {request.operation}"
)
)
await self.storage_response_producer.send(response)
except Exception as e:
logger.error(f"Error handling storage management request: {e}", exc_info=True)
# Send error response
from .... schema import Error
response = StorageManagementResponse(
error=Error(
type="processing_error",
message=str(e)
)
)
await self.storage_response_producer.send(response)
async def create_collection(self, user: str, collection: str):
"""Create/verify collection exists in Cassandra object store""" """Create/verify collection exists in Cassandra object store"""
# Connect if not already connected # Connect if not already connected
self.connect_cassandra() self.connect_cassandra()