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
32 changes: 31 additions & 1 deletion backend/app/api/admin_routes/knowledge_base/graph/routes.py
Original file line number Diff line number Diff line change
@@ -1,7 +1,10 @@
import logging
from typing import List
import json

from fastapi import APIRouter, HTTPException, status
from fastapi.responses import StreamingResponse
from fastapi.encoders import jsonable_encoder

from app.api.admin_routes.knowledge_base.graph.models import (
SynopsisEntityCreate,
Expand Down Expand Up @@ -259,4 +262,31 @@ def get_entire_knowledge_graph(session: SessionDep, kb_id: int):
raise e
except Exception as e:
# TODO: throw InternalServerError
raise e
raise e

@router.get("/admin/knowledge_bases/{kb_id}/graph/entire_graph/stream")
def stream_entire_knowledge_graph(session: SessionDep, kb_id: int):
try:
kb = knowledge_base_repo.must_get(session, kb_id)
graph_store = get_kb_tidb_graph_store(session, kb)

def generate():
for chunk in graph_store.stream_entire_knowledge_graph(chunk_size=5000):
yield f"data: {json.dumps(jsonable_encoder(chunk))}\n\n"
yield f"data: {json.dumps({'type': 'complete'})}\n\n"

return StreamingResponse(
generate(),
media_type="text/event-stream",
headers={
"Cache-Control": "no-cache",
"Connection": "keep-alive",
"Access-Control-Allow-Origin": "*",
}
)

except KBNotFound as e:
raise e
except Exception as e:
logger.exception(e)
raise InternalServerError()
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@
from llama_index.embeddings.openai import OpenAIEmbedding, OpenAIEmbeddingModelType
import sqlalchemy
from sqlmodel import Session, asc, func, select, text, SQLModel
from sqlalchemy.orm import aliased, defer, joinedload
from sqlalchemy.orm import aliased, defer, joinedload, noload
from tidb_vector.sqlalchemy import VectorAdaptor
from sqlalchemy import or_, desc

Expand Down Expand Up @@ -1182,3 +1182,91 @@ def get_entire_knowledge_graph(self) -> RetrievedKnowledgeGraph:
entities=entities,
relationships=relationships,
)

def stream_entire_knowledge_graph(self, chunk_size: int = 5000):
"""Stream entire knowledge graph in chunks

Args:
chunk_size: Number of entities/relationships per chunk

Yields:
Dict containing chunk type and data
"""
# Stream entities
entity_query = (
select(self._entity_model)
.options(
defer(self._entity_model.description_vec),
defer(self._entity_model.meta_vec),
)
.order_by(self._entity_model.id)
)
last_entity_id = 0

while True:
chunk_query = entity_query.where(
self._entity_model.id > last_entity_id
).limit(chunk_size)
db_entities = self._session.exec(chunk_query).all()

if not db_entities:
break

entities = []
for entity in db_entities:
entities.append(
RetrievedEntity(
id=entity.id,
knowledge_base_id=self.knowledge_base.id,
name=entity.name,
description=entity.description,
meta=entity.meta,
entity_type=entity.entity_type,
)
)

last_entity_id = db_entities[-1].id
yield {"type": "entities", "data": entities}

# Stream relationships
relationship_query = (
select(self._relationship_model)
.options(
defer(self._relationship_model.description_vec),
defer(self._relationship_model.chunk_id),
noload(self._relationship_model.source_entity),
noload(self._relationship_model.target_entity),
)
.order_by(self._relationship_model.id)
)
logger.info(f"Relationship query: {relationship_query}")
last_relationship_id = 0

while True:
chunk_query = relationship_query.where(
self._relationship_model.id > last_relationship_id
).limit(chunk_size)
logger.info(f"Executing relationship chunk query: {chunk_query}")
db_relationships = self._session.exec(chunk_query).all()

if not db_relationships:
break

relationships = []
for rel in db_relationships:
relationships.append(
RetrievedRelationship(
id=rel.id,
knowledge_base_id=self.knowledge_base.id,
source_entity_id=rel.source_entity_id,
target_entity_id=rel.target_entity_id,
description=rel.description,
rag_description=None, # Skip rag_description for streaming performance
meta=rel.meta,
weight=rel.weight,
last_modified_at=rel.last_modified_at,
)
)

last_relationship_id = db_relationships[-1].id
yield {"type": "relationships", "data": relationships}
55 changes: 55 additions & 0 deletions frontend/app/src/api/graph.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import { authenticationHeaders, handleResponse, requestUrl } from '@/lib/request';
import { zodJsonDate } from '@/lib/zod';
import { bufferedReadableStreamTransformer } from '@/lib/buffered-readable-stream';
import { z, type ZodType } from 'zod';

