Coverage for src/lilbee/server/routes/documents.py: 100%

66 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-09-28 17:20 +0000

1"""Document management route handlers: add, list, remove, sync, export. 

2 

3Every route needs the token, reads included: the listing names the user's 

4files and ``/api/export`` serializes the whole corpus. 

5""" 

6 

7from __future__ import annotations 

8 

9import asyncio 

10from typing import Annotated 

11 

12from litestar import Request, Response, get, post 

13from litestar.datastructures import UploadFile 

14from litestar.exceptions import ValidationException 

15from litestar.params import FromQuery, MultipartBody, QueryParameter 

16from litestar.response import Stream 

17from pydantic import BaseModel, Field 

18 

19from lilbee.core.config import validate_ocr_timeout 

20from lilbee.server import handlers 

21from lilbee.server.content_disposition import CONTENT_DISPOSITION 

22from lilbee.server.handlers.sse import SSE_MEDIA_TYPE 

23from lilbee.server.models import ( 

24 AddRequest, 

25 DocumentListResponse, 

26 DocumentRemoveResponse, 

27 SyncRequest, 

28) 

29 

30 

31class RemoveRequest(BaseModel): 

32 """Request body for /api/documents/remove.""" 

33 

34 names: list[str] = Field(max_length=100) 

35 

36 

37@post("/api/sync", media_type=SSE_MEDIA_TYPE) 

38async def sync_route(data: SyncRequest | None = None) -> Stream: 

39 """Re-index changed documents with streaming SSE progress events. 

40 

41 Pass ``{"force_rebuild": true}`` to wipe the store and re-ingest every file 

42 under the current ``cfg.embedding_model``. This is the recovery path after 

43 a ``PUT /api/models/embedding`` that returned ``reindex_required=true``. 

44 Pass ``{"retry_skipped": true}`` for the lighter path: retry the files that 

45 failed a previous sync without dropping the store. 

46 Pass ``{"prune_ignored": true}`` to also drop sources a ``.lilbeeignore`` 

47 now excludes; without it, sync leaves already-indexed sources alone. 

48 """ 

49 enable_ocr = data.enable_ocr if data else None 

50 ocr_timeout = data.ocr_timeout if data else None 

51 force_rebuild = data.force_rebuild if data else False 

52 retry_skipped = data.retry_skipped if data else False 

53 prune_ignored = data.prune_ignored if data else False 

54 try: 

55 validate_ocr_timeout(ocr_timeout) 

56 except ValueError as exc: 

57 raise ValidationException(str(exc)) from exc 

58 return Stream( 

59 handlers.sync_stream( 

60 enable_ocr=enable_ocr, 

61 ocr_timeout=ocr_timeout, 

62 force_rebuild=force_rebuild, 

63 retry_skipped=retry_skipped, 

64 prune_ignored=prune_ignored, 

65 ), 

66 media_type=SSE_MEDIA_TYPE, 

67 ) 

68 

69 

70@post("/api/add", media_type=SSE_MEDIA_TYPE) 

71async def add_route(data: AddRequest) -> Stream: 

72 """Add files to the knowledge base with streaming SSE progress.""" 

73 try: 

74 paths, force, enable_ocr, ocr_timeout = handlers.validate_add_paths(data.model_dump()) 

75 except ValueError as exc: 

76 raise ValidationException(str(exc)) from exc 

77 return Stream( 

78 handlers.add_files_stream( 

79 paths, force=force, enable_ocr=enable_ocr, ocr_timeout=ocr_timeout 

80 ), 

81 media_type=SSE_MEDIA_TYPE, 

82 status_code=201, 

83 ) 

84 

85 

86@post("/api/add/upload", media_type=SSE_MEDIA_TYPE) 

87async def add_upload_route( 

88 data: MultipartBody[list[UploadFile]], 

89 enable_ocr: FromQuery[bool | None] = None, 

90 ocr_timeout: FromQuery[float | None] = None, 

91) -> Stream: 

