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
6 changes: 4 additions & 2 deletions context-graph/sessions-graph/src/sessions_graph/core.py
Original file line number Diff line number Diff line change
Expand Up @@ -64,7 +64,9 @@ def combined_text(self) -> str:

def document(self, user_id: str) -> Document:
"""combined_text as an unstructured2graph Document: one segment per deduped
source, carrying its speaker and timestamp, and the session's user."""
source, carrying its speaker, timestamp and node id, and the session's user.

A text repeated across sources is one segment, attributed to its first source."""
from unstructured2graph import Document, Segment

first: dict[str, ReconciliationSource] = {}
Expand All @@ -73,7 +75,7 @@ def document(self, user_id: str) -> Document:
segments, cursor = [], 0
for digest, text in self.unique_texts.items():
source = first[digest]
segments.append(Segment(cursor, cursor + len(text), source.role, source.valid_at))
segments.append(Segment(cursor, cursor + len(text), source.role, source.valid_at, source.node_id))
cursor += len(text) + 2
return Document(text=self.combined_text, segments=tuple(segments), user_id=user_id)

Expand Down
15 changes: 14 additions & 1 deletion context-graph/sessions-graph/tests/test_reconciliation.py
Original file line number Diff line number Diff line change
Expand Up @@ -794,7 +794,8 @@ async def test_reconcile_session_with_gliner2_binds_the_user_and_stamps_the_turn
):
"""The typed relation model end to end: turns become segments, the user's
own "I" binds to a synthesized (:User), the edge carries the turn's
timestamp as a datetime, and the integrity counts are zero."""
timestamp as a datetime and the turn it came from, mentions record their
turns, and the integrity counts are zero."""
from actions_graph import Session
from unstructured2graph import load_ontology
from unstructured2graph.gliner2_backend import GLiNER2Backend
Expand Down Expand Up @@ -844,3 +845,15 @@ async def test_reconcile_session_with_gliner2_binds_the_user_and_stamps_the_turn
)
assert rows == [{"user_id": "anon-s-1", "place": "Paris", "valid_at": "2023-05-30T17:27:00.000000+00:00"}]
assert backend.stats.mentions_dropped == 1 # the assistant's "I"