export interface KnowledgeGraph {
Expand Down Expand Up @@ -181,6 +182,60 @@ export async function getEntireKnowledgeGraph (kbId: number, params: KBRetrieveK
.then(handleResponse(knowledgeGraphSchema));
}

export async function streamEntireKnowledgeGraph (kbId: number): Promise<KnowledgeGraph> {
const entities: KnowledgeGraphEntity[] = [];
const relationships: KnowledgeGraphRelationship[] = [];

const response = await fetch(requestUrl(`/api/v1/admin/knowledge_bases/${kbId}/graph/entire_graph/stream`), {
method: 'GET',
headers: await authenticationHeaders(),
credentials: 'include',
});

if (!response.ok) {
throw new Error(`${response.status} ${response.statusText}`);
}

if (!response.body) {
throw new Error('Empty response body');
}

const reader = response.body.pipeThrough(bufferedReadableStreamTransformer()).getReader();

try {
while (true) {
const { done, value } = await reader.read();
if (done) break;

if (value.trim() && value.startsWith('data: ')) {
const dataStr = value.substring(6).trim();
if (dataStr) {
try {
const data = JSON.parse(dataStr);

if (data.type === 'entities') {
entities.push(...data.data);
// console.log(`Received ${data.data.length} entities, total: ${entities.length}`);
} else if (data.type === 'relationships') {
relationships.push(...data.data);
// console.log(`Received ${data.data.length} relationships, total: ${relationships.length}`);
} else if (data.type === 'complete') {
// console.log(`Streaming complete. Final counts - entities: ${entities.length}, relationships: ${relationships.length}`);
return { entities, relationships };
}
} catch (error) {
console.warn('Failed to parse streaming data:', error, 'Data:', dataStr);
}
}
}
}
} finally {
reader.releaseLock();
}

return { entities, relationships };
}

export async function getRelationship (kbId: number, id: number) {
return await fetch(requestUrl(`/api/v1/admin/knowledge_bases/${kbId}/graph/relationships/${id}`), {
headers: {
Expand Down
33 changes: 22 additions & 11 deletions frontend/app/src/components/graph/GraphEditor.tsx
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
'use client';

import { getChatMessageSubgraph } from '@/api/chats';
import { getEntitySubgraph, getEntireKnowledgeGraph, type KnowledgeGraph, search } from '@/api/graph';
import { getEntitySubgraph, streamEntireKnowledgeGraph, type KnowledgeGraph, search } from '@/api/graph';
import { LinkDetails } from '@/components/graph/components/LinkDetails';
import { NetworkViewer, type NetworkViewerDetailsProps } from '@/components/graph/components/NetworkViewer';
import { NodeDetails } from '@/components/graph/components/NodeDetails';
Expand Down Expand Up @@ -93,6 +93,9 @@ function SubgraphSelector ({ knowledgeBaseId, query, onQueryChange }: { knowledg
<Select value={type} onValueChange={type => {
setType(type);
setInput('');
if (type === 'entire-knowledge-graph') {
onQueryChange(`${type}:`);
}
}}>
<SelectTrigger className="w-max">
<SelectValue />
Expand All @@ -103,18 +106,23 @@ function SubgraphSelector ({ knowledgeBaseId, query, onQueryChange }: { knowledg
<SelectItem value="message-subgraph">Message Subgraph</SelectItem>
<SelectItem value="trace" disabled>Langfuse Trace ID (UUID)</SelectItem>
<SelectItem value="document" disabled>Document URI</SelectItem>
<SelectItem value="entire-knowledge-graph">Entire Knowledge Graph</SelectItem>
</SelectContent>
</Select>
<Input
className="flex-1"
value={input}
onChange={event => setInput(event.target.value)}
onKeyDown={event => {
if (isHotkey('Enter', event)) {
onQueryChange(`${type}:${input}`);
}
}}
/>
{type !== 'entire-knowledge-graph' && (
<>
<Input
className="flex-1"
value={input}
onChange={event => setInput(event.target.value)}
onKeyDown={event => {
if (isHotkey('Enter', event)) {
onQueryChange(`${type}:${input}`);
}
}}
/>
</>
)}
<Link className={buttonVariants({})} href={`/knowledge-bases/${knowledgeBaseId}/knowledge-graph-explorer/create-synopsis-entity`}>
Create Synopsis Entity
</Link>
Expand Down Expand Up @@ -164,6 +172,7 @@ function getFetchInfo (kbId: number, query: string | null): [string | false, ()

const param = parsedQuery[1];


switch (parsedQuery[0]) {
// case 'trace':
// return ['get', `/api/v1/traces/${parsedQuery[1]}/knowledge-graph-retrieval`];
Expand All @@ -175,6 +184,8 @@ function getFetchInfo (kbId: number, query: string | null): [string | false, ()
return [`api.knowledge-bases.${kbId}.graph.search?query=${param}`, () => search(kbId, { query: param })];
case 'message-subgraph':
return [`api.chats.get-message-subgraph?id=${param}`, () => getChatMessageSubgraph(parseInt(param))];
case 'entire-knowledge-graph':
return [`api.knowledge-bases.${kbId}.graph.entire-knowledge-graph`, () => streamEntireKnowledgeGraph(kbId)];
}

return [false, () => Promise.reject()];
Expand Down
Loading
Loading