92 """Ingest uploaded file content with streaming SSE progress. 

93 

94 Unlike /api/add, which reads server-side paths, this accepts the client's raw 

95 file bytes. That lets a client whose files the server cannot read by path -- 

96 e.g. the plugin or CLI in external mode against a remote lilbee / GPU box -- 

97 ingest its own local files by uploading them straight to the server. 

98 ``enable_ocr`` and ``ocr_timeout`` are query parameters because the request 

99 body is the upload's raw multipart file list. 

100 """ 

101 # Names first, bytes second: reading every part before validating cost a 

102 # full in-memory copy of a payload that was going to be rejected anyway. 

103 try: 

104 names = handlers.validate_upload_names([upload.filename for upload in data]) 

105 validate_ocr_timeout(ocr_timeout) 

106 except ValueError as exc: 

107 raise ValidationException(str(exc)) from exc 

108 cleaned = [(name, await upload.read()) for name, upload in zip(names, data, strict=True)] 

109 return Stream( 

110 handlers.add_uploads_stream(cleaned, enable_ocr=enable_ocr, ocr_timeout=ocr_timeout), 

111 media_type=SSE_MEDIA_TYPE, 

112 status_code=201, 

113 ) 

114 

115 

116@get("/api/documents") 

117async def documents_list_route( 

118 search: FromQuery[str] = "", 

119 limit: Annotated[int, QueryParameter(ge=1, le=1000)] = 50, 

120 offset: Annotated[int, QueryParameter(ge=0)] = 0, 

121) -> DocumentListResponse: 

122 """List indexed documents with metadata, paginated and searchable.""" 

123 return await handlers.list_documents(search=search, limit=limit, offset=offset) 

124 

125 

126@post("/api/documents/remove") 

127async def documents_remove_route(data: RemoveRequest) -> DocumentRemoveResponse: 

128 """Remove documents from the knowledge base by source name.""" 

129 return await handlers.delete_documents(data.names) 

130 

131 

132@get("/api/export", media_type="application/octet-stream") 

133async def export_route( 

134 fmt: Annotated[str, QueryParameter(name="format")] = "", 

135 source: FromQuery[str] = "", 

136) -> Response[bytes]: 

137 """Download the per-page text dataset as a file (parquet by default). 

138 

139 The media type is declared on the decorator as well as on the returned 

140 Response for the same reason the streaming routes declare theirs: litestar 

141 documents the content type from the decorator, so without it the schema 

142 promises JSON and a generated client parses a parquet file as text. 

143 """ 

144 from lilbee.app.dataset import DatasetError, export_to_bytes 

145 

146 try: 

147 # export_to_bytes serializes the whole per-page dataset into memory; 

148 # offload so a large export doesn't stall every other request, matching 

149 # get_source_content's own off-loop read. 

150 payload = await asyncio.to_thread(export_to_bytes, fmt, source or None) 

151 except DatasetError as exc: 

152 raise ValidationException(str(exc)) from exc 

153 return Response( 

154 content=payload.data, 

155 media_type="application/octet-stream", 

156 headers={CONTENT_DISPOSITION: f'attachment; filename="pages.{payload.fmt}"'}, 

157 ) 

158 

159 

160@post("/api/import", media_type=SSE_MEDIA_TYPE) 

161async def import_route( 

162 request: Request, 

163 fmt: Annotated[str, QueryParameter(name="format")] = "", 

164) -> Stream: 

165 """Import an uploaded per-page dataset with streaming SSE progress events. 

166 

167 The request body is the raw dataset bytes; ``?format=parquet|jsonl`` is 

168 required since there is no filename to infer from. Bounded by the server's 

169 body-size limit; larger datasets use the path-based CLI/MCP import. 

170 """ 

171 from lilbee.app.dataset import DatasetError, require_format 

172 

173 try: 

174 require_format(fmt) 

175 except DatasetError as exc: 

176 raise ValidationException(str(exc)) from exc 

177 return Stream( 

178 handlers.import_stream(await request.body(), fmt), 

179 media_type=SSE_MEDIA_TYPE, 

180 status_code=201, 

181 )