diff --git a/activated_dids.txt b/activated_dids.txt index 489659f..6c20009 100644 --- a/activated_dids.txt +++ b/activated_dids.txt @@ -3,4 +3,19 @@ # newlines. cameron.pfiffer.org -neuromute.ai \ No newline at end of file +neuromute.ai + +# deanshamess.ca +did:plc:ycrg7g2vkickmvxurd7bnidj + +# goat.navy +goat.navy + +# dulanyw.bsky.social +dulanyw.bsky.social + +# baileytownsend.dev +did:plc:rnpkyqnmsw4ipey6eotbdnnf + +# the dame +dame.is diff --git a/docker-compose.yml b/docker-compose.yml index c47fabb..a47ad3b 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -26,7 +26,7 @@ services: ports: - "8002:8000" command: > - --model noxneural/Hermes-3-Llama-3.2-3B-awq-4-bit + --model NousResearch/Hermes-3-Llama-3.1-8B --max_model_len 20000 --guided-decoding-backend outlines diff --git a/src/graph_sync.py b/src/graph_sync.py index a5a13c5..65d34f1 100644 --- a/src/graph_sync.py +++ b/src/graph_sync.py @@ -157,12 +157,14 @@ class GraphSyncService: logger.info(f"Syncing collection: {collection}") try: - records = self.record_manager.list_records(collection) + # Use list_all_records to handle pagination + records = self.record_manager.list_all_records(collection) if not records: logger.info(f"No records found in collection: {collection}") return 0 + logger.info(f"Found {len(records)} records in collection: {collection}") synced_count = 0 for record in records: @@ -454,31 +456,84 @@ class GraphSyncService: target_uri = value.get("target", "") relationship_type = value.get("relationship", "RELATES_TO") - query = """ - 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) - 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() - """ + # First, check if the concept exists and has text + # If not, try to fetch it from the repository + concept_text = None + if self.record_manager and target_uri: + try: + # Parse the URI to get collection and rkey + parts = target_uri.split('/') + if len(parts) >= 5 and target_uri.startswith('at://'): + repo = parts[2] + collection = parts[3] + rkey = '/'.join(parts[4:]) + + # Fetch the concept record + concept_record = self.record_manager.client.com.atproto.repo.get_record({ + 'collection': collection, + 'repo': repo, + 'rkey': rkey + }) + + if concept_record and concept_record.value: + concept_text = concept_record.value.get('concept', None) + logger.debug(f"Fetched concept text for {target_uri}: {concept_text}") + except Exception as e: + logger.warning(f"Failed to fetch concept record {target_uri}: {e}") + + # Create the relationship, ensuring the concept has its text if we found it + if concept_text: + query = """ + MERGE (source {uri: $source_uri}) + ON CREATE SET source.createdAt = datetime() + MERGE (target:Concept {uri: $target_uri}) + ON CREATE SET target.createdAt = datetime() + SET target.text = $concept_text + MERGE (source)-[r:CONCEPT_RELATION {uri: $uri}]->(target) + 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() + """ + params = { + "uri": uri, + "cid": cid, + "source_uri": source_uri, + "target_uri": target_uri, + "concept_text": concept_text, + "relationship": relationship_type, + "createdAt": value.get("createdAt", ""), + } + else: + # Fall back to original query without setting text + query = """ + 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) + 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() + """ + params = { + "uri": uri, + "cid": cid, + "source_uri": source_uri, + "target_uri": target_uri, + "relationship": relationship_type, + "createdAt": value.get("createdAt", ""), + } with self.driver.session() as session: - session.run( - query, - uri=uri, - cid=cid, - source_uri=source_uri, - target_uri=target_uri, - relationship=relationship_type, - createdAt=value.get("createdAt", ""), - ) + session.run(query, **params) def _create_link_relationship(self, uri: str, cid: str, value: Dict): """Create a general link relationship between nodes.""" diff --git a/src/jetstream_consumer.py b/src/jetstream_consumer.py index da9674b..02ead50 100644 --- a/src/jetstream_consumer.py +++ b/src/jetstream_consumer.py @@ -178,18 +178,22 @@ def is_did(text: str) -> bool: def resolve_handle_to_did(client, handle: str, user_info_cache: UserInfoCache) -> Optional[str]: """Resolve a handle to a DID using the ATProto client""" - if user_info_cache.contains(handle): - user_info = user_info_cache.get_user_info(handle) - did = user_info.did - else: - user_info = client.get_profile(handle) - handle = user_info.handle - display_name = user_info.display_name - did = user_info.did - description = user_info.description - user_info_cache.add_user_info(did, UserInfo(did=did, handle=handle, display_name=display_name, description=description)) - - return did + try: + if user_info_cache.contains(handle): + user_info = user_info_cache.get_user_info(handle) + did = user_info.did + else: + user_info = client.get_profile(handle) + handle = user_info.handle + display_name = user_info.display_name + did = user_info.did + description = user_info.description + user_info_cache.add_user_info(did, UserInfo(did=did, handle=handle, display_name=display_name, description=description)) + + return did + except Exception as e: + logger.warning(f"Failed to resolve handle '{handle}' to DID: {e}") + return None def load_activated_dids_from_file(client: Client, file_path: str, user_info_cache: UserInfoCache) -> List[str]: """Load activated DIDs from a text file @@ -204,6 +208,7 @@ def load_activated_dids_from_file(client: Client, file_path: str, user_info_cach List of DIDs """ dids = [] + failed_identifiers = [] try: # Create the file if it doesn't exist @@ -214,7 +219,9 @@ def load_activated_dids_from_file(client: Client, file_path: str, user_info_cach return [] with open(file_path, 'r') as f: + line_number = 0 for line in f: + line_number += 1 identifier = line.strip() # Skip empty lines and comments @@ -224,39 +231,69 @@ def load_activated_dids_from_file(client: Client, file_path: str, user_info_cach # If it's already a DID, add it directly if is_did(identifier): dids.append(identifier) + logger.debug(f"Line {line_number}: Added DID directly: {identifier}") # Otherwise, resolve the handle to a DID else: did = resolve_handle_to_did(client, identifier, user_info_cache) if did: dids.append(did) + logger.debug(f"Line {line_number}: Resolved handle '{identifier}' to DID: {did}") + else: + failed_identifiers.append(f"Line {line_number}: {identifier}") + logger.warning(f"Line {line_number}: Failed to resolve identifier: {identifier}") + + # Log summary + if failed_identifiers: + logger.warning(f"Failed to resolve {len(failed_identifiers)} identifiers:") + for failed in failed_identifiers: + logger.warning(f" {failed}") - logger.info(f"Loaded {len(dids)} activated DIDs from {file_path}") + logger.info(f"Successfully loaded {len(dids)} activated DIDs from {file_path}") - # If no DIDs were loaded, raise an error + # If no DIDs were loaded, log an error but don't raise an exception if len(dids) == 0: - logger.error(f"No activated DIDs found in {file_path}") + logger.error(f"No activated DIDs could be loaded from {file_path}") + logger.error("Please check that the file contains valid DIDs or handles") with open(file_path, 'r') as f: - print(f.read()) - raise Exception(f"No activated DIDs found in {file_path}") - + content = f.read() + if content.strip(): + logger.error("File contents:") + print(content) + else: + logger.error("File is empty") + # Return empty list instead of raising exception + return [] return dids except Exception as e: - logger.error(f"Error loading activated DIDs from {file_path}: {e}") - return [] + logger.error(f"Unexpected error while loading activated DIDs from {file_path}: {e}") + # Return whatever DIDs we managed to load before the error + if dids: + logger.info(f"Returning {len(dids)} DIDs that were successfully loaded before the error") + return dids def update_activated_dids(client: Client, file_path: str, user_info_cache: UserInfoCache) -> None: """Update the list of activated DIDs from the file""" global activated_dids try: - activated_dids = load_activated_dids_from_file(client, file_path, user_info_cache) - logger.info(f"Observing {len(activated_dids)} repositories") + new_dids = load_activated_dids_from_file(client, file_path, user_info_cache) + if new_dids: + activated_dids = new_dids + logger.info(f"Observing {len(activated_dids)} repositories") + else: + logger.warning("No DIDs were loaded from the file") + if activated_dids: + logger.info(f"Keeping existing list of {len(activated_dids)} DIDs") + else: + logger.warning("No activated DIDs available - will process all posts but content will be marked as [NOT AVAILABLE]") except Exception as e: logger.error(f"Failed to update activated DIDs: {e}") # Keep existing list if update fails + if activated_dids: + logger.info(f"Keeping existing list of {len(activated_dids)} DIDs") async def process_event(