Files
2026-09-29 10:42:39 +04:00

419 lines
16 KiB
Python

"""Source document management chunk management."""
import math
from flask import current_app, jsonify, make_response, request
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.resource_access import AccessDenied
from docsgpt.api.user.sources.access import denied_response, load_source
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"
)
def _resolve_source(doc_id: str, user: str, action: str = "use") -> dict:
"""Resolve a source (UUID or legacy ObjectId) the caller may ``action`` on.
``use`` (browse chunks) is open to every role; ``edit`` (add / delete /
update chunks) needs owner or team editor. The vector partition is keyed
by source id, so a team editor's write needs no owner id.
Args:
doc_id: Source id from the request.
user: The caller's ``sub``.
action: ``use`` for reads, ``edit`` for chunk writes.
Returns:
dict: The source row (PG UUID in ``id``).
Raises:
AccessDenied: 404 when not visible, 403 when the role can't do it.
"""
with db_readonly() as conn:
doc, _ra = load_source(conn, doc_id, user, action)
return doc
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.
Stores round-trip metadata differently -- pgvector keeps JSON types, the
Mongo backend can hand back strings -- so a numeric string counts as
recorded. Anything else (missing, empty, non-numeric, zero, negative, or
non-finite) does not: ``float("inf")`` is greater than zero but is not a
number of tokens, and it reaches ``toLocaleString`` in the UI as "∞".
Args:
metadata: A chunk's metadata mapping.
Returns:
True when ``token_count`` holds a finite positive number.
"""
raw = metadata.get("token_count")
if isinstance(raw, bool) or not isinstance(raw, (int, float, str)):
return False
try:
value = float(raw)
except (TypeError, ValueError):
return False
return math.isfinite(value) and value > 0
def _with_token_counts(chunks: list) -> list:
"""Fill in ``metadata.token_count`` for chunks that were stored without one.
Ingestion records a per-chunk count in the embedding model's tokenizer,
but chunks indexed before that was written -- and any path that rebuilt a
chunk's metadata from scratch -- reach the UI without the key, which then
renders a bare "-". The count recomputed here is cl100k rather than the
embedding model's tokenizer: it is a display fallback, and loading the
model's tokenizer would put a Hugging Face download in the request path.
Only the page being returned is counted, so the cost is bounded by
``per_page`` rather than by the size of the index.
Args:
chunks: The chunk dicts about to be serialised.
Returns:
The same list, with each chunk's metadata normalised to a dict that
carries a ``token_count``.
"""
for chunk in chunks:
metadata = chunk.get("metadata") or {}
if not _has_usable_token_count(metadata):
metadata["token_count"] = num_tokens_from_string(chunk.get("text") or "")
chunk["metadata"] = metadata
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(
description="Retrieves chunks from a document, optionally filtered by file path and search term",
params={
"id": "The document ID",
"page": "Page number for pagination",
"per_page": "Number of chunks per page",
"path": "Optional: Filter chunks by relative file path",
"search": "Optional: Search term to filter chunks by title or content",
},
)
def get(self):
decoded_token = request.decoded_token
if not decoded_token:
return make_response(jsonify({"success": False}), 401)
user = decoded_token.get("sub")
doc_id = request.args.get("id")
page = int(request.args.get("page", 1))
per_page = int(request.args.get("per_page", 10))
path = request.args.get("path")
search_term = request.args.get("search", "").strip().lower()
if not doc_id:
return make_response(jsonify({"error": "Invalid doc_id"}), 400)
try:
doc = _resolve_source(doc_id, user)
except AccessDenied as err:
return denied_response(err)
except Exception as e:
current_app.logger.error(f"Error resolving source: {e}", exc_info=True)
return make_response(jsonify({"error": "Invalid doc_id"}), 400)
resolved_id = str(doc["id"])
try:
store = get_vector_store(resolved_id)
chunks = store.get_chunks()
filtered_chunks = []
for chunk in chunks:
metadata = chunk.get("metadata", {})
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()
if not (text_match or title_match):
continue
filtered_chunks.append(chunk)
chunks = filtered_chunks
total_chunks = len(chunks)
start = (page - 1) * per_page
end = start + per_page
paginated_chunks = _with_token_counts(chunks[start:end])
return make_response(
jsonify(
{
"page": page,
"per_page": per_page,
"total": total_chunks,
"chunks": paginated_chunks,
"path": path if path else None,
"search": search_term if search_term else None,
}
),
200,
)
except Exception as e:
current_app.logger.error(f"Error getting chunks: {e}", exc_info=True)
return make_response(jsonify({"success": False}), 500)
@sources_chunks_ns.route("/add_chunk")
class AddChunk(Resource):
@api.expect(
api.model(
"AddChunkModel",
{
"id": fields.String(required=True, description="Document ID"),
"text": fields.String(required=True, description="Text of the chunk"),
"metadata": fields.Raw(
required=False,
description="Metadata associated with the chunk",
),
},
)
)
@api.doc(
description="Adds a new chunk to the document",
)
def post(self):
decoded_token = request.decoded_token
if not decoded_token:
return make_response(jsonify({"success": False}), 401)
user = decoded_token.get("sub")
data = request.get_json()
required_fields = ["id", "text"]
missing_fields = check_required_fields(data, required_fields)
if missing_fields:
return missing_fields
doc_id = data.get("id")
text = data.get("text")
metadata = data.get("metadata", {})
token_count = num_tokens_from_string(text)
metadata["token_count"] = token_count
try:
doc = _resolve_source(doc_id, user, "edit")
except AccessDenied as err:
return denied_response(err)
except Exception as e:
current_app.logger.error(f"Error resolving source: {e}", exc_info=True)
return make_response(jsonify({"error": "Invalid doc_id"}), 400)
try:
store = get_vector_store(str(doc["id"]))
chunk_id = store.add_chunk(text, metadata)
return make_response(
jsonify({"message": "Chunk added successfully", "chunk_id": chunk_id}),
201,
)
except Exception as e:
current_app.logger.error(f"Error adding chunk: {e}", exc_info=True)
return make_response(jsonify({"success": False}), 500)
@sources_chunks_ns.route("/delete_chunk")
class DeleteChunk(Resource):
@api.doc(
description="Deletes a specific chunk from the document.",
params={"id": "The document ID", "chunk_id": "The ID of the chunk to delete"},
)
def delete(self):
decoded_token = request.decoded_token
if not decoded_token:
return make_response(jsonify({"success": False}), 401)
user = decoded_token.get("sub")
doc_id = request.args.get("id")
chunk_id = request.args.get("chunk_id")
try:
doc = _resolve_source(doc_id, user, "edit")
except AccessDenied as err:
return denied_response(err)
except Exception as e:
current_app.logger.error(f"Error resolving source: {e}", exc_info=True)
return make_response(jsonify({"error": "Invalid doc_id"}), 400)
try:
store = get_vector_store(str(doc["id"]))
deleted = store.delete_chunk(chunk_id)
if deleted:
return make_response(
jsonify({"message": "Chunk deleted successfully"}), 200
)
else:
return make_response(
jsonify({"message": "Chunk not found or could not be deleted"}),
404,
)
except Exception as e:
current_app.logger.error(f"Error deleting chunk: {e}", exc_info=True)
return make_response(jsonify({"success": False}), 500)
@sources_chunks_ns.route("/update_chunk")
class UpdateChunk(Resource):
@api.expect(
api.model(
"UpdateChunkModel",
{
"id": fields.String(required=True, description="Document ID"),
"chunk_id": fields.String(
required=True, description="Chunk ID to update"
),
"text": fields.String(
required=False, description="New text of the chunk"
),
"metadata": fields.Raw(
required=False,
description="Updated metadata associated with the chunk",
),
},
)
)
@api.doc(
description="Updates an existing chunk in the document.",
)
def put(self):
decoded_token = request.decoded_token
if not decoded_token:
return make_response(jsonify({"success": False}), 401)
user = decoded_token.get("sub")
data = request.get_json()
required_fields = ["id", "chunk_id"]
missing_fields = check_required_fields(data, required_fields)
if missing_fields:
return missing_fields
doc_id = data.get("id")
chunk_id = data.get("chunk_id")
text = data.get("text")
metadata = data.get("metadata")
if text is not None:
token_count = num_tokens_from_string(text)
if metadata is None:
metadata = {}
metadata["token_count"] = token_count
try:
doc = _resolve_source(doc_id, user, "edit")
except AccessDenied as err:
return denied_response(err)
except Exception as e:
current_app.logger.error(f"Error resolving source: {e}", exc_info=True)
return make_response(jsonify({"error": "Invalid doc_id"}), 400)
try:
store = get_vector_store(str(doc["id"]))
chunks = store.get_chunks()
existing_chunk = next((c for c in chunks if c["doc_id"] == chunk_id), None)
if not existing_chunk:
return make_response(jsonify({"error": "Chunk not found"}), 404)
new_text = text if text is not None else existing_chunk["text"]
if metadata is not None:
new_metadata = existing_chunk["metadata"].copy()
new_metadata.update(metadata)
else:
new_metadata = existing_chunk["metadata"].copy()
if text is not None:
new_metadata["token_count"] = num_tokens_from_string(new_text)
try:
# 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(
{
"message": "Chunk updated successfully",
"chunk_id": new_chunk_id,
"original_chunk_id": chunk_id,
}
),
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 update chunk {chunk_id}: {add_error}")
return make_response(
jsonify({"error": "Failed to update chunk - addition failed"}), 500
)
except Exception as e:
current_app.logger.error(f"Error updating chunk: {e}", exc_info=True)
return make_response(jsonify({"success": False}), 500)