From 028c29df48f43bfe9f896d5be3b65f3ed175b9b1 Mon Sep 17 00:00:00 2001 From: Cameron Pfiffer Date: Tue, 27 May 2025 09:53:34 -0700 Subject: [PATCH] refactor: Enhance graph sync and service management functionality MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 🤖 Generated with [Claude Code](https://claude.ai/code) Co-Authored-By: Claude --- scripts/services.sh | 6 +- src/graph_sync.py | 255 +++++++++++++++++++++++++------------------- 2 files changed, 148 insertions(+), 113 deletions(-) diff --git a/scripts/services.sh b/scripts/services.sh index d926592..10a696c 100755 --- a/scripts/services.sh +++ b/scripts/services.sh @@ -119,13 +119,13 @@ sync_graph() { fi # Run the graph sync script - python scripts/graph_sync.py --setup-schema --sync-all + python3 scripts/graph_sync.py --setup-schema --sync-all echo "" echo "Graph sync complete! You can now:" echo " - Browse the graph: http://localhost:7474" - echo " - Query concepts: python scripts/graph_sync.py --concept-network 'your concept'" - echo " - Find clusters: python scripts/graph_sync.py --concept-clusters" + echo " - Query concepts: python3 scripts/graph_sync.py --concept-network 'your concept'" + echo " - Find clusters: python3 scripts/graph_sync.py --concept-clusters" } # Main command handling diff --git a/src/graph_sync.py b/src/graph_sync.py index b8375fa..19e787b 100644 --- a/src/graph_sync.py +++ b/src/graph_sync.py @@ -25,16 +25,16 @@ logger = logging.getLogger("graph_sync") class GraphSyncService: """ Syncs ATProto records to Neo4j graph database. - + This service creates a graph representation of Comind's knowledge network, enabling sophisticated queries and analysis of concepts, relationships, and content patterns. """ - + # Comind collections to sync COMIND_COLLECTIONS = [ "me.comind.concept", - "me.comind.thought", + "me.comind.thought", "me.comind.emotion", "me.comind.sphere.core", "me.comind.relationship.concept", @@ -42,19 +42,19 @@ class GraphSyncService: "me.comind.relationship.sphere", "me.comind.relationship.similarity" ] - + # External collections that may be referenced EXTERNAL_COLLECTIONS = [ "app.bsky.feed.post", "app.bsky.feed.like", "app.bsky.graph.follow" ] - - def __init__(self, neo4j_uri: str, neo4j_user: str, neo4j_password: str, + + def __init__(self, neo4j_uri: str, neo4j_user: str, neo4j_password: str, record_manager: RecordManager): """ Initialize the Graph Sync Service. - + Args: neo4j_uri: Neo4j connection URI (e.g., "bolt://localhost:7687") neo4j_user: Neo4j username @@ -63,7 +63,7 @@ class GraphSyncService: """ self.record_manager = record_manager self.driver = GraphDatabase.driver(neo4j_uri, auth=(neo4j_user, neo4j_password)) - + # Verify connection try: self.driver.verify_connectivity() @@ -71,27 +71,27 @@ class GraphSyncService: except Exception as e: logger.error(f"Failed to connect to Neo4j: {e}") raise - + def close(self): """Close the Neo4j driver connection.""" if self.driver: self.driver.close() logger.info("Neo4j connection closed") - + def setup_schema(self): """ Set up Neo4j schema with constraints and indexes for optimal performance. """ logger.info("Setting up Neo4j schema...") - + schema_queries = [ # Constraints for uniqueness "CREATE CONSTRAINT concept_uri IF NOT EXISTS FOR (c:Concept) REQUIRE c.uri IS UNIQUE", - "CREATE CONSTRAINT thought_uri IF NOT EXISTS FOR (t:Thought) REQUIRE t.uri IS UNIQUE", + "CREATE CONSTRAINT thought_uri IF NOT EXISTS FOR (t:Thought) REQUIRE t.uri IS UNIQUE", "CREATE CONSTRAINT emotion_uri IF NOT EXISTS FOR (e:Emotion) REQUIRE e.uri IS UNIQUE", "CREATE CONSTRAINT sphere_uri IF NOT EXISTS FOR (s:Sphere) REQUIRE s.uri IS UNIQUE", "CREATE CONSTRAINT post_uri IF NOT EXISTS FOR (p:Post) REQUIRE p.uri IS UNIQUE", - + # Indexes for common queries "CREATE INDEX concept_text IF NOT EXISTS FOR (c:Concept) ON (c.text)", "CREATE INDEX thought_type IF NOT EXISTS FOR (t:Thought) ON (t.thoughtType)", @@ -99,7 +99,7 @@ class GraphSyncService: "CREATE INDEX sphere_title IF NOT EXISTS FOR (s:Sphere) ON (s.title)", "CREATE INDEX created_at IF NOT EXISTS FOR (n) ON (n.createdAt)" ] - + with self.driver.session() as session: for query in schema_queries: try: @@ -107,24 +107,24 @@ class GraphSyncService: logger.debug(f"Executed schema query: {query[:50]}...") except Exception as e: logger.warning(f"Schema query failed (may already exist): {e}") - + logger.info("Neo4j schema setup complete") - + def sync_all_records(self, include_external: bool = False): """ Sync all Comind records to Neo4j. - + Args: include_external: Whether to include external collections (posts, likes, etc.) """ logger.info("Starting full sync of all records...") - + collections_to_sync = self.COMIND_COLLECTIONS.copy() if include_external: collections_to_sync.extend(self.EXTERNAL_COLLECTIONS) - + total_synced = 0 - + for collection in collections_to_sync: try: synced_count = self.sync_collection(collection) @@ -132,57 +132,70 @@ class GraphSyncService: logger.info(f"Synced {synced_count} records from {collection}") except Exception as e: logger.error(f"Failed to sync collection {collection}: {e}") - + logger.info(f"Full sync complete. Total records synced: {total_synced}") return total_synced - + def sync_collection(self, collection: str) -> int: """ Sync all records from a specific collection. - + Args: collection: The collection NSID to sync - + Returns: Number of records synced """ logger.info(f"Syncing collection: {collection}") - + try: records = self.record_manager.list_records(collection) - + if not records: logger.info(f"No records found in collection: {collection}") return 0 - + synced_count = 0 - + for record in records: try: self.sync_record(record, collection) synced_count += 1 except Exception as e: + import traceback logger.error(f"Failed to sync record {record.uri}: {e}") - + logger.error(f"Full traceback: {traceback.format_exc()}") + # Also debug the record structure + logger.error(f"Record type: {type(record)}") + logger.error(f"Record attributes: {dir(record)}") + return synced_count - + except Exception as e: logger.error(f"Failed to list records in collection {collection}: {e}") raise - + def sync_record(self, record: Any, collection: str): """ Sync a single ATProto record to Neo4j. - + Args: record: ATProto record object collection: The collection this record belongs to """ + # Convert pydantic record to dict + record_dict = record.model_dump() + # Extract basic properties - uri = record.uri - cid = record.cid - value = record.value - + uri = record_dict.get('uri') + cid = record_dict.get('cid') + value = record_dict.get('value') + + # Check for None values + if uri is None or cid is None or value is None: + logger.error(f"Record has None attributes: uri={uri}, cid={cid}, value={value}") + return + # Determine record type and create appropriate node if collection == "me.comind.concept": self._create_concept_node(uri, cid, value) @@ -202,27 +215,28 @@ class GraphSyncService: self._create_post_node(uri, cid, value) else: logger.warning(f"Unknown collection type: {collection}") - + def _create_concept_node(self, uri: str, cid: str, value: Dict): """Create a Concept node in Neo4j.""" + query = """ MERGE (c:Concept {uri: $uri}) SET c.cid = $cid, c.text = $text, c.updatedAt = datetime() """ - + with self.driver.session() as session: - session.run(query, + session.run(query, uri=uri, - cid=cid, + cid=cid, text=value.get('concept', '') ) - + def _create_thought_node(self, uri: str, cid: str, value: Dict): """Create a Thought node in Neo4j.""" generated = value.get('generated', {}) - + query = """ MERGE (t:Thought {uri: $uri}) SET t.cid = $cid, @@ -233,7 +247,7 @@ class GraphSyncService: t.createdAt = $createdAt, t.updatedAt = datetime() """ - + with self.driver.session() as session: session.run(query, uri=uri, @@ -244,11 +258,11 @@ class GraphSyncService: confidence=generated.get('confidence'), createdAt=value.get('createdAt', '') ) - + def _create_emotion_node(self, uri: str, cid: str, value: Dict): """Create an Emotion node in Neo4j.""" generated = value.get('generated', {}) - + query = """ MERGE (e:Emotion {uri: $uri}) SET e.cid = $cid, @@ -257,7 +271,7 @@ class GraphSyncService: e.createdAt = $createdAt, e.updatedAt = datetime() """ - + with self.driver.session() as session: session.run(query, uri=uri, @@ -266,7 +280,7 @@ class GraphSyncService: emotionType=generated.get('emotionType', ''), createdAt=value.get('createdAt', '') ) - + def _create_sphere_node(self, uri: str, cid: str, value: Dict): """Create a Sphere node in Neo4j.""" query = """ @@ -278,7 +292,7 @@ class GraphSyncService: s.createdAt = $createdAt, s.updatedAt = datetime() """ - + with self.driver.session() as session: session.run(query, uri=uri, @@ -288,7 +302,7 @@ class GraphSyncService: description=value.get('description', ''), createdAt=value.get('createdAt', '') ) - + def _create_post_node(self, uri: str, cid: str, value: Dict): """Create a Post node in Neo4j.""" query = """ @@ -298,7 +312,7 @@ class GraphSyncService: p.createdAt = $createdAt, p.updatedAt = datetime() """ - + with self.driver.session() as session: session.run(query, uri=uri, @@ -306,23 +320,32 @@ class GraphSyncService: text=value.get('text', ''), createdAt=value.get('createdAt', '') ) - + def _create_concept_relationship(self, uri: str, cid: str, value: Dict): """Create a relationship between source and concept.""" source_uri = value.get('source', '') target_uri = value.get('target', '') relationship_type = value.get('relationship', 'RELATES_TO') - + + print("Target URI:", target_uri ) + print("Source URI:", source_uri ) + print("Relationship Type:", relationship_type) + query = """ - MATCH (source {uri: $source_uri}) - MATCH (target:Concept {uri: $target_uri}) + MERGE (source {uri: $source_uri}) + ON CREATE SET source.createdAt = datetime() + MERGE (target:Concept {uri: $target_uri}) + ON CREATE SET target.createdAt = datetime() MERGE (source)-[r:CONCEPT_RELATION {uri: $uri}]->(target) - SET r.cid = $cid, - r.relationship = $relationship, - r.createdAt = $createdAt, - r.updatedAt = datetime() + ON CREATE SET r.cid = $cid, + r.relationship = $relationship, + r.createdAt = $createdAt, + r.updatedAt = datetime() + ON MATCH SET r.cid = $cid, + r.relationship = $relationship, + r.updatedAt = datetime() """ - + with self.driver.session() as session: session.run(query, uri=uri, @@ -332,26 +355,34 @@ class GraphSyncService: relationship=relationship_type, createdAt=value.get('createdAt', '') ) - + def _create_link_relationship(self, uri: str, cid: str, value: Dict): """Create a general link relationship between nodes.""" - source_uri = value.get('source', '') + source = value.get('source', {}) + source_uri = source.get('uri', '') if isinstance(source, dict) else source target_uri = value.get('target', '') generated = value.get('generated', {}) relationship_type = generated.get('relationship', 'LINKS_TO') - + query = """ - MATCH (source {uri: $source_uri}) - MATCH (target {uri: $target_uri}) + MERGE (source {uri: $source_uri}) + ON CREATE SET source.createdAt = datetime() + MERGE (target {uri: $target_uri}) + ON CREATE SET target.createdAt = datetime() MERGE (source)-[r:LINK {uri: $uri}]->(target) - SET r.cid = $cid, - r.relationship = $relationship, - r.strength = $strength, - r.note = $note, - r.createdAt = $createdAt, - r.updatedAt = datetime() + ON CREATE SET r.cid = $cid, + r.relationship = $relationship, + r.strength = $strength, + r.note = $note, + r.createdAt = $createdAt, + r.updatedAt = datetime() + ON MATCH SET r.cid = $cid, + r.relationship = $relationship, + r.strength = $strength, + r.note = $note, + r.updatedAt = datetime() """ - + with self.driver.session() as session: session.run(query, uri=uri, @@ -363,21 +394,25 @@ class GraphSyncService: note=generated.get('note', ''), createdAt=value.get('createdAt', '') ) - + def _create_sphere_relationship(self, uri: str, cid: str, value: Dict): """Create a relationship between content and sphere.""" target_uri = value.get('target', '') sphere_uri = value.get('sphere_uri', '') - + query = """ - MATCH (target {uri: $target_uri}) - MATCH (sphere:Sphere {uri: $sphere_uri}) + MERGE (target {uri: $target_uri}) + ON CREATE SET target.createdAt = datetime() + MERGE (sphere:Sphere {uri: $sphere_uri}) + ON CREATE SET sphere.createdAt = datetime() MERGE (target)-[r:IN_SPHERE {uri: $uri}]->(sphere) - SET r.cid = $cid, - r.createdAt = $createdAt, - r.updatedAt = datetime() + ON CREATE SET r.cid = $cid, + r.createdAt = $createdAt, + r.updatedAt = datetime() + ON MATCH SET r.cid = $cid, + r.updatedAt = datetime() """ - + with self.driver.session() as session: session.run(query, uri=uri, @@ -386,15 +421,15 @@ class GraphSyncService: sphere_uri=sphere_uri, createdAt=value.get('createdAt', '') ) - + def get_concept_network(self, concept_text: str, depth: int = 2) -> Dict: """ Get the network of concepts connected to a given concept. - + Args: concept_text: The concept to explore depth: How many relationship hops to include - + Returns: Dictionary with nodes and relationships """ @@ -402,13 +437,13 @@ class GraphSyncService: MATCH path = (c:Concept {{text: $concept_text}})-[*1..{depth}]-(connected) RETURN path """ - + with self.driver.session() as session: result = session.run(query, concept_text=concept_text) - + nodes = set() relationships = [] - + for record in result: path = record['path'] for node in path.nodes: @@ -420,19 +455,19 @@ class GraphSyncService: 'type': rel.type, 'properties': dict(rel) }) - + return { 'nodes': [{'id': node_id, 'properties': props} for node_id, props in nodes], 'relationships': relationships } - + def get_sphere_concepts(self, sphere_title: str) -> List[Dict]: """ Get all concepts associated with a sphere. - + Args: sphere_title: The sphere title to query - + Returns: List of concept dictionaries """ @@ -442,19 +477,19 @@ class GraphSyncService: RETURN DISTINCT c.text as concept, count(*) as frequency ORDER BY frequency DESC """ - + with self.driver.session() as session: result = session.run(query, sphere_title=sphere_title) - return [{'concept': record['concept'], 'frequency': record['frequency']} + return [{'concept': record['concept'], 'frequency': record['frequency']} for record in result] - + def find_concept_clusters(self, min_connections: int = 3) -> List[Dict]: """ Find clusters of highly connected concepts. - + Args: min_connections: Minimum number of connections for a concept to be included - + Returns: List of concept clusters """ @@ -464,30 +499,30 @@ class GraphSyncService: WHERE connections >= $min_connections MATCH (c)<-[:CONCEPT_RELATION]-(source)-[:CONCEPT_RELATION]->(related:Concept) WHERE related <> c - RETURN c.text as concept, + RETURN c.text as concept, connections, collect(DISTINCT related.text) as related_concepts ORDER BY connections DESC """ - + with self.driver.session() as session: result = session.run(query, min_connections=min_connections) return [dict(record) for record in result] def create_graph_sync_service(neo4j_uri: str = "bolt://localhost:7687", - neo4j_user: str = "neo4j", + neo4j_user: str = "neo4j", neo4j_password: str = "comind123", record_manager: Optional[RecordManager] = None) -> GraphSyncService: """ Factory function to create a GraphSyncService instance. - + Args: neo4j_uri: Neo4j connection URI - neo4j_user: Neo4j username + neo4j_user: Neo4j username neo4j_password: Neo4j password record_manager: RecordManager instance (will create one if None) - + Returns: Configured GraphSyncService instance """ @@ -495,16 +530,16 @@ def create_graph_sync_service(neo4j_uri: str = "bolt://localhost:7687", from session_reuse import default_login client = default_login() record_manager = RecordManager(client) - + return GraphSyncService(neo4j_uri, neo4j_user, neo4j_password, record_manager) if __name__ == "__main__": # Example usage import argparse - + parser = argparse.ArgumentParser(description="Sync ATProto records to Neo4j") - parser.add_argument("--neo4j-uri", default="bolt://localhost:7687", + parser.add_argument("--neo4j-uri", default="bolt://localhost:7687", help="Neo4j URI") parser.add_argument("--neo4j-user", default="neo4j", help="Neo4j username") @@ -516,25 +551,25 @@ if __name__ == "__main__": help="Sync all records") parser.add_argument("--collection", type=str, help="Sync specific collection") - + args = parser.parse_args() - + try: # Create sync service sync_service = create_graph_sync_service( args.neo4j_uri, args.neo4j_user, args.neo4j_password ) - + if args.setup_schema: sync_service.setup_schema() - + if args.sync_all: sync_service.sync_all_records() elif args.collection: sync_service.sync_collection(args.collection) - + sync_service.close() - + except Exception as e: logger.error(f"Graph sync failed: {e}") - raise \ No newline at end of file + raise -- 2.51.2