项目文件夹

文件
wehub-resource-sync c889a57b6b
Test Suites / Build CI Environment (push) Has been cancelled
Test Suites / Basic Tests (push) Has been cancelled
Test Suites / End-to-End Tests (push) Has been cancelled
Test Suites / CLI Tests (push) Has been cancelled
Test Suites / Slow End-to-End Tests (push) Has been cancelled
Test Suites / Graph Database Tests (push) Has been cancelled
Test Suites / Vector DB Tests (push) Has been cancelled
Test Suites / Temporal Graph Test (push) Has been cancelled
Test Suites / Search Test on Different DBs (push) Has been cancelled
Test Suites / Example Tests (push) Has been cancelled
Test Suites / Notebook Tests (push) Has been cancelled
Test Suites / OS and Python Tests Ubuntu (push) Has been cancelled
Test Suites / OS and Python Tests Extended (push) Has been cancelled
Test Suites / LLM Test Suite (push) Has been cancelled
Test Suites / S3 File Storage Test (push) Has been cancelled
Test Suites / Run Integration Tests (push) Has been cancelled
Test Suites / MCP Tests (push) Has been cancelled
Test Suites / Docker Compose Test (push) Has been cancelled
Test Suites / Docker CI test (push) Has been cancelled
Test Suites / Relational DB Migration Tests (push) Has been cancelled
Test Suites / Distributed Cognee Test (push) Has been cancelled
Test Suites / DB Examples Tests (push) Has been cancelled
Test Suites / Test Completion Status (push) Has been cancelled
Test Suites / Claude Code Review (push) Has been cancelled
Test Suites / basic checks (push) Has been cancelled
build | Build and Push Cognee MCP Docker Image to dockerhub / docker-build-and-push (push) Has been cancelled
Scorecard supply-chain security / Scorecard analysis (push) Has been cancelled
build | Build and Push Docker Image to dockerhub / docker-build-and-push (push) Has been cancelled
Weighted Edges Tests / Test Weighted Edges Core Functionality (3.11) (push) Has been cancelled
Weighted Edges Tests / Test Weighted Edges Core Functionality (3.12) (push) Has been cancelled
Weighted Edges Tests / Test Weighted Edges with Different Graph Databases (kuzu, kuzu) (push) Has been cancelled
Weighted Edges Tests / Test Weighted Edges with Different Graph Databases (neo4j, neo4j) (push) Has been cancelled
Weighted Edges Tests / Test Weighted Edges Examples (push) Has been cancelled
Weighted Edges Tests / Code Quality for Weighted Edges (push) Has been cancelled
chore: import upstream snapshot with attribution
2026-07-13 13:02:24 +08:00

148 行
5.0 KiB
Python

"""Utilities for tracking data access in retrievers."""
from datetime import datetime, timezone
from typing import Any, Optional
from uuid import UUID
import os
from cognee.infrastructure.databases.graph import get_graph_engine
from cognee.infrastructure.databases.relational import get_relational_engine
from cognee.modules.data.models import Data
from cognee.shared.logging_utils import get_logger
from sqlalchemy import update
from cognee.modules.graph.cognee_graph.CogneeGraph import CogneeGraph
from cognee.modules.search.utils.transform_triplets_to_graph import transform_triplets_to_graph
from cognee.modules.graph.cognee_graph.CogneeGraphElements import Edge
logger = get_logger(__name__)
async def update_node_access_timestamps(items: Any):
if os.getenv("ENABLE_LAST_ACCESSED", "false").lower() != "true":
return
# In case there are no retrievable node ids, skip processing.
if not items:
return
node_ids = _extract_access_node_ids(items)
if not node_ids:
logger.debug("No valid items to update access timestamps for.")
return
graph_engine = await get_graph_engine()
timestamp_dt = datetime.now(timezone.utc)
# Focus on document-level tracking via projection
try:
doc_ids = await _find_origin_documents_via_projection(graph_engine, node_ids)
if doc_ids:
await _update_sql_records(doc_ids, timestamp_dt)
except Exception as e:
logger.error(f"Failed to update SQL timestamps: {e}")
raise
async def _find_origin_documents_via_projection(graph_engine, node_ids):
"""Find origin documents using graph projection instead of DB queries"""
# Project the entire graph with necessary properties
memory_fragment = CogneeGraph()
await memory_fragment.project_graph_from_db(
graph_engine,
node_properties_to_project=["id", "type"],
edge_properties_to_project=["relationship_name"],
)
# Find origin documents by traversing the in-memory graph
doc_ids = set()
for node_id in node_ids:
node = memory_fragment.get_node(node_id)
if node and node.get_attribute("type") == "DocumentChunk":
# Traverse edges to find connected documents
for edge in node.get_skeleton_edges():
# Get the neighbor node
neighbor = (
edge.get_destination_node()
if edge.get_source_node().id == node_id
else edge.get_source_node()
)
if neighbor and neighbor.get_attribute("type") in ["TextDocument", "Document"]:
doc_ids.add(neighbor.id)
return list(doc_ids)
def _extract_access_node_ids(items: Any) -> list[str]:
if isinstance(items, list) and all(isinstance(item, Edge) for item in items):
return _extract_node_ids_from_edges(items)
if isinstance(items, dict):
return _extract_node_ids_from_retrieval_dict(items)
return []
def _extract_node_ids_from_edges(edges: list[Edge]) -> list[str]:
graph_items = transform_triplets_to_graph(edges)
node_ids = set()
for item in graph_items.get("nodes"):
item_id = item.payload.get("id") if hasattr(item, "payload") else item.get("id")
item_id = _display_id(item_id)
if item_id:
node_ids.add(item_id)
return sorted(node_ids)
def _extract_node_ids_from_retrieval_dict(items: dict) -> list[str]:
node_ids = set()
for chunk in items.get("chunks", []) or []:
chunk_id = _result_id(chunk)
if chunk_id:
node_ids.add(chunk_id)
for entity in items.get("entities", []) or []:
if not isinstance(entity, dict):
continue
entity_id = _display_id(entity.get("id"))
if entity_id:
node_ids.add(entity_id)
for edge in entity.get("edges", []) or []:
if not isinstance(edge, dict):
continue
for key in ("source_id", "target_id"):
edge_node_id = _display_id(edge.get(key))
if edge_node_id:
node_ids.add(edge_node_id)
return sorted(node_ids)
def _result_id(result: Any) -> Optional[str]:
payload = _payload(result)
return _display_id(payload.get("id")) or _display_id(getattr(result, "id", None))
def _payload(result: Any) -> dict:
if isinstance(result, dict):
return result
payload = getattr(result, "payload", None)
return payload if isinstance(payload, dict) else {}
def _display_id(value: Any) -> Optional[str]:
if value is None:
return None
text = str(value).strip()
return text or None
async def _update_sql_records(doc_ids, timestamp_dt):
"""Update SQL Data table (same for all providers)"""
db_engine = get_relational_engine()
async with db_engine.get_async_session() as session:
stmt = (
update(Data)
.where(Data.id.in_([UUID(doc_id) for doc_id in doc_ids]))
.values(last_accessed=timestamp_dt)
)
await session.execute(stmt)
await session.commit()