turns = {
row["role"]: row["id"]
for row in memgraph.query(
"MATCH (a:UserMessage) RETURN 'user' AS role, a.action_id AS id "
"UNION MATCH (a:AssistantMessage) RETURN 'assistant' AS role, a.action_id AS id"
)
}
edge = memgraph.query("MATCH ()-[r:visited]->() RETURN r.source_id AS source_id, r.role AS role, r.text AS text")
assert edge == [{"source_id": turns["user"], "role": "user", "text": "user: I visited Paris."}]
sources = memgraph.query("MATCH (:Location {text: 'Paris'})-[m:MENTIONED_IN]->(:Chunk) RETURN m.sources AS sources")
assert sources == [{"sources": sorted(turns.values())}]
6 changes: 3 additions & 3 deletions unstructured2graph/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -103,15 +103,15 @@ await from_unstructured(

It extracts on gliner2's joint path, where each relation type's `start_labels`/`end_labels` constrain decoding. It extracts one window per segment, splitting only a segment longer than `chunk_size` words. It writes:

- entity nodes, merged per the type's `identity` (see the ontology section), each linked to every chunk that mentions it by `MENTIONED_IN`;
- relationships keyed by source chunk, carrying `chunk`, `confidence`, and `valid_at` (a Memgraph `datetime`, from the segment's timestamp);
- entity nodes, merged per the type's `identity` (see the ontology section), each linked to every chunk that mentions it by `MENTIONED_IN`, whose `sources` lists the `source_id` of each segment that mentions it;
- relationships keyed by source chunk and, when the segment has one, its `source_id`, so a fact said in two turns is two relationships. Each carries `chunk`, `confidence`, `valid_at` (a Memgraph `datetime`, from the segment's timestamp), `source_id`, `role` (the speaker) and `text` (the sentence or sentences covering both endpoints, at most 300 characters);
- the chunk's user's own mentions onto `(:User {user_id})`, decided by the `mention_resolver` (default `resolve_user_mentions`). The caller must create that node.

Self-loops left after identity resolution are dropped. `backend.stats` counts windows, re-typed and dropped mentions, and self-loops.

### Ingesting documents verbatim, with segments

`from_documents` stores each `Document`'s text exactly as given, as one `Chunk`, instead of re-chunking it through `unstructured` (whose partitioner rewrites text). Its `segments` travel with the chunk: for a conversation, one per turn, with `role` and `valid_at`. The GLiNER2 backend uses them; LightRAG reads the text alone.
`from_documents` stores each `Document`'s text exactly as given, as one `Chunk`, instead of re-chunking it through `unstructured` (whose partitioner rewrites text). Its `segments` travel with the chunk: for a conversation, one per turn, with `role`, `valid_at` and an opaque `source_id` naming the turn's node. The GLiNER2 backend uses them; LightRAG reads the text alone.

```python
from unstructured2graph import Document, Segment, from_documents
Expand Down
74 changes: 64 additions & 10 deletions unstructured2graph/src/unstructured2graph/gliner2_backend.py
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,11 @@
- A relationship carries the source chunk, the source turn's timestamp as
`valid_at` (#364) and the model's confidence. Self-loops left after identity
resolution are dropped (#355).
- Provenance is exact, from the spans the model returns (#392): a segment's
`source_id` (the turn) is recorded in `MENTIONED_IN.sources` for every
entity it mentions, and each relationship extracted from it is its own edge
carrying `source_id`, the speaker as `role`, and the sentence covering both
endpoints as `text`.
"""

import asyncio
Expand Down Expand Up @@ -222,6 +227,34 @@ def _word_windows(text: str, start: int, end: int, size: int, overlap: int) -> l
return ranges


#: A sentence ends at terminal punctuation followed by whitespace, or at a line break.
_SENTENCE_END = re.compile(r"(?<=[.!?])\s+|\n+")

#: The longest source text a relationship carries; the full turn is one hop
#: away through its source_id.
SOURCE_TEXT_CAP = 300


def _source_text(text: str, bounds: tuple[int, int], head: tuple[int, int], tail: tuple[int, int]) -> str:
"""The sentence or sentences of text[bounds] covering both endpoint spans, at most SOURCE_TEXT_CAP characters.

Past the cap the text is trimmed around the spans, keeping both whole.
"""
lo, hi = min(head[0], tail[0]), max(head[1], tail[1])
start, end = bounds
for boundary in _SENTENCE_END.finditer(text, bounds[0], bounds[1]):
if boundary.end() <= lo:
start = boundary.end()
elif boundary.start() >= hi:
end = boundary.start()
break
if end - start > SOURCE_TEXT_CAP:
slack = max(SOURCE_TEXT_CAP - (hi - lo), 0)
start = max(start, lo - slack // 2)
end = min(end, max(hi, start + SOURCE_TEXT_CAP))
return text[start:end].strip()


class GLiNER2Backend:
"""Local, LLM-free ExtractionBackend backed by a GLiNER2 model.

Expand Down Expand Up @@ -427,8 +460,8 @@ async def aingest_chunk(self, memgraph: Memgraph, chunk: "Chunk") -> None:
self.stats.infeasible_windows += extracted.infeasible

endpoints: list[Endpoint | None] = []
valid_at: list[str | None] = []
nodes: dict[str, dict[str, Any]] = {}
sources: dict[str, set[str]] = {}
for mention, segment in extracted.mentions:
self.stats.mentions += 1
if (
Expand All @@ -437,7 +470,6 @@ async def aingest_chunk(self, memgraph: Memgraph, chunk: "Chunk") -> None:
and mention.confidence < self.entity_confidence_threshold
):
endpoints.append(None)
valid_at.append(None)
continue
resolution = self._resolve(mention, segment, chunk)
if resolution.action == "drop":
Expand All @@ -448,15 +480,24 @@ async def aingest_chunk(self, memgraph: Memgraph, chunk: "Chunk") -> None:
self.stats.mentions_retyped += 1
resolved = self._endpoint(mention, resolution, chunk)
endpoints.append(resolved[0] if resolved else None)
valid_at.append(segment.valid_at if segment is not None else None)
if resolved and resolved[1] is not None:
nodes.setdefault(resolved[1]["entity_id"], resolved[1])
entity_id = resolved[1]["entity_id"]
nodes.setdefault(entity_id, resolved[1])
mentioned_in = sources.setdefault(entity_id, set())
if segment is not None and segment.source_id is not None:
mentioned_in.add(segment.source_id)

if nodes:
create_nodes_from_list(memgraph, list(nodes.values()), self._workspace, 100, merge_key="entity_id")
link_mentions(memgraph, self._workspace, "entity_id", chunk.hash, list(nodes))
link_mentions(
memgraph,
self._workspace,
"entity_id",
chunk.hash,
{entity_id: list(ids) for entity_id, ids in sources.items()},
)

merged: dict[tuple[str, Endpoint, Endpoint], dict[str, Any]] = {}
merged: dict[tuple[str, Endpoint, Endpoint, str | None], dict[str, Any]] = {}
for relation_type, head_index, tail_index, confidence in extracted.relations:
head, tail = endpoints[head_index], endpoints[tail_index]
if head is None or tail is None:
Expand All @@ -472,22 +513,35 @@ async def aingest_chunk(self, memgraph: Memgraph, chunk: "Chunk") -> None:
# resolution then makes them one node (#355).
self.stats.self_loops_dropped += 1
continue
when = valid_at[head_index]
key = (relation_type, head, tail)
head_mention, segment = extracted.mentions[head_index]
tail_mention = extracted.mentions[tail_index][0]
# Relations never cross windows, so both endpoints share the head's segment.
when = segment.valid_at if segment is not None else None
source_id = segment.source_id if segment is not None else None
key = (relation_type, head, tail, source_id)
text = _source_text(
chunk.text,
(segment.start, segment.end) if segment is not None else (0, len(chunk.text)),
(head_mention.start, head_mention.end),
(tail_mention.start, tail_mention.end),
)
previous = merged.get(key)
if previous is None:
merged[key] = {
"type": relation_type,
"head": head,
"tail": tail,
"chunk": chunk.hash,
"source_id": source_id,
"valid_at": when,
"confidence": confidence,
"text": text,
"role": segment.role if segment is not None else None,
}
continue
# The same fact twice in one chunk: when it was first said, how sure the model ever was.
# The same fact twice under one key: when and where it was first said, how sure the model ever was.
if when is not None and (previous["valid_at"] is None or when < previous["valid_at"]):
previous["valid_at"] = when
previous.update(valid_at=when, text=text, role=segment.role if segment is not None else None)
if confidence is not None and (previous["confidence"] is None or confidence > previous["confidence"]):
previous["confidence"] = confidence

Expand Down
5 changes: 5 additions & 0 deletions unstructured2graph/src/unstructured2graph/loaders.py
Original file line number Diff line number Diff line change
Expand Up @@ -45,12 +45,17 @@ class Segment:
that is nobody's utterance (a tool result).
valid_at: ISO-8601 timestamp of when it was produced. Relationships
extracted from it carry this as `valid_at` (#364).
source_id: Opaque id of the node the segment's text came from (e.g. a
conversation turn), never interpreted here. When set, mentions
record it in `MENTIONED_IN.sources` and each relationship extracted
from the segment is its own edge carrying it as `source_id` (#392).
"""

start: int
end: int
role: str | None = None
valid_at: str | None = None
source_id: str | None = None


@dataclass
Expand Down
63 changes: 47 additions & 16 deletions unstructured2graph/src/unstructured2graph/memgraph.py
Original file line number Diff line number Diff line change
Expand Up @@ -143,25 +143,44 @@ def _entities(workspace_label: str, chunk_hashes: list[str] | None) -> str:
)


def link_mentions(memgraph: Memgraph, entity_label: str, match_key: str, chunk_hash: str, entity_ids: list[str]):
"""MERGE (entity)-[:MENTIONED_IN]->(:Chunk {hash: chunk_hash}) for each id, at ingest.
def link_mentions(
memgraph: Memgraph,
entity_label: str,
match_key: str,
chunk_hash: str,
sources_by_entity: dict[str, list[str]],
):
"""MERGE (entity)-[:MENTIONED_IN]->(:Chunk {hash: chunk_hash}) for each entity id, at ingest.

Needed once an entity's identity outlives its chunk (#346): a merged node
keeps only its first chunk's file_path, so connect_chunks_to_entities alone
would link it to that one chunk and lose every later mention.

Args:
sources_by_entity: entity id -> the `Segment.source_id`s of the
segments mentioning it in this chunk, unioned into
`MENTIONED_IN.sources` so re-ingesting adds nothing. An empty list
links the entity without touching `sources`.
"""
_require_valid_identifier(entity_label, "entity_label")
_require_valid_identifier(match_key, "match_key")
if not entity_ids:
if not sources_by_entity:
return
rows = [
{"id": entity_id, "sources": sorted(set(sources))} for entity_id, sources in sorted(sources_by_entity.items())
]
memgraph.query(
f"""
MATCH (c:Chunk {{hash: $chunk_hash}})
UNWIND $ids AS id
MATCH (n:{entity_label} {{{match_key}: id}})
MERGE (n)-[:MENTIONED_IN]->(c)
UNWIND $rows AS row
MATCH (n:{entity_label} {{{match_key}: row.id}})
MERGE (n)-[r:MENTIONED_IN]->(c)
SET r.sources = CASE
WHEN size(row.sources) = 0 THEN r.sources
ELSE [s IN coalesce(r.sources, []) WHERE NOT s IN row.sources] + row.sources
END
""",
params={"chunk_hash": chunk_hash, "ids": sorted(set(entity_ids))},
params={"chunk_hash": chunk_hash, "rows": rows},
)


Expand Down Expand Up @@ -472,35 +491,42 @@ class Endpoint:

def upsert_extracted_relationships(memgraph: Memgraph, relationships: list[dict[str, Any]]) -> None:
"""
MERGE extracted relationships, one per (type, head, tail, source chunk).
MERGE extracted relationships, one per (type, head, tail, source chunk, source_id).

Keyed on the source chunk as well as the endpoints, so a fact asserted in
two sessions is two relationships with their own valid_at: superseded
facts are retained, timestamps only (#347), and a question asking for the
*initial* value still has it. Re-ingesting a chunk is idempotent.
*initial* value still has it. With a `source_id` the key also takes the
source turn, so a fact said in two turns of one session is two
relationships too (#392). Re-ingesting a chunk is idempotent.

Args:
relationships: dicts with `type`, `head`/`tail` (Endpoint), `chunk`
(the source chunk hash), `valid_at` (ISO-8601 string or None, stored
as a Memgraph datetime so date arithmetic is a subtraction, #364)
and `confidence` (float or None).
as a Memgraph datetime so date arithmetic is a subtraction, #364),
`confidence` (float or None), and optionally `source_id`, `text`
(the source sentence) and `role` (the speaker), each str or None.

Raises:
ValueError: if a relationship type, endpoint label or key isn't a valid Cypher identifier.
"""
groups: dict[tuple[str, str, str, str, str], list[dict[str, Any]]] = defaultdict(list)
groups: dict[tuple[str, str, str, str, str, bool], list[dict[str, Any]]] = defaultdict(list)
for rel in relationships:
head, tail = rel["head"], rel["tail"]
groups[(rel["type"], head.label, head.key, tail.label, tail.key)].append(
has_source = rel.get("source_id") is not None
groups[(rel["type"], head.label, head.key, tail.label, tail.key, has_source)].append(
{
"from": head.value,
"to": tail.value,
"chunk": rel["chunk"],
"source_id": rel.get("source_id"),
"valid_at": rel.get("valid_at"),
"confidence": rel.get("confidence"),
"text": rel.get("text"),
"role": rel.get("role"),
}
)
for (relation_type, head_label, head_key, tail_label, tail_key), rows in groups.items():
for (relation_type, head_label, head_key, tail_label, tail_key, has_source), rows in groups.items():
for value, role in (
(relation_type, "relation type"),
(head_label, "endpoint label"),
Expand All @@ -509,14 +535,19 @@ def upsert_extracted_relationships(memgraph: Memgraph, relationships: list[dict[
(tail_key, "endpoint key"),
):
_require_valid_identifier(value, role)
# MERGE can't key on a null property, so a relationship with no source
# turn keeps the chunk-only key.
key = "{chunk: rel.chunk, source_id: rel.source_id}" if has_source else "{chunk: rel.chunk}"
memgraph.query(
f"""
UNWIND $rows AS rel
MATCH (a:{head_label} {{{head_key}: rel.from}})
MATCH (b:{tail_label} {{{tail_key}: rel.to}})
MERGE (a)-[r:{relation_type} {{chunk: rel.chunk}}]->(b)
MERGE (a)-[r:{relation_type} {key}]->(b)
SET r.valid_at = CASE WHEN rel.valid_at IS NULL THEN null ELSE datetime(rel.valid_at) END,
r.confidence = rel.confidence
r.confidence = rel.confidence,
r.text = rel.text,
r.role = rel.role
""",
params={"rows": rows},
)
Expand Down
Loading
Loading