Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions context-graph/sessions-graph/pyproject.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[project]
name = "sessions-graph"
version = "0.2.0"
version = "0.3.0"
description = "Cross-session memory store for agents, backed by Memgraph"
readme = "README.md"
license = { text = "MIT" }
Expand All @@ -21,7 +21,7 @@ agent-context-graph = [
]
reconciliation = [
"actions-graph>=0.1.1",
"unstructured2graph>=0.4.0",
"unstructured2graph>=0.5.0",
]
test = [
"pytest>=9.0.3",
Expand Down
2 changes: 1 addition & 1 deletion integrations/lightrag-memgraph/pyproject.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[project]
name = "lightrag-memgraph"
version = "0.3.0"
version = "0.3.1"
description = "LightRAG integration with Memgraph"
readme = "README.md"
requires-python = ">=3.10"
Expand Down
20 changes: 18 additions & 2 deletions integrations/lightrag-memgraph/src/lightrag_memgraph/core.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@
import os

from lightrag import LightRAG
from lightrag.kg.shared_storage import initialize_pipeline_status
from lightrag.kg.shared_storage import finalize_share_data, initialize_pipeline_status
from lightrag.llm.openai import gpt_4o_mini_complete, openai_embed
from lightrag.utils import logger, setup_logger

Expand Down Expand Up @@ -114,6 +114,22 @@ async def ainsert(self, **kwargs) -> None:
await self.rag.ainsert(**kwargs)

async def afinalize(self) -> None:
"""
Finalize storages and reset LightRAG's process-global shared-storage
state (locks, shared dicts). Without the latter reset, a second
LightRAG instance created later in the same process reuses asyncio
locks bound to this instance's (by-then-closed) event loop and raises
"bound to a different event loop" -- e.g. two real-LightRAG tests
run back-to-back under pytest-asyncio's function-scoped event loop.

The reset runs in a finally block: if finalize_storages() raises
(e.g. a transient network error), the shared-storage state would
otherwise stay registered as initialized and break every subsequent
LightRAG instance in this process, not just this one's cleanup.
"""
if self.rag is None:
raise RuntimeError("LightRAG not initialized. Call initialize() first.")
await self.rag.finalize_storages()
try:
await self.rag.finalize_storages()
finally:
finalize_share_data()
34 changes: 33 additions & 1 deletion integrations/lightrag-memgraph/tests/test_core.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,12 +9,13 @@
from __future__ import annotations

import os
from unittest.mock import AsyncMock, patch

import pytest
from lightrag.llm.openai import gpt_4o_mini_complete, openai_embed
from lightrag.utils import logger as lightrag_logger

from lightrag_memgraph.core import _apply_lightrag_defaults, _bridge_lightrag_env_names
from lightrag_memgraph.core import MemgraphLightRAGWrapper, _apply_lightrag_defaults, _bridge_lightrag_env_names
from lightrag_memgraph.embeddings import memgraph_sentence_embed

_MEMGRAPH_ENV_NAMES = ("MEMGRAPH_URL", "MEMGRAPH_URI", "MEMGRAPH_USER", "MEMGRAPH_USERNAME")
Expand Down Expand Up @@ -136,3 +137,34 @@ def test_no_openai_key_required_with_non_openai_llm_and_embedding_default(monkey
kwargs = {"llm_model_func": lambda *a, **kw: None}
_apply_lightrag_defaults(kwargs) # must not raise
assert kwargs["embedding_func"] is memgraph_sentence_embed


@pytest.mark.asyncio
async def test_afinalize_resets_shared_data_on_success():
wrapper = MemgraphLightRAGWrapper()
wrapper.rag = AsyncMock()

with patch("lightrag_memgraph.core.finalize_share_data") as mock_finalize_share_data:
await wrapper.afinalize()

wrapper.rag.finalize_storages.assert_awaited_once()
mock_finalize_share_data.assert_called_once()


@pytest.mark.asyncio
async def test_afinalize_resets_shared_data_even_if_finalize_storages_raises():
"""A transient error in finalize_storages() must not leave the
process-global shared-storage state stuck as initialized -- that would
break every subsequent LightRAG instance in this process, not just this
one's cleanup."""
wrapper = MemgraphLightRAGWrapper()
wrapper.rag = AsyncMock()
wrapper.rag.finalize_storages.side_effect = RuntimeError("transient network error")

with (
patch("lightrag_memgraph.core.finalize_share_data") as mock_finalize_share_data,
pytest.raises(RuntimeError, match="transient network error"),
):
await wrapper.afinalize()

mock_finalize_share_data.assert_called_once()
3 changes: 2 additions & 1 deletion unstructured2graph/pyproject.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[project]
name = "unstructured2graph"
version = "0.4.0"
version = "0.5.0"
description = "Convert unstructured documents into knowledge graphs"
readme = "README.md"
requires-python = ">=3.10"
Expand All @@ -25,6 +25,7 @@ dependencies = [
"unstructured>=0.18.18",
"lightrag-memgraph>=0.3.0",
"memgraph-toolbox>=0.1.11",
"pyyaml>=6.0",
]
[project.optional-dependencies]
all-docs = [
Expand Down
10 changes: 9 additions & 1 deletion unstructured2graph/src/unstructured2graph/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -23,12 +23,18 @@
create_unique_constraint,
create_vector_search_index,
link_nodes_in_order,
promote_entity_types_to_labels,
)
from .ontology import DEFAULT_ONTOLOGY, DEFAULT_ONTOLOGY_PATH, EntityType, Ontology, load_ontology

__version__ = "0.4.0"
__version__ = "0.5.0"
__all__ = [
"DEFAULT_ONTOLOGY",
"DEFAULT_ONTOLOGY_PATH",
"Chunk",
"ChunkedDocument",
"EntityType",
"Ontology",
"compute_embeddings",
"connect_chunks_to_entities",
"create_label_index",
Expand All @@ -39,7 +45,9 @@
"from_texts",
"from_unstructured",
"link_nodes_in_order",
"load_ontology",
"make_chunks",
"parse_source",
"parse_text",
"promote_entity_types_to_labels",
]
27 changes: 27 additions & 0 deletions unstructured2graph/src/unstructured2graph/default_ontology.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,27 @@
# Default entity-type ontology. Mirrors LightRAG's own built-in vocabulary
# (lightrag.prompt's default_entity_types_guidance), so label promotion
# matches what LightRAG extracts by default even before this file is
# customized or overridden with a project-specific ontology_path.
entity_types:
- label: Person
description: Human individuals, real or fictional
- label: Creature
description: Non-human living beings (animals, mythical beings, etc.)
- label: Organization
description: Companies, institutions, government bodies, groups
- label: Location
description: Geographic places (cities, countries, buildings, regions)
- label: Event
description: Occurrences, incidents, ceremonies, meetings
- label: Concept
description: Abstract ideas, theories, principles, beliefs
- label: Method
description: Procedures, techniques, algorithms, workflows
- label: Content
description: Creative or informational works (books, articles, films, reports)
- label: Data
description: Quantitative or structured information (statistics, datasets, measurements)
- label: Artifact
description: Physical or digital objects created by humans (tools, software, devices)
- label: NaturalObject
description: Natural non-living objects (minerals, celestial bodies, chemical compounds)
47 changes: 45 additions & 2 deletions unstructured2graph/src/unstructured2graph/loaders.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,9 @@
create_nodes_from_list,
create_unique_constraint,
link_nodes_in_order,
promote_entity_types_to_labels,
)
from .ontology import DEFAULT_ONTOLOGY, load_ontology

SCRIPT_DIR = os.path.dirname(os.path.realpath(__file__))
logger = logging.getLogger(__name__)
Expand Down Expand Up @@ -165,12 +167,15 @@ async def _ingest_chunks(
only_chunks: bool = False,
link_chunks: bool = False,
entity_workspace: str | None = None,
enforce_ontology: bool = False,
ontology_path: str | Path | None = None,
) -> list[Chunk]:
"""
Ingest an already-produced flat list of chunks into Memgraph: upsert Chunk
nodes, optionally chain them with NEXT, and (unless only_chunks) run
LightRAG entity extraction and connect the resulting entities back to
their chunks via MENTIONED_IN.
LightRAG entity extraction, connect the resulting entities back to their
chunks via MENTIONED_IN, and (if enforce_ontology) promote entity_type to
a real label for entities that match the ontology.

Internal helper shared by from_unstructured() and from_texts(). Not
exported: it relies on its caller having already ensured the Chunk.hash
Expand All @@ -188,6 +193,11 @@ async def _ingest_chunks(
only_chunks: If True, only create chunk nodes without LightRAG processing.
link_chunks: If True, link chunks in order with NEXT relationship.
entity_workspace: Node label LightRAG entities were written under.
enforce_ontology: If True, promote entity_type to labels per ontology_path (or
DEFAULT_ONTOLOGY_PATH). If False (default), entities are left exactly as
LightRAG wrote them -- no label promotion, no ontology_conformant flagging.
ontology_path: Path to an ontology YAML config file. Only consulted when
enforce_ontology=True; defaults to DEFAULT_ONTOLOGY_PATH.
Returns:
The same chunks that were passed in, for convenience chaining.
"""
Expand All @@ -198,6 +208,9 @@ async def _ingest_chunks(
if not only_chunks and lightrag_wrapper is None:
raise ValueError("lightrag_wrapper is required when only_chunks=False")

if ontology_path and not enforce_ontology:
logger.warning("ontology_path was provided but enforce_ontology=False; ignoring ontology_path")

memgraph_node_props = []
for chunk in chunks:
logger.debug(f"Chunk: {chunk.hash} - {chunk.text}")
Expand All @@ -214,6 +227,9 @@ async def _ingest_chunks(
for chunk in chunks:
await lightrag_wrapper.ainsert(input=chunk.text, file_paths=[chunk.hash])
connect_chunks_to_entities(memgraph, "Chunk", entity_workspace)
if enforce_ontology:
ontology = load_ontology(ontology_path) if ontology_path else DEFAULT_ONTOLOGY
promote_entity_types_to_labels(memgraph, entity_workspace, ontology)

return chunks

Expand All @@ -224,6 +240,8 @@ async def from_texts(
lightrag_wrapper: MemgraphLightRAGWrapper | None = None,
only_chunks: bool = False,
entity_workspace: str | None = None,
enforce_ontology: bool = False,
ontology_path: str | Path | None = None,
) -> list[list[Chunk]]:
"""
Ingest raw in-memory strings (not files or URLs) into Memgraph.
Expand All @@ -241,6 +259,20 @@ async def from_texts(
entity_workspace: Node label LightRAG entities were written under. If None
(default), auto-derived from lightrag_wrapper's resolved LightRAG
workspace, falling back to "base" if that fails.
enforce_ontology: If False (default), entities are left exactly as LightRAG
wrote them -- no label promotion, no ontology_conformant flagging. If True,
entity_type gets promoted to a real Memgraph label (e.g. entity_type="person"
-> :Person) in addition to the entity_workspace label every entity already
gets, per ontology_path.
ontology_path: Path to an ontology YAML config file (see load_ontology()). Only
consulted when enforce_ontology=True; defaults to DEFAULT_ONTOLOGY_PATH,
which mirrors LightRAG's own built-in type vocabulary. entity_type values
outside the ontology are never rejected -- the node and its entity_type
property are kept, stamped ontology_conformant=false instead of getting a
label. To also steer LightRAG's extraction itself toward the same
vocabulary, load the same path with load_ontology() and pass its
addon_params() into MemgraphLightRAGWrapper.initialize() -- using the same
path at both call sites is what keeps them in sync.
Returns:
One list of Chunks per input text, in input order. A text that
parse_text() splits into several pieces contributes several Chunks in
Expand Down Expand Up @@ -269,6 +301,8 @@ async def from_texts(
only_chunks=only_chunks,
link_chunks=False,
entity_workspace=resolved_entity_workspace,
enforce_ontology=enforce_ontology,
ontology_path=ontology_path,
)
return grouped_chunks

Expand All @@ -281,6 +315,8 @@ async def from_unstructured(
link_chunks: bool = False,
entity_workspace: str | None = None,
partition_kwargs: dict[str, Any] | None = None,
enforce_ontology: bool = False,
ontology_path: str | Path | None = None,
) -> list[list[Chunk]]:
"""
Process unstructured sources and ingest them into Memgraph using LightRAG.
Expand All @@ -297,6 +333,11 @@ async def from_unstructured(
partition_kwargs: Additional keyword arguments to pass to unstructured's
partition function (e.g., strategy, languages, pdf_infer_table_structure,
ocr_languages, headers, ssl_verify, etc.)
enforce_ontology: If False (default), no label promotion or ontology_conformant
flagging happens. If True, entity_type gets promoted to a real Memgraph
label per ontology_path. See from_texts() for details.
ontology_path: Path to an ontology YAML config file. Only consulted when
enforce_ontology=True; defaults to DEFAULT_ONTOLOGY_PATH.
Returns:
One list of Chunks per source, in `sources` order — the same
grouped-return contract as from_texts(). A source that produced no
Expand Down Expand Up @@ -328,6 +369,8 @@ async def from_unstructured(
only_chunks=only_chunks,
link_chunks=link_chunks,
entity_workspace=resolved_entity_workspace,
enforce_ontology=enforce_ontology,
ontology_path=ontology_path,
)
grouped_chunks.append(document.chunks)

Expand Down
41 changes: 41 additions & 0 deletions unstructured2graph/src/unstructured2graph/memgraph.py
Original file line number Diff line number Diff line change
@@ -1,9 +1,13 @@
import logging
import time
from typing import TYPE_CHECKING

from lightrag_memgraph import DEFAULT_EMBEDDING_DIM
from memgraph_toolbox.api.memgraph import Memgraph

if TYPE_CHECKING:
from .ontology import Ontology

logger = logging.getLogger(__name__)


Expand Down Expand Up @@ -69,6 +73,43 @@ def connect_chunks_to_entities(memgraph: Memgraph, chunk_label: str, entity_labe
)


def promote_entity_types_to_labels(memgraph: Memgraph, workspace_label: str, ontology: "Ontology") -> None:
"""
Additively promote each entity's `entity_type` property to a real
Memgraph label (e.g. entity_type="person" -> :Person), for entity_type
values that match the given ontology.

The workspace label is never touched: LightRAG's own upsert_node()
re-MERGEs future updates by matching on it, so removing it would break
LightRAG's ability to recognize this node on subsequent re-ingestion.
Entities whose entity_type doesn't match any type in the ontology are
never rejected -- the node and its raw entity_type are always kept, and
are instead stamped `ontology_conformant: false` so what the ontology
doesn't recognize stays visible and queryable rather than silently
indistinguishable from an unprocessed node. Re-running this (e.g. after
the ontology grows a new type) clears the flag on anything that now
conforms.
"""
labels = ontology.allowed_labels()
for label in labels:
memgraph.query(
f"""
MATCH (n:{workspace_label})
WHERE toLower(n.entity_type) = toLower($label) AND NOT n:{label}
SET n:{label}
""",
params={"label": label},
)
Comment thread
Copilot marked this conversation as resolved.

if not labels:
memgraph.query(f"MATCH (n:{workspace_label}) SET n.ontology_conformant = false")
return

conforms_clause = " OR ".join(f"n:{label}" for label in labels)
memgraph.query(f"MATCH (n:{workspace_label}) WHERE NOT ({conforms_clause}) SET n.ontology_conformant = false")
memgraph.query(f"MATCH (n:{workspace_label}) WHERE {conforms_clause} REMOVE n.ontology_conformant")


def link_nodes_in_order(
memgraph: Memgraph,
find_label: str,
Expand Down
Loading
Loading