From 706a0cb2b2872a2336f1b6459ba9a4adca23a447 Mon Sep 17 00:00:00 2001 From: Pavel Date: Mon, 28 Sep 2026 17:35:16 +0400 Subject: [PATCH] Big source revamp --- docsgpt/api/pat/rules.py | 4 + docsgpt/api/user/sources/chunks.py | 105 +- docsgpt/api/user/sources/routes.py | 182 ++- docsgpt/graphrag/store.py | 410 +++++- docsgpt/parser/remote/base.py | 150 +- docsgpt/parser/remote/crawler_loader.py | 31 +- docsgpt/parser/remote/crawler_markdown.py | 32 +- docsgpt/parser/remote/sitemap_loader.py | 18 +- docsgpt/parser/remote/web_loader.py | 20 +- docsgpt/vectorstore/base.py | 34 + docsgpt/vectorstore/faiss.py | 52 + docsgpt/vectorstore/mongodb.py | 64 +- docsgpt/vectorstore/pgvector.py | 54 +- docsgpt/vectorstore/qdrant.py | 44 + frontend/DESIGN.md | 249 +++- frontend/src/api/endpoints.ts | 13 + frontend/src/api/services/userService.ts | 9 + frontend/src/components/ArtifactSidebar.tsx | 3 +- frontend/src/components/Chunks.test.tsx | 857 ++++++++++-- frontend/src/components/Chunks.tsx | 1231 +++++++++-------- .../src/components/ConnectorTree.test.tsx | 147 ++ frontend/src/components/ConnectorTree.tsx | 56 +- frontend/src/components/FileTree.tsx | 22 +- frontend/src/components/GraphView.tsx | 843 +++++++---- frontend/src/components/MarkdownPreview.tsx | 3 +- .../src/components/MultiSelectPopover.tsx | 2 +- .../src/components/SourceMarkdown.test.tsx | 55 + frontend/src/components/SourceMarkdown.tsx | 128 ++ frontend/src/components/WikiViewer.test.tsx | 407 ++++++ frontend/src/components/WikiViewer.tsx | 554 +++++--- frontend/src/components/chunkUtils.test.ts | 75 +- frontend/src/components/chunkUtils.ts | 46 +- .../components/graph/GraphCanvasControls.tsx | 58 + .../components/graph/GraphChunkSheet.test.tsx | 214 +++ .../src/components/graph/GraphChunkSheet.tsx | 237 ++++ .../src/components/graph/GraphEntities.tsx | 350 +++++ .../components/graph/GraphEntitySearch.tsx | 191 +++ .../components/graph/GraphNodePanel.test.tsx | 459 ++++++ .../src/components/graph/GraphNodePanel.tsx | 419 ++++++ .../components/graph/GraphSourceView.test.tsx | 573 ++++++++ .../src/components/graph/GraphSourceView.tsx | 293 ++++ .../src/components/graph/GraphTypeDot.tsx | 66 + .../components/graph/graphCanvasUtils.test.ts | 270 ++++ .../src/components/graph/graphCanvasUtils.ts | 281 ++++ .../components/graph/useGraphNodeDetail.ts | 69 + .../src/components/graphViewUtils.test.ts | 242 +++- frontend/src/components/graphViewUtils.ts | 247 +++- .../src/components/tree/PathHeader.test.tsx | 59 +- frontend/src/components/tree/PathHeader.tsx | 143 +- frontend/src/components/tree/ReaderPanel.tsx | 45 + .../components/tree/SourceEditSheet.test.tsx | 112 ++ .../src/components/tree/SourceEditSheet.tsx | 183 +++ .../components/tree/SourceNavigator.test.tsx | 175 +++ .../src/components/tree/SourceNavigator.tsx | 277 ++++ .../src/components/tree/TreeBrowser.test.tsx | 388 +++++- frontend/src/components/tree/TreeBrowser.tsx | 750 +++++----- .../components/tree/navigatorUtils.test.ts | 109 ++ .../src/components/tree/navigatorUtils.ts | 174 +++ frontend/src/components/tree/types.ts | 6 - frontend/src/components/ui/breadcrumb.tsx | 7 +- frontend/src/components/ui/command.test.tsx | 44 + frontend/src/components/ui/command.tsx | 32 +- frontend/src/components/ui/form-field.tsx | 9 +- frontend/src/components/ui/input.tsx | 12 +- frontend/src/components/ui/list-row.test.tsx | 11 + frontend/src/components/ui/list-row.tsx | 13 +- frontend/src/components/ui/multi-select.tsx | 1 - .../src/components/ui/pagination.test.tsx | 16 + frontend/src/components/ui/pagination.tsx | 10 +- frontend/src/components/ui/table.test.tsx | 21 + frontend/src/components/ui/table.tsx | 10 +- .../src/components/wikiViewerUtils.test.ts | 65 +- frontend/src/components/wikiViewerUtils.ts | 77 ++ frontend/src/conversation/MarkdownAnswer.tsx | 33 +- frontend/src/design/DesignSystem.tsx | 90 +- frontend/src/lib/markdown.test.tsx | 55 +- frontend/src/lib/markdown.tsx | 52 + frontend/src/locale/de.json | 110 +- frontend/src/locale/en.json | 110 +- frontend/src/locale/es.json | 110 +- frontend/src/locale/jp.json | 102 +- frontend/src/locale/ru.json | 126 +- frontend/src/locale/zh-TW.json | 102 +- frontend/src/locale/zh.json | 102 +- frontend/src/settings/Sources.tsx | 6 +- .../src/settings/components/UsageQuota.tsx | 6 +- frontend/src/utils/dateTimeUtils.test.ts | 24 +- frontend/src/utils/dateTimeUtils.ts | 18 + tests/api/user/sources/test_chunks.py | 274 +++- tests/api/user/sources/test_graph_view.py | 321 +++++ tests/api/user/sources/test_routes.py | 43 + tests/graphrag/test_store.py | 614 ++++++++ tests/parser/remote/test_base.py | 133 ++ tests/parser/remote/test_crawler_loader.py | 25 + tests/parser/remote/test_crawler_markdown.py | 25 + tests/parser/remote/test_sitemap_loader.py | 23 +- tests/parser/remote/test_web_loader.py | 47 +- tests/vectorstore/test_base.py | 52 + tests/vectorstore/test_faiss.py | 85 ++ tests/vectorstore/test_mongodb.py | 122 ++ tests/vectorstore/test_pgvector.py | 86 ++ .../vectorstore/test_pgvector_live_schema.py | 33 + tests/vectorstore/test_qdrant.py | 77 ++ 103 files changed, 13687 insertions(+), 1937 deletions(-) create mode 100644 frontend/src/components/ConnectorTree.test.tsx create mode 100644 frontend/src/components/SourceMarkdown.test.tsx create mode 100644 frontend/src/components/SourceMarkdown.tsx create mode 100644 frontend/src/components/WikiViewer.test.tsx create mode 100644 frontend/src/components/graph/GraphCanvasControls.tsx create mode 100644 frontend/src/components/graph/GraphChunkSheet.test.tsx create mode 100644 frontend/src/components/graph/GraphChunkSheet.tsx create mode 100644 frontend/src/components/graph/GraphEntities.tsx create mode 100644 frontend/src/components/graph/GraphEntitySearch.tsx create mode 100644 frontend/src/components/graph/GraphNodePanel.test.tsx create mode 100644 frontend/src/components/graph/GraphNodePanel.tsx create mode 100644 frontend/src/components/graph/GraphSourceView.test.tsx create mode 100644 frontend/src/components/graph/GraphSourceView.tsx create mode 100644 frontend/src/components/graph/GraphTypeDot.tsx create mode 100644 frontend/src/components/graph/graphCanvasUtils.test.ts create mode 100644 frontend/src/components/graph/graphCanvasUtils.ts create mode 100644 frontend/src/components/graph/useGraphNodeDetail.ts create mode 100644 frontend/src/components/tree/ReaderPanel.tsx create mode 100644 frontend/src/components/tree/SourceEditSheet.test.tsx create mode 100644 frontend/src/components/tree/SourceEditSheet.tsx create mode 100644 frontend/src/components/tree/SourceNavigator.test.tsx create mode 100644 frontend/src/components/tree/SourceNavigator.tsx create mode 100644 frontend/src/components/tree/navigatorUtils.test.ts create mode 100644 frontend/src/components/tree/navigatorUtils.ts create mode 100644 tests/parser/remote/test_base.py diff --git a/docsgpt/api/pat/rules.py b/docsgpt/api/pat/rules.py index c1d7683d..cea33e79 100644 --- a/docsgpt/api/pat/rules.py +++ b/docsgpt/api/pat/rules.py @@ -201,6 +201,9 @@ RULES: dict[tuple[str, str], Rule] = { ("/api/agents//schedules", "GET"): _rule( "schedules:read", refs=(("agents", (VIEW, "agent_id")),), blocked_by=_NON_AGENT_FAMILIES ), + ("/api/agents//schedules/stats", "GET"): _rule( + "schedules:read", refs=(("agents", (VIEW, "agent_id")),), blocked_by=_NON_AGENT_FAMILIES + ), ("/api/agents//schedules", "POST"): _rule( "schedules:write", refs=(("agents", (VIEW, "agent_id")),), blocked_by=_NON_AGENT_FAMILIES ), @@ -225,6 +228,7 @@ RULES: dict[tuple[str, str], Rule] = { ("/api/sources//graph/node/", "GET"): _rule( "sources:read", (VIEW, "source_id") ), + ("/api/sources//graph/nodes", "GET"): _rule("sources:read", (VIEW, "source_id")), # Ingestion and attachment extraction both report through this poll. ("/api/task_status", "GET"): _rule(any_of=("sources:read", "sources:write", "chat:run"), open=True), ("/api/upload", "POST"): _rule("sources:write"), diff --git a/docsgpt/api/user/sources/chunks.py b/docsgpt/api/user/sources/chunks.py index ae289aa6..240ed571 100644 --- a/docsgpt/api/user/sources/chunks.py +++ b/docsgpt/api/user/sources/chunks.py @@ -7,10 +7,11 @@ from flask_restx import fields, Namespace, Resource from docsgpt.api import api from docsgpt.api.user.base import get_vector_store -from docsgpt.api.user.team_sharing import effective_write_owner +from docsgpt.api.user.team_sharing import can_access, effective_write_owner from docsgpt.storage.db.repositories.sources import SourcesRepository from docsgpt.storage.db.session import db_readonly from docsgpt.utils import check_required_fields, num_tokens_from_string +from docsgpt.vectorstore.base import InvalidChunkMetadataError sources_chunks_ns = Namespace( "sources", description="Source document management operations", path="/api" @@ -18,12 +19,18 @@ sources_chunks_ns = Namespace( def _resolve_source(doc_id: str, user: str): - """Resolve a source (UUID or legacy ObjectId) for the caller. + """Resolve a source (UUID or legacy ObjectId) the caller may READ. - Returns the row dict (with PG UUID in ``id``) or ``None`` if missing. + Read access = owner or any team grant (viewer/editor). Returns the row + dict (with PG UUID in ``id``) or ``None`` if missing or not visible. """ with db_readonly() as conn: - return SourcesRepository(conn).get_any(doc_id, user) + doc = SourcesRepository(conn).get_any(doc_id, user) + if doc is not None: + return doc + if not can_access(conn, "source", doc_id, user): + return None + return SourcesRepository(conn).get_by_id(doc_id) def _resolve_source_for_write(doc_id: str, user: str): @@ -45,6 +52,35 @@ def _resolve_source_for_write(doc_id: str, user: str): return SourcesRepository(conn).get_any(doc_id, owner) +def _remap_graph_chunk(doc: dict, old_chunk_id: str, new_chunk_id: str) -> None: + """Move a graphrag source's links from an edited chunk's old id to its new one. + + Only needed when the store's ``update_chunk`` fell back to re-adding the + chunk under a new id (stores that update in place keep the id). Without + this the graph keeps pointing at the deleted row: the entity loses the + chunk and retrieval stops returning it. A failure is logged, not raised, + because the edit itself has already been saved. + + Args: + doc: The resolved source row. + old_chunk_id: The edited chunk's previous id. + new_chunk_id: The id the edit was saved under. + """ + from docsgpt.storage.db.source_config import SourceConfig + + if SourceConfig.parse(doc.get("config")).kind != "graphrag": + return + try: + from docsgpt.graphrag.store import GraphStore + + GraphStore().remap_chunk(str(doc["id"]), old_chunk_id, new_chunk_id) + except Exception as e: + current_app.logger.error( + f"Failed to remap graph links from chunk {old_chunk_id} to {new_chunk_id}: {e}", + exc_info=True, + ) + + def _has_usable_token_count(metadata: dict) -> bool: """Whether ``metadata`` already carries a count worth showing. @@ -98,6 +134,38 @@ def _with_token_counts(chunks: list) -> list: return chunks +def _path_ends_with(value: str, path: str) -> bool: + """Return whether ``value`` is ``path`` or ends with it at a ``/`` boundary. + + A bare ``endswith`` let a root ``setup.md`` also claim the chunks of + ``guides/setup.md`` (and ``a.md`` those of ``data.md``). + """ + return bool(value) and (value == path or value.endswith(f"/{path}")) + + +def _chunk_matches_path(metadata: dict, path: str) -> bool: + """Return whether a chunk belongs to the tree file at ``path``. + + Args: + metadata: The chunk's stored metadata. + path: The file's key path in the source's ``directory_structure``. + + Returns: + True when the chunk's ``source`` or ``file_path`` names that file, or + when the worker could only have keyed it by title: a remote ingest + (web page, Reddit post) whose chunks carry no ``file_path`` or ``key`` + (see ``remote_worker``). Sources ingested that way stay browsable + without a re-ingest. + """ + source = metadata.get("source") or "" + file_path = metadata.get("file_path") or "" + if _path_ends_with(source, path) or _path_ends_with(file_path, path): + return True + if "://" in source and not file_path and not metadata.get("key"): + return metadata.get("title") == path + return False + + @sources_chunks_ns.route("/get_chunks") class GetChunks(Resource): @api.doc( @@ -141,14 +209,8 @@ class GetChunks(Resource): for chunk in chunks: metadata = chunk.get("metadata", {}) - if path: - chunk_source = metadata.get("source", "") - chunk_file_path = metadata.get("file_path", "") - source_match = chunk_source and chunk_source.endswith(path) - file_path_match = chunk_file_path and chunk_file_path.endswith(path) - - if not (source_match or file_path_match): - continue + if path and not _chunk_matches_path(metadata, path): + continue if search_term: text_match = search_term in chunk.get("text", "").lower() title_match = search_term in metadata.get("title", "").lower() @@ -339,13 +401,11 @@ class UpdateChunk(Resource): if text is not None: new_metadata["token_count"] = num_tokens_from_string(new_text) try: - new_chunk_id = store.add_chunk(new_text, new_metadata) - - deleted = store.delete_chunk(chunk_id) - if not deleted: - current_app.logger.warning( - f"Failed to delete old chunk {chunk_id}, but new chunk {new_chunk_id} was created" - ) + # In place where the store supports it (same id, same list + # position); the base fallback re-adds under a new id. + new_chunk_id = store.update_chunk(chunk_id, new_text, new_metadata) + if new_chunk_id != chunk_id: + _remap_graph_chunk(doc, chunk_id, new_chunk_id) return make_response( jsonify( { @@ -356,8 +416,13 @@ class UpdateChunk(Resource): ), 200, ) + except InvalidChunkMetadataError as meta_error: + current_app.logger.warning( + f"Rejected metadata for chunk {chunk_id}: {meta_error}" + ) + return make_response(jsonify({"error": "Invalid metadata"}), 400) except Exception as add_error: - current_app.logger.error(f"Failed to add updated chunk: {add_error}") + current_app.logger.error(f"Failed to update chunk {chunk_id}: {add_error}") return make_response( jsonify({"error": "Failed to update chunk - addition failed"}), 500 ) diff --git a/docsgpt/api/user/sources/routes.py b/docsgpt/api/user/sources/routes.py index 6826c9e9..0930d689 100644 --- a/docsgpt/api/user/sources/routes.py +++ b/docsgpt/api/user/sources/routes.py @@ -545,7 +545,7 @@ class DirectoryStructure(Resource): return make_response(jsonify({"error": "Document ID is required"}), 400) try: with db_readonly() as conn: - doc = SourcesRepository(conn).get_any(doc_id, user) + doc = _resolve_readable_source(conn, doc_id, user) if not doc: return make_response( jsonify({"error": "Document not found or access denied"}), 404 @@ -1247,18 +1247,50 @@ class EnableSourceGraphRAG(Resource): ) -def _graph_overview_payload(source_id, limit): - """Return a bounded ``{nodes, edges}`` overview for a graphrag source. +def _graph_overview_payload(source_id: str, limit: int) -> dict: + """Return a bounded ``{nodes, edges, stats}`` overview for a graphrag source. An empty graph (extraction pending/capped/failed) yields empty lists rather - than an error, mirroring the ClassicRAG degradation guarantee. + than an error, mirroring the ClassicRAG degradation guarantee, and skips + the overview query entirely. ``stats`` carries the whole graph's totals, + not the bounded overview's. + + Args: + source_id: The resolved source id. + limit: Requested overview size; the store clamps it. + + Returns: + dict: ``{"nodes": [...], "edges": [...], "stats": {"nodes": int, + "edges": int}}``. """ from docsgpt.graphrag.store import GraphStore store = GraphStore() - if store.count_nodes(source_id) == 0: - return {"nodes": [], "edges": []} - return store.get_graph_overview(source_id, limit) + node_count = store.count_nodes(source_id) + if node_count == 0: + return {"nodes": [], "edges": [], "stats": {"nodes": 0, "edges": 0}} + overview = store.get_graph_overview(source_id, limit) + return { + "nodes": overview["nodes"], + "edges": overview["edges"], + "stats": {"nodes": node_count, "edges": store.count_edges(source_id)}, + } + + +def _int_arg(name: str, default: int) -> int: + """Read an integer query arg, falling back to ``default`` when unparsable. + + Args: + name: Query-string parameter name. + default: Value used when the arg is missing or not an integer. + + Returns: + int: The parsed value or ``default``. + """ + try: + return int(request.args.get(name, default)) + except (TypeError, ValueError): + return default @sources_ns.route("/sources//graph") @@ -1310,6 +1342,142 @@ class SourceGraph(Resource): "success": True, "nodes": overview["nodes"], "edges": overview["edges"], + "stats": overview["stats"], + } + ), + 200, + ) + + +def _graph_node_page( + source_id: str, + query: str | None, + type_key: str | None, + page: int, + per_page: int, +) -> tuple[dict, list]: + """Return ``(listing, type_facets)`` for one page of a source's graph nodes. + + A graph store that is not configured, or whose tables do not exist yet + (no graph has ever been built on this deployment), yields an empty page, + matching :func:`_graph_overview_payload`. Every other failure propagates, + so a broken query is never shown as an empty graph. + + Args: + source_id: The resolved source id. + query: Name substring filter, or ``None``. + type_key: Type key filter, or ``None`` for no filter. + page: 1-based page, already clamped. + per_page: Page size, already clamped. + + Returns: + tuple: ``({"nodes": [...], "total": int}, [{"key", "label", "count"}])``. + + Raises: + Exception: Any store failure other than a missing store or schema. + """ + import psycopg + + from docsgpt.graphrag.store import GraphStore + + empty = ({"nodes": [], "total": 0}, []) + try: + store = GraphStore() + except (ValueError, ImportError) as err: + current_app.logger.info("Graph store unavailable, listing no nodes: %s", err) + return empty + try: + listing = store.list_nodes( + source_id, + query=query, + type_key=type_key, + offset=(page - 1) * per_page, + limit=per_page, + ) + types = store.node_type_facets(source_id) + except psycopg.errors.UndefinedTable as err: + current_app.logger.info("Graph tables missing, listing no nodes: %s", err) + return empty + return listing, types + + +@sources_ns.route("/sources//graph/nodes") +class SourceGraphNodes(Resource): + @api.doc( + description="Paged, searchable list of a graphrag source's nodes, " + "highest degree first, plus type facets over all its nodes (read " + "access: owner or shared). Query params: q (name substring), type " + "(type key), page (1-based), per_page (1-100, default 25)." + ) + def get(self, source_id: str): + """List one page of a source's graph nodes with type facets. + + Args: + source_id: The source id from the URL. + + Returns: + Response: ``{success, nodes, total, page, per_page, types}``, or + 401/404/400 on missing auth, no read access or a failure. ``page`` + is clamped to ``1..GRAPH_NODE_LIST_MAX_PAGE`` and ``per_page`` to + ``1..GRAPH_NODE_LIST_MAX_LIMIT``. A graph store that is not + configured or has no tables yet reads as an empty page, like the + overview; any other store query failure is a 400, never an empty + 200. + """ + decoded_token = request.decoded_token + if not decoded_token: + return make_response(jsonify({"success": False}), 401) + user = decoded_token.get("sub") + from docsgpt.graphrag.store import ( + GRAPH_NODE_LIST_DEFAULT_LIMIT, + GRAPH_NODE_LIST_MAX_LIMIT, + GRAPH_NODE_LIST_MAX_PAGE, + ) + + page = max(1, min(_int_arg("page", 1), GRAPH_NODE_LIST_MAX_PAGE)) + per_page = max( + 1, + min( + _int_arg("per_page", GRAPH_NODE_LIST_DEFAULT_LIMIT), + GRAPH_NODE_LIST_MAX_LIMIT, + ), + ) + query = (request.args.get("q") or "").strip() or None + type_key = request.args.get("type") + try: + with db_readonly() as conn: + doc = _resolve_readable_source(conn, source_id, user) + if doc is None: + return make_response( + jsonify({"success": False, "message": "Source not found"}), + 404, + ) + resolved_source_id = str(doc["id"]) + except Exception as err: + current_app.logger.error( + f"Error resolving source {source_id} for graph nodes: {err}", + exc_info=True, + ) + return make_response(jsonify({"success": False}), 400) + try: + listing, types = _graph_node_page( + resolved_source_id, query, type_key, page, per_page + ) + except Exception as err: + current_app.logger.error( + f"Error listing graph nodes for {source_id}: {err}", + exc_info=True, + ) + return make_response(jsonify({"success": False}), 400) + return make_response( + jsonify( + { + "success": True, + "nodes": listing["nodes"], + "total": listing["total"], + "page": page, + "per_page": per_page, + "types": types, } ), 200, diff --git a/docsgpt/graphrag/store.py b/docsgpt/graphrag/store.py index 0b73d031..d19c94d1 100644 --- a/docsgpt/graphrag/store.py +++ b/docsgpt/graphrag/store.py @@ -39,8 +39,49 @@ MAX_SUBGRAPH_EDGES = 2000 GRAPH_OVERVIEW_DEFAULT_LIMIT = 100 GRAPH_OVERVIEW_MAX_LIMIT = 250 +GRAPH_NODE_LIST_DEFAULT_LIMIT = 25 +GRAPH_NODE_LIST_MAX_LIMIT = 100 +# Highest 1-based page the node-list route accepts. Keeps ``OFFSET`` (at most +# ~1e8 rows at the max page size) far inside Postgres' bigint range, so a huge +# ``page`` query arg yields an empty page with the real total, not an error. +GRAPH_NODE_LIST_MAX_PAGE = 1_000_000 + +MAX_NODE_RELATIONSHIPS = 300 + PGVECTOR_SOURCE_COLUMN = "source_id" +def graph_type_key(type_value: Optional[str]) -> str: + """Return the grouping key for a node type. + + ``"Person"``, ``"PERSON"`` and ``"per son"`` all fold to ``"person"``, so + the extractor's inconsistent spellings land in one facet. Computed in + Python on purpose, never in SQL: ``lower()`` and ``[:alnum:]`` follow the + database's ``LC_CTYPE``, and on a ``C``-locale database (common on managed + Postgres) they fold every Cyrillic or CJK type to ``""``. Must agree with + the frontend, which filters and colours by this key with + ``/[^\\p{L}\\p{N}]/gu`` and ``toLowerCase``. + + Args: + type_value: The raw node type, possibly ``None`` or empty. + + Returns: + str: The lower-cased type with non-alphanumerics removed; ``""`` for a + missing type. + """ + return "".join(ch for ch in (type_value or "").lower() if ch.isalnum()) + + +def _escape_like(value: str) -> str: + """Escape ``value`` so ``LIKE``/``ILIKE`` treats it as a literal. + + Args: + value: Raw user text. + + Returns: + str: ``value`` with ``\\``, ``%`` and ``_`` backslash-escaped. + """ + return value.replace("\\", "\\\\").replace("%", "\\%").replace("_", "\\_") + def _safe_identifier(name: str) -> str: """Return ``name`` if it is a bare SQL identifier, else raise. @@ -857,6 +898,174 @@ class GraphStore: cursor.close() conn.rollback() + def count_edges(self, source_id: str) -> int: + """Number of edges for a source. + + Args: + source_id: Source whose edges to count. + + Returns: + int: The edge count, or ``0`` when the query fails. + """ + conn = self._get_connection() + cursor = conn.cursor() + try: + cursor.execute( + "SELECT count(*) FROM graph_edges WHERE source_id = %s;", + (source_id,), + ) + return int(cursor.fetchone()[0]) + except Exception as e: + logging.error(f"Error counting edges: {e}") + return 0 + finally: + cursor.close() + conn.rollback() + + def list_nodes( + self, + source_id: str, + query: Optional[str] = None, + type_key: Optional[str] = None, + offset: int = 0, + limit: int = GRAPH_NODE_LIST_DEFAULT_LIMIT, + ) -> Dict[str, Any]: + """One page of a source's nodes, highest degree first, with the match count. + + Args: + source_id: Source whose nodes to list. + query: Case-insensitive substring of the node name; ``%``, ``_`` + and ``\\`` match literally. Blank means no name filter. + type_key: Only nodes whose :func:`graph_type_key` equals this; + ``""`` selects untyped nodes and ``None`` applies no filter. + offset: Rows to skip (floored at 0). + limit: Page size, clamped to ``1..GRAPH_NODE_LIST_MAX_LIMIT``. + + Returns: + dict: ``{"nodes": [{id, name, type, degree, doc_freq}], "total": int}`` + where ``total`` counts every match, not just this page. A source + with no graph yields ``{"nodes": [], "total": 0}``. + + Raises: + Exception: Any query failure, so a caller never mistakes a broken + query for an empty graph. + """ + limit = max(1, min(int(limit), GRAPH_NODE_LIST_MAX_LIMIT)) + offset = max(0, int(offset)) + clauses = ["source_id = %s"] + params: List[Any] = [source_id] + clean = (query or "").strip() + if clean: + clauses.append("name ILIKE %s ESCAPE '\\'") + params.append(f"%{_escape_like(clean)}%") + conn = self._get_connection() + cursor = conn.cursor() + try: + if type_key is not None: + # The key is folded in Python (see ``graph_type_key``), so the + # filter matches the raw spellings whose key equals it. + cursor.execute( + "SELECT DISTINCT type FROM graph_nodes WHERE source_id = %s;", + (source_id,), + ) + raw_types = sorted( + row[0] + for row in cursor.fetchall() + if row[0] is not None and graph_type_key(row[0]) == type_key + ) + if type_key == "": + clauses.append("(type IS NULL OR type = ANY(%s))") + elif not raw_types: + return {"nodes": [], "total": 0} + else: + clauses.append("type = ANY(%s)") + params.append(raw_types) + where = " AND ".join(clauses) + cursor.execute( + f"SELECT count(*) FROM graph_nodes WHERE {where};", tuple(params) + ) + total = int(cursor.fetchone()[0]) + cursor.execute( + f""" + SELECT id, name, type, degree, doc_freq + FROM graph_nodes + WHERE {where} + ORDER BY degree DESC, id + LIMIT %s OFFSET %s; + """, + (*params, limit, offset), + ) + nodes = [ + { + "id": str(row[0]), + "name": row[1], + "type": row[2], + "degree": row[3], + "doc_freq": row[4], + } + for row in cursor.fetchall() + ] + return {"nodes": nodes, "total": total} + finally: + cursor.close() + conn.rollback() + + def node_type_facets(self, source_id: str) -> List[Dict[str, Any]]: + """Node counts per type key over every node of a source. + + Spellings that share a :func:`graph_type_key` are folded into one + facet labelled with the most frequent spelling (ties go to the + alphabetically first). Untyped nodes form the ``""`` facet with a + ``None`` label. + + Args: + source_id: Source whose node types to count. + + Returns: + list: ``[{"key", "label", "count"}]`` sorted by count descending, + then label. Empty for a source with no graph. + + Raises: + Exception: Any query failure, so a caller never mistakes a broken + query for an empty graph. + """ + conn = self._get_connection() + cursor = conn.cursor() + try: + # Group by the raw spelling only; the key is folded in Python so it + # does not depend on the database locale (see ``graph_type_key``). + cursor.execute( + """ + SELECT type, count(*) + FROM graph_nodes + WHERE source_id = %s + GROUP BY type; + """, + (source_id,), + ) + rows = cursor.fetchall() + finally: + cursor.close() + conn.rollback() + spellings: Dict[str, List[tuple]] = {} + for raw_type, n in rows: + spellings.setdefault(graph_type_key(raw_type), []).append( + (int(n), raw_type) + ) + facets = [] + for key, variants in spellings.items(): + label = None + if key: + # Most frequent spelling; a tie goes to the alphabetically first. + label = min(variants, key=lambda v: (-v[0], v[1]))[1] + facets.append( + {"key": key, "label": label, "count": sum(v[0] for v in variants)} + ) + facets.sort( + key=lambda f: (-f["count"], f["label"] is None, f["label"] or "") + ) + return facets + def count_nodes_many(self, source_ids: List[str]) -> Dict[str, int]: """Node counts for several sources in one round trip. @@ -1390,11 +1599,26 @@ class GraphStore: def get_node_detail( self, source_id: str, node_id: str, max_chunks: int = 20 ) -> Optional[Dict[str, Any]]: - """A node's full record plus a bounded list of its linked chunks. + """A node's full record, its relationships and a bounded list of its chunks. - Returns ``None`` when the node does not belong to the source. Chunk texts - are read from the co-located pgvector table; at most ``max_chunks`` are - returned so a hub node never streams an unbounded payload. + Chunk texts are read from the co-located pgvector table; at most + ``max_chunks`` are returned so a hub node never streams an unbounded + payload. ``relationships`` holds every edge touching the node in either + direction (self-loops skipped), joined to the other endpoint, strongest + neighbour first and capped at ``MAX_NODE_RELATIONSHIPS``; + ``relationships_total`` is how many such edges there are, so a capped + list can say what it leaves out. + + Args: + source_id: Source the node must belong to. + node_id: The node's id. + max_chunks: Most linked chunks to return. + + Returns: + dict | None: ``{id, name, type, description, degree, doc_freq, + relationships: [{id, name, type, degree, edge_type, direction}], + relationships_total, chunks: [{chunk_id, text, metadata}]}``; ``None`` when the node does + not belong to the source or the read fails. """ conn = self._get_connection() try: @@ -1421,6 +1645,15 @@ class GraphStore: "degree": row[4], "doc_freq": row[5], } + relationships = self._node_relationships(conn, source_id, node_id) + node["relationships"] = relationships + node["relationships_total"] = ( + self._node_relationships_total( + conn, source_id, node_id, len(relationships) + ) + if len(relationships) >= MAX_NODE_RELATIONSHIPS + else len(relationships) + ) chunk_ids = self.get_chunk_ids_for_nodes(source_id, [node_id]).get( str(node_id), [] @@ -1443,6 +1676,107 @@ class GraphStore: finally: conn.rollback() + def _node_relationships( + self, conn, source_id: str, node_id: str + ) -> List[Dict[str, Any]]: + """Edges touching ``node_id``, joined to the neighbour at the other end. + + Args: + conn: Open connection; the caller owns the transaction. + source_id: Source the edges and neighbours must belong to. + node_id: The node whose relationships to read. + + Returns: + list: ``[{id, name, type, degree, edge_type, direction}]`` where + ``direction`` is ``"out"`` when the node is the edge's source, + ordered by neighbour degree descending, neighbour id, edge type. + Empty when the read fails (the transaction is rolled back). + """ + cursor = conn.cursor() + try: + cursor.execute( + """ + SELECT o.id, o.name, o.type, o.degree, e.type, + CASE WHEN e.src_node_id = %s THEN 'out' ELSE 'in' END AS direction + FROM graph_edges e + JOIN graph_nodes o + ON o.id = CASE WHEN e.src_node_id = %s + THEN e.dst_node_id ELSE e.src_node_id END + WHERE e.source_id = %s + AND (e.src_node_id = %s OR e.dst_node_id = %s) + AND e.src_node_id <> e.dst_node_id + AND o.source_id = %s + ORDER BY o.degree DESC, o.id, e.type + LIMIT %s; + """, + ( + node_id, node_id, source_id, node_id, node_id, source_id, + MAX_NODE_RELATIONSHIPS, + ), + ) + return [ + { + "id": str(row[0]), + "name": row[1], + "type": row[2], + "degree": row[3], + "edge_type": row[4], + "direction": row[5], + } + for row in cursor.fetchall() + ] + except Exception as e: + # Degrade to no relationships rather than failing the whole node + # detail; roll back so the caller's chunk read can reuse ``conn``. + logging.error(f"Error reading node relationships: {e}") + conn.rollback() + return [] + finally: + cursor.close() + + def _node_relationships_total( + self, conn, source_id: str, node_id: str, fallback: int + ) -> int: + """How many edges ``_node_relationships`` would return without its cap. + + Same filters as the list: either direction, self-loops skipped, the + neighbour in the same source. + + Args: + conn: Open connection; the caller owns the transaction. + source_id: Source the edges and neighbours must belong to. + node_id: The node whose relationships to count. + fallback: What to return when the count fails (the list's length). + + Returns: + int: The edge count, or ``fallback`` when the read fails (the + transaction is rolled back). + """ + cursor = conn.cursor() + try: + cursor.execute( + """ + SELECT COUNT(*) + FROM graph_edges e + JOIN graph_nodes o + ON o.id = CASE WHEN e.src_node_id = %s + THEN e.dst_node_id ELSE e.src_node_id END + WHERE e.source_id = %s + AND (e.src_node_id = %s OR e.dst_node_id = %s) + AND e.src_node_id <> e.dst_node_id + AND o.source_id = %s; + """, + (node_id, source_id, node_id, node_id, source_id), + ) + row = cursor.fetchone() + return int(row[0]) if row else fallback + except Exception as e: + logging.error(f"Error counting node relationships: {e}") + conn.rollback() + return fallback + finally: + cursor.close() + def set_node_degrees(self, source_id: str): """Recompute every node's degree from its incident edges for a source. @@ -1506,6 +1840,74 @@ class GraphStore: return self._write_with_reconnect(_write) + def remap_chunk(self, source_id: str, old_chunk_id: str, new_chunk_id: str) -> None: + """Point every graph reference to a chunk at its replacement id. + + A chunk edit re-adds the chunk under a new id and deletes the old row, + so without this the node links, edge provenance and extraction + checkpoint would all name a chunk that no longer exists. The graph + itself is not re-extracted from the edited text. + + Args: + source_id: Source the chunk belongs to. + old_chunk_id: The id the graph currently references. + new_chunk_id: The id of the chunk that replaced it. + + Raises: + Exception: The underlying write failure, after a rollback. + """ + self._ensure_tables_once() + old_id, new_id = str(old_chunk_id), str(new_chunk_id) + + def _write(conn): + cursor = conn.cursor() + try: + _lock_source(cursor, source_id) + cursor.execute( + """ + UPDATE graph_node_chunks SET chunk_id = %s + WHERE source_id = %s AND chunk_id = %s; + """, + (new_id, source_id, old_id), + ) + # Rewrite the matching element in place, keeping array order. + cursor.execute( + """ + UPDATE graph_edges + SET source_chunk_ids = ( + SELECT jsonb_agg( + CASE WHEN t.v #>> '{}' = %s THEN to_jsonb(%s::text) ELSE t.v END + ORDER BY t.ord + ) + FROM jsonb_array_elements(source_chunk_ids) + WITH ORDINALITY AS t(v, ord) + ) + WHERE source_id = %s + AND jsonb_typeof(source_chunk_ids) = 'array' + AND EXISTS ( + SELECT 1 FROM jsonb_array_elements(source_chunk_ids) AS s(v) + WHERE s.v #>> '{}' = %s + ); + """, + (old_id, new_id, source_id, old_id), + ) + cursor.execute( + """ + UPDATE graph_ingest_progress SET chunk_id = %s + WHERE source_id = %s AND chunk_id = %s; + """, + (new_id, source_id, old_id), + ) + conn.commit() + except Exception as e: + _safe_rollback(conn) + logging.error(f"Error remapping graph chunk: {e}") + raise + finally: + cursor.close() + + return self._write_with_reconnect(_write) + def pending_chunks(self, source_id: str, all_chunk_ids: List[str]) -> List[str]: """Chunk ids from ``all_chunk_ids`` not yet marked ``done`` for the source.""" if not all_chunk_ids: diff --git a/docsgpt/parser/remote/base.py b/docsgpt/parser/remote/base.py index 5030ccc9..2cbfaf37 100644 --- a/docsgpt/parser/remote/base.py +++ b/docsgpt/parser/remote/base.py @@ -1,10 +1,158 @@ """Base reader class.""" +import hashlib +import os +import re from abc import abstractmethod -from typing import Any, List +from typing import Any, Dict, List +from urllib.parse import parse_qsl, urldefrag, urlencode, urlparse from docsgpt.parser.schema.base import Document from docsgpt.vectorstore.document_class import Document as VectorDocument +_PAGE_EXTENSIONS = (".html", ".htm", ".php", ".asp", ".aspx", ".jsp") + +# Joins a page's path to its encoded query string in the virtual file name. +_QUERY_SEPARATOR = "__" +# Longest encoded query kept verbatim; longer ones keep a prefix plus a hash. +MAX_QUERY_SEGMENT_LENGTH = 64 +_QUERY_HASH_LENGTH = 10 +# Anything but letters, digits, ``_``, ``.`` and ``-`` in a query key or value. +_UNSAFE_QUERY_CHARS = re.compile(r"[^\w.-]+") + + +def _query_segment(query: str) -> str: + """Encode a URL query string as a deterministic, tree-safe name segment. + + Parameters are sorted so the same page maps to the same name across + re-syncs whatever order its links list them in. Keys and values are + percent-decoded and every run of characters outside ``[\\w.-]`` becomes + ``-``, so the segment carries no ``/``, ``?``, ``#`` or ``%``. A segment + longer than ``MAX_QUERY_SEGMENT_LENGTH`` keeps a prefix and a hash of the + full query, which keeps it bounded and still distinct. + + Args: + query: The raw query string, without the leading ``?``. + + Returns: + The encoded segment, or ``""`` when the query has no parameters. + """ + pairs = sorted(parse_qsl(query, keep_blank_values=True)) + if not pairs: + return "" + parts = [] + for key, value in pairs: + key = _UNSAFE_QUERY_CHARS.sub("-", key) + value = _UNSAFE_QUERY_CHARS.sub("-", value) + parts.append(f"{key}={value}" if value else key) + segment = "&".join(parts) + if len(segment) <= MAX_QUERY_SEGMENT_LENGTH: + return segment + digest = hashlib.sha1(urlencode(pairs).encode("utf-8")).hexdigest() + keep = MAX_QUERY_SEGMENT_LENGTH - _QUERY_HASH_LENGTH - 1 + return f"{segment[:keep]}-{digest[:_QUERY_HASH_LENGTH]}" + + +def url_to_virtual_path(url: str, include_host: bool = False) -> str: + """Convert a page URL to the virtual ``.md`` path used as its tree key. + + Web ingests store this as the chunk's ``file_path``: the worker keys the + source's file tree by it and the chunks view filters by it, so the two + agree on which chunks belong to which file. + + The fragment is dropped (it names a spot on the same page), but the query + string is kept, encoded into the file name (see ``_query_segment``), so + ``/p?page=1`` and ``/p?page=2`` stay two files. A URL without a query maps + exactly as before. Paths that still collide, such as ``/a`` and + ``/a.html``, are told apart by ``dedupe_virtual_paths``. + + Args: + url: Page URL, e.g. ``"https://docs.docsgpt.cloud/guides/setup"``. + include_host: Prefix the host, for ingests spanning several hosts + whose paths would otherwise collide (every root is ``index.md``). + + Returns: + A relative path such as ``"index.md"``, ``"guides/setup.md"``, + ``"list__page=2.md"`` or, with ``include_host``, + ``"docs.docsgpt.cloud/guides/setup.md"``. + """ + parsed = urlparse(url) + path = parsed.path.strip("/") + + if not path: + path = "index.md" + else: + base, ext = os.path.splitext(path) + if ext.lower() in _PAGE_EXTENSIONS: + path = base + if not path.endswith(".md"): + path = f"{path}.md" + + query = _query_segment(parsed.query) + if query: + path = f"{path[:-len('.md')]}{_QUERY_SEPARATOR}{query}.md" + + if include_host and parsed.netloc: + return f"{parsed.netloc}/{path}" + return path + + +def spans_multiple_hosts(urls: List[str]) -> bool: + """Return whether ``urls`` point at more than one host. + + Args: + urls: Page URLs about to be ingested together. + + Returns: + True when paths alone could collide and need a host prefix. + """ + return len({urlparse(u).netloc for u in urls}) > 1 + + +def dedupe_virtual_paths(documents: List[Document]) -> List[Document]: + """Give distinct pages that share a virtual ``file_path`` distinct ones. + + ``url_to_virtual_path`` folds some different URLs onto one path (``/a`` and + ``/a.html`` are both ``a.md``), and the worker would then merge those pages + into one tree entry. Within a path, the pages are ordered by URL (fragment + dropped): the first keeps the path and each later one gets ``-2``, ``-3`` + and so on before ``.md``, skipping any path another page already has. + Ordering by URL rather than by fetch order keeps the names stable across + re-syncs of the same pages. Documents for the same URL (one page reached + via two fragments) are one page and keep one path. + + Args: + documents: Loaded documents; each ``extra_info`` carries ``source`` + (the page URL) and ``file_path``. Others are left untouched. + + Returns: + ``documents``, with colliding ``file_path`` values rewritten in place. + """ + pages_by_path: Dict[str, Dict[str, List[Document]]] = {} + for doc in documents: + info = doc.extra_info or {} + path = info.get("file_path") + if not path: + continue + page = urldefrag(str(info.get("source") or path))[0] + pages_by_path.setdefault(path, {}).setdefault(page, []).append(doc) + + taken = set(pages_by_path) + for path in sorted(pages_by_path): + pages = pages_by_path[path] + if len(pages) < 2: + continue + stem = path[: -len(".md")] if path.endswith(".md") else path + suffix = 1 + for page in sorted(pages)[1:]: + suffix += 1 + while f"{stem}-{suffix}.md" in taken: + suffix += 1 + new_path = f"{stem}-{suffix}.md" + taken.add(new_path) + for doc in pages[page]: + doc.extra_info["file_path"] = new_path + return documents + class BaseRemote: """Utilities for loading data from a directory.""" diff --git a/docsgpt/parser/remote/crawler_loader.py b/docsgpt/parser/remote/crawler_loader.py index f75f054a..e6d2b8d5 100644 --- a/docsgpt/parser/remote/crawler_loader.py +++ b/docsgpt/parser/remote/crawler_loader.py @@ -1,9 +1,8 @@ import logging -import os from bs4 import BeautifulSoup from urllib.parse import urljoin, urlparse -from docsgpt.parser.remote.base import BaseRemote +from docsgpt.parser.remote.base import BaseRemote, dedupe_virtual_paths, url_to_virtual_path from docsgpt.parser.schema.base import Document from docsgpt.core.url_validation import validate_url, SSRFError from docsgpt.security.safe_url import pinned_request @@ -73,30 +72,8 @@ class CrawlerLoader(BaseRemote): if self.limit is not None and len(visited_urls) >= self.limit: break - return loaded_content + return dedupe_virtual_paths(loaded_content) def _url_to_virtual_path(self, url): - """ - Convert a URL to a virtual file path ending with .md. - - Examples: - https://docs.docsgpt.cloud/ -> index.md - https://docs.docsgpt.cloud/guides/setup -> guides/setup.md - https://docs.docsgpt.cloud/guides/setup/ -> guides/setup.md - https://example.com/page.html -> page.md - """ - parsed = urlparse(url) - path = parsed.path.strip("/") - - if not path: - return "index.md" - - # Remove common file extensions and add .md - base, ext = os.path.splitext(path) - if ext.lower() in [".html", ".htm", ".php", ".asp", ".aspx", ".jsp"]: - path = base - - if not path.endswith(".md"): - path = f"{path}.md" - - return path + """Convert a URL to a virtual ``.md`` path; see ``url_to_virtual_path``.""" + return url_to_virtual_path(url) diff --git a/docsgpt/parser/remote/crawler_markdown.py b/docsgpt/parser/remote/crawler_markdown.py index b2de5cbb..db8693cc 100644 --- a/docsgpt/parser/remote/crawler_markdown.py +++ b/docsgpt/parser/remote/crawler_markdown.py @@ -1,13 +1,12 @@ from urllib.parse import urlparse, urljoin from bs4 import BeautifulSoup -from docsgpt.parser.remote.base import BaseRemote +from docsgpt.parser.remote.base import BaseRemote, dedupe_virtual_paths, url_to_virtual_path from docsgpt.core.url_validation import validate_url, SSRFError from docsgpt.security.safe_url import UnsafeUserUrlError, pinned_request import re from markdownify import markdownify from docsgpt.parser.schema.base import Document import tldextract -import os # The bundled public-suffix snapshot is enough for domain matching; the # default extractor would fetch the live list on first use and cache it on @@ -91,7 +90,7 @@ class CrawlerLoader(BaseRemote): if self.limit is not None and len(visited_urls) >= self.limit: break - return documents + return dedupe_virtual_paths(documents) def _fetch_page(self, url): try: @@ -160,28 +159,5 @@ class CrawlerLoader(BaseRemote): return filtered def _url_to_virtual_path(self, url): - """ - Convert a URL to a virtual file path ending with .md. - - Examples: - https://docs.docsgpt.cloud/ -> index.md - https://docs.docsgpt.cloud/guides/setup -> guides/setup.md - https://docs.docsgpt.cloud/guides/setup/ -> guides/setup.md - https://example.com/page.html -> page.md - """ - parsed = urlparse(url) - path = parsed.path.strip("/") - - if not path: - return "index.md" - - # Remove common file extensions and add .md - base, ext = os.path.splitext(path) - if ext.lower() in [".html", ".htm", ".php", ".asp", ".aspx", ".jsp"]: - path = base - - # Ensure path ends with .md - if not path.endswith(".md"): - path = path + ".md" - - return path \ No newline at end of file + """Convert a URL to a virtual ``.md`` path; see ``url_to_virtual_path``.""" + return url_to_virtual_path(url) diff --git a/docsgpt/parser/remote/sitemap_loader.py b/docsgpt/parser/remote/sitemap_loader.py index a8500c79..fa686b28 100644 --- a/docsgpt/parser/remote/sitemap_loader.py +++ b/docsgpt/parser/remote/sitemap_loader.py @@ -4,7 +4,12 @@ import re import defusedxml.ElementTree as ET from bs4 import BeautifulSoup -from docsgpt.parser.remote.base import BaseRemote +from docsgpt.parser.remote.base import ( + BaseRemote, + dedupe_virtual_paths, + spans_multiple_hosts, + url_to_virtual_path, +) from docsgpt.parser.schema.base import Document from docsgpt.core.url_validation import validate_url, SSRFError from docsgpt.security.safe_url import UnsafeUserUrlError, pinned_request @@ -32,6 +37,7 @@ class SitemapLoader(BaseRemote): return [] # Load content of extracted URLs + include_host = spans_multiple_hosts(urls) documents = [] processed_urls = 0 # Counter for processed URLs for url in urls: @@ -50,7 +56,13 @@ class SitemapLoader(BaseRemote): documents.append( Document( soup.get_text(separator="\n", strip=True), - extra_info={"source": url}, + # Without file_path the worker had no tree key (no + # title, key or doc_id), so sitemap pages never + # appeared in the file tree. + extra_info={ + "source": url, + "file_path": url_to_virtual_path(url, include_host), + }, ) ) processed_urls += 1 # Increment the counter after processing each URL @@ -58,7 +70,7 @@ class SitemapLoader(BaseRemote): logging.error(f"Error processing URL {url}: {e}", exc_info=True) continue - return documents + return dedupe_virtual_paths(documents) def _extract_urls(self, sitemap_url): try: diff --git a/docsgpt/parser/remote/web_loader.py b/docsgpt/parser/remote/web_loader.py index a5a19e49..1eacbc8b 100644 --- a/docsgpt/parser/remote/web_loader.py +++ b/docsgpt/parser/remote/web_loader.py @@ -3,7 +3,12 @@ import logging from bs4 import BeautifulSoup from docsgpt.core.url_validation import SSRFError, validate_url -from docsgpt.parser.remote.base import BaseRemote +from docsgpt.parser.remote.base import ( + BaseRemote, + dedupe_virtual_paths, + spans_multiple_hosts, + url_to_virtual_path, +) from docsgpt.parser.schema.base import Document from docsgpt.security.safe_url import pinned_request @@ -24,15 +29,17 @@ class WebLoader(BaseRemote): urls = inputs if isinstance(urls, str): urls = [urls] - documents = [] + valid_urls = [] for url in urls: try: - url = validate_url(url) + valid_urls.append(validate_url(url)) except SSRFError as e: logging.warning( f"Skipping URL due to SSRF validation failure: {url} - {e}" ) - continue + include_host = spans_multiple_hosts(valid_urls) + documents = [] + for url in valid_urls: try: response = pinned_request("GET", url, headers=headers, timeout=30) response.raise_for_status() @@ -45,6 +52,9 @@ class WebLoader(BaseRemote): html_tag = soup.find("html") if html_tag and html_tag.get("lang"): metadata["language"] = html_tag.get("lang") + # The worker keys the file tree by file_path; without it the + # tree fell back to the title, which the chunks view can't match. + metadata["file_path"] = url_to_virtual_path(url, include_host) documents.append( Document( soup.get_text(separator="\n", strip=True), @@ -54,4 +64,4 @@ class WebLoader(BaseRemote): except Exception as e: logging.error(f"Error processing URL {url}: {e}", exc_info=True) continue - return documents + return dedupe_virtual_paths(documents) diff --git a/docsgpt/vectorstore/base.py b/docsgpt/vectorstore/base.py index d0c63c04..348f99e5 100644 --- a/docsgpt/vectorstore/base.py +++ b/docsgpt/vectorstore/base.py @@ -362,6 +362,14 @@ def build_local_embeddings( return embedding_instance +class InvalidChunkMetadataError(ValueError): + """Chunk metadata the store cannot write, such as a key it reserves. + + A client-input error, distinct from a store or embedding failure: the + chunk routes answer it with a 400 rather than a 500. + """ + + class BaseVectorStore(ABC): def __init__(self): pass @@ -431,6 +439,32 @@ class BaseVectorStore(ABC): """Delete a specific chunk from the vectorstore""" pass + def update_chunk(self, chunk_id: str, text: str, metadata: dict) -> str: + """Replace a chunk's text and metadata, returning the id it is now under. + + Stores that can rewrite a row in place override this and keep both the + id and the chunk's position in :meth:`get_chunks`. This default works + for any store but re-adds the chunk and deletes the old one, so the + returned id differs and the chunk moves; callers holding the old id + (graph links, for one) must follow the returned id. + + Args: + chunk_id: Id of the chunk to replace. + text: The chunk's new text. + metadata: The chunk's complete new metadata. + + Returns: + The id the updated chunk is stored under. + """ + new_chunk_id = self.add_chunk(text, metadata) + if not self.delete_chunk(chunk_id): + logging.warning( + "Failed to delete old chunk %s, but new chunk %s was created", + chunk_id, + new_chunk_id, + ) + return new_chunk_id + def delete_chunks_by_source_path(self, path) -> int: """Delete every chunk whose ``metadata.source`` equals ``path``. diff --git a/docsgpt/vectorstore/faiss.py b/docsgpt/vectorstore/faiss.py index 03dc4dd6..0c2d1aba 100644 --- a/docsgpt/vectorstore/faiss.py +++ b/docsgpt/vectorstore/faiss.py @@ -306,6 +306,58 @@ class FaissStore(BaseVectorStore): self._save_to_storage() return ids[0] + def update_chunk(self, chunk_id: str, text: str, metadata: Dict[str, Any]) -> str: + """Replace a chunk in place, keeping its id and its place in the docstore. + + The new vector is computed before anything changes, so a failed embed + leaves the chunk as it was. Its old row is removed from the index and + the new vector appended; the row mapping is renumbered to match, and + the docstore entry is replaced where it stands so :meth:`get_chunks` + keeps its order. Saved to storage once. + + Args: + chunk_id: Id of the chunk to replace. + text: The chunk's new text. + metadata: The chunk's complete new metadata. + + Returns: + ``chunk_id``, unchanged. + + Raises: + KeyError: If ``chunk_id`` is not in this index. + ValueError: If the new vector's width does not match the index. + """ + if chunk_id not in self.documents: + raise KeyError(f"Chunk id not found in index: {chunk_id}") + rows_by_id = {doc_id: row for row, doc_id in self.index_to_docstore_id.items()} + if chunk_id not in rows_by_id: + raise KeyError(f"Chunk id has no row in the FAISS index: {chunk_id}") + + vector = np.array(self.embeddings.embed_documents([text]), dtype=np.float32) + if vector.ndim != 2 or vector.shape != (1, self.index.d): + raise ValueError( + f"Embedding for chunk {chunk_id} has shape {vector.shape}, " + f"expected (1, {self.index.d})" + ) + + row = rows_by_id[chunk_id] + self.index.remove_ids(np.array([row], dtype=np.int64)) + self.index.add(vector) + # remove_ids compacts the index and add appends, so the edited chunk's + # vector is now the last row; renumber as delete_index does. + remaining = [ + doc_id + for current_row, doc_id in sorted(self.index_to_docstore_id.items()) + if current_row != row + ] + self.index_to_docstore_id = dict(enumerate(remaining + [chunk_id])) + self.documents[chunk_id] = { + "page_content": text, + "metadata": dict(metadata or {}), + } + self._save_to_storage() + return chunk_id + def delete_chunk(self, chunk_id: str) -> bool: """Delete a chunk and save to storage.""" self.delete_index([chunk_id]) diff --git a/docsgpt/vectorstore/mongodb.py b/docsgpt/vectorstore/mongodb.py index cb735310..96c58cb1 100644 --- a/docsgpt/vectorstore/mongodb.py +++ b/docsgpt/vectorstore/mongodb.py @@ -2,7 +2,7 @@ import logging from functools import cached_property from docsgpt.core.settings import settings -from docsgpt.vectorstore.base import BaseVectorStore +from docsgpt.vectorstore.base import BaseVectorStore, InvalidChunkMetadataError from docsgpt.vectorstore.document_class import Document @@ -224,6 +224,68 @@ class MongoDBVectorStore(BaseVectorStore): result = self._collection.insert_one(chunk_data) return str(result.inserted_id) + def update_chunk(self, chunk_id: str, text: str, metadata: dict) -> str: + """Rewrite a chunk's document in place, keeping its ``_id``. + + Metadata lives as top-level fields, so the new metadata is ``$set`` + and any field the record carries but the new metadata lacks is + ``$unset``; the record then holds exactly the new metadata. ``_id``, + the text, the embedding and ``source_id`` are reserved and cannot be + overwritten through ``metadata``. The embedding is computed before + anything is written. + + Args: + chunk_id: Id of the chunk to replace. + text: The chunk's new text. + metadata: The chunk's complete new metadata. + + Returns: + ``chunk_id``, unchanged. + + Raises: + KeyError: If this source has no chunk with that id. + InvalidChunkMetadataError: If a metadata key is empty, contains + ``.`` or starts with ``$``. Mongo reads those as a path or an + operator, so ``$set`` would fail or write a nested field. + ValueError: If no embedding could be generated. + """ + from bson.objectid import ObjectId + + for key in metadata or {}: + if not isinstance(key, str) or not key or "." in key or key.startswith("$"): + raise InvalidChunkMetadataError( + f"Metadata key {key!r} is not allowed: keys must be non-empty, " + "contain no '.' and not start with '$'" + ) + + query = {"_id": ObjectId(chunk_id), "source_id": self._source_id} + existing = self._collection.find_one(query) + if existing is None: + raise KeyError(f"Chunk {chunk_id} not found for source {self._source_id}") + + embeddings = self._embedding.embed_documents([text]) + if not embeddings: + raise ValueError("Could not generate embedding for chunk") + + reserved = {"_id", self._text_key, self._embedding_key, "source_id"} + fields = {k: v for k, v in (metadata or {}).items() if k not in reserved} + update = { + "$set": { + self._text_key: text, + self._embedding_key: embeddings[0], + "source_id": self._source_id, + **fields, + } + } + stale = {k: "" for k in existing if k not in reserved and k not in fields} + if stale: + update["$unset"] = stale + + result = self._collection.update_one(query, update) + if result.matched_count == 0: + raise KeyError(f"Chunk {chunk_id} not found for source {self._source_id}") + return chunk_id + def delete_chunk(self, chunk_id): try: from bson.objectid import ObjectId diff --git a/docsgpt/vectorstore/pgvector.py b/docsgpt/vectorstore/pgvector.py index 60bb5262..ca657b7d 100644 --- a/docsgpt/vectorstore/pgvector.py +++ b/docsgpt/vectorstore/pgvector.py @@ -637,7 +637,7 @@ class PGVectorStore(BaseVectorStore): select_query = f""" SELECT id, {self._text_column}, {self._metadata_column} FROM {self._table_name} - WHERE source_id = %s; + WHERE source_id = %s ORDER BY id; """ cursor.execute(select_query, (self._source_id,)) results = cursor.fetchall() @@ -704,6 +704,58 @@ class PGVectorStore(BaseVectorStore): finally: cursor.close() + def update_chunk(self, chunk_id: str, text: str, metadata: Dict[str, Any]) -> str: + """Rewrite a chunk's row in place, keeping its id. + + One ``UPDATE`` scoped to this source. The embedding is computed first, + so a failed embed writes nothing. ``get_chunks`` orders by id, so the + chunk also keeps its place in the list. + + Args: + chunk_id: Id of the chunk to replace. + text: The chunk's new text. + metadata: The chunk's complete new metadata; ``source_id`` is + stamped on it as :meth:`add_chunk` does. + + Returns: + ``chunk_id``, unchanged. + + Raises: + KeyError: If this source has no chunk with that id. + ValueError: If no embedding could be generated. + """ + final_metadata = dict(metadata or {}) + final_metadata["source_id"] = self._source_id + + embeddings = self._embedding.embed_documents([text]) + if not embeddings: + raise ValueError("Could not generate embedding for chunk") + + conn = self._get_connection() + cursor = conn.cursor() + + try: + update_query = f""" + UPDATE {self._table_name} + SET {self._text_column} = %s, {self._vector_column} = %s, {self._metadata_column} = %s + WHERE id = %s AND source_id = %s; + """ + cursor.execute( + update_query, + (text, embeddings[0], Jsonb(final_metadata), int(chunk_id), self._source_id), + ) + if cursor.rowcount == 0: + raise KeyError(f"Chunk {chunk_id} not found for source {self._source_id}") + conn.commit() + return str(chunk_id) + + except Exception as e: + conn.rollback() + logging.error(f"Error updating chunk: {e}") + raise + finally: + cursor.close() + def delete_chunk(self, chunk_id: str) -> bool: """Delete a specific chunk by its ID""" conn = self._get_connection() diff --git a/docsgpt/vectorstore/qdrant.py b/docsgpt/vectorstore/qdrant.py index a3822e9b..5b169416 100644 --- a/docsgpt/vectorstore/qdrant.py +++ b/docsgpt/vectorstore/qdrant.py @@ -206,6 +206,50 @@ class QdrantStore(BaseVectorStore): ids = self.add_texts([text], [metadata or {}]) return ids[0] + def update_chunk(self, chunk_id: str, text: str, metadata: Dict[str, Any]) -> str: + """Overwrite a chunk's point in place, keeping its id. + + Upserting the same id replaces the vector and payload; ``scroll`` + orders by point id, so the chunk also keeps its place in + :meth:`get_chunks`. The payload has the shape :meth:`add_texts` writes. + + Args: + chunk_id: Id of the point to replace. + text: The chunk's new text. + metadata: The chunk's complete new metadata; ``source_id`` is + stamped on it for source scoping. + + Returns: + ``chunk_id``, unchanged. + + Raises: + KeyError: If this source has no point with that id. + """ + records = self._client.retrieve( + collection_name=self._collection, + ids=[chunk_id], + with_payload=True, + with_vectors=False, + ) + payload = (records[0].payload or {}) if records else {} + if (payload.get("metadata") or {}).get("source_id") != self._source_id: + raise KeyError(f"Chunk {chunk_id} not found for source {self._source_id}") + + vector = self._embeddings.embed_documents([text])[0] + payload_metadata = dict(metadata or {}) + payload_metadata["source_id"] = self._source_id + self._client.upsert( + collection_name=self._collection, + points=[ + self._models.PointStruct( + id=chunk_id, + vector=vector, + payload={"page_content": text, "metadata": payload_metadata}, + ) + ], + ) + return chunk_id + def delete_chunk(self, chunk_id: str) -> bool: """Delete a single chunk by id.""" try: diff --git a/frontend/DESIGN.md b/frontend/DESIGN.md index 078e4ed8..eaf5651a 100644 --- a/frontend/DESIGN.md +++ b/frontend/DESIGN.md @@ -71,7 +71,12 @@ ones `warning`. Series with no meaning (models, agents, sources) take `chart-1` to `chart-5` in order. When there are more than five, the four largest keep `chart-1` to `chart-4` and the rest are summed into one "Other" series in `chart-5`, so no colour repeats (`settings/foldSeries.ts`). -`secondary` is never a chart colour. +`secondary` is never a chart colour. The knowledge graph is the one +place Other is not `chart-5`: its entity types have no meaning to map, and +half a graph often folds into Other, so red would read as "these failed". +There the four largest types take `chart-1` to `chart-4` and Other is +`muted-foreground` (`foldGraphTypes`, `readGraphPalette` in +`components/graphViewUtils.ts`). The status set is `success | warning | destructive | info`. `default` is the component's own base tone (brand on a Badge or Button, quiet on an Alert or @@ -266,8 +271,17 @@ aria-current="page">` with `data-active`. The the same pixels, plus `role="tablist"`/`"tab"`, `aria-selected` and arrow-key focus; wrap the panel in `TabsContent` (FilePicker's My Files / Shared with Me, Schedules' Recurring / One-time in - `agents/schedules/SchedulesView.tsx`). The `default` variant is a pill tab, - unused in the app. + `agents/schedules/SchedulesView.tsx`, the graph source view, the source + edit drawer). Each trigger keeps its `px-4`, so the first label sits 16px + in from the content edge and the underline runs past the label on both + sides. That inset is deliberate (decided 2026-09-28): the tab row reads as + its own strip, with a wider target per tab, so don't pull it flush with + `-ml-4` or strip the padding. Only the workflow builder's toolbar, where + the tabs share a row with a breadcrumb, uses the padding-free `Button +variant="tab" size="inline"`. The `default` variant is a pill tab, unused + in the app. Panels unmount when hidden, except one whose state is costly to + rebuild (the graph source view's laid-out canvas and zoom): that + `TabsContent` takes `forceMount` plus `data-[state=inactive]:hidden`. - A section panel's disclosure header (NewAgent's Advanced and Guardrails panels) is `variant="section-toggle" size="sm"` with `-ml-3 w-fit justify-start` and `aria-expanded`: a lucide `ChevronRight` first @@ -499,7 +513,10 @@ that already draws the frame (the renaming sidebar row, a search strip in a bordered panel) is `variant="bare"`; the host shows focus, and any inset padding goes on the host, not the field. A field on a muted panel (the ImportSpec Base URL box) is `variant="filled"`, so it keeps the card fill -instead of showing the panel through; never pass `bg-card` for it. +instead of showing the panel through; never pass `bg-card` for it. A floating +label rests at the field's own text size (16px, 14px from `md`), so a +labelled search and a placeholder-only one read the same, and with a +`leftIcon` it rests where the text starts (40px in). ### SelectTrigger (`ui/select.tsx`) @@ -520,9 +537,29 @@ lg` and CommandInput. The `sm` sizes stay 14px. A highlighted list row is `none`, `vertical` (default), `both`. Same border, ring and invalid styling as Input. `variant`: `default` (transparent), `filled` (card fill), with the same rule as Input: a textarea on a muted panel is `variant="filled"`, -never `bg-card`. Three raw `