diff --git a/.gitignore b/.gitignore index 29f9787..a81f60d 100644 --- a/.gitignore +++ b/.gitignore @@ -9,6 +9,9 @@ build/ # Frontend node_modules/ +# 前端构建产物(本地残留的 Vue 构建 / 符号链接,勿提交;线上由 frontend/dist 挂载) +app/static/admin + # Env .env .env.local diff --git a/CLAUDE.md b/CLAUDE.md index 0c62ad9..d05d86a 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -90,12 +90,15 @@ QMDSearch/ | POST | `/api/v1/search` | 分层检索 | 免登录 | | POST | `/api/v1/documents` | 文档入库(202 异步入库,返回 task_id) | Bearer | | POST | `/api/v1/documents/upload` | multipart 文件上传入库(202 异步,支持 .txt/.md/.html/.htm/.pdf/.docx) | Bearer | +| POST | `/api/v1/documents/upload-batch` | 批量文件上传入库(多文件,202 异步) | Bearer | +| GET | `/api/v1/documents/tasks` | 入库任务列表(近期,按时间降序) | Bearer | | GET | `/api/v1/documents/tasks/{task_id}` | 入库任务状态查询(done 附 result,failed 附 error) | 免登录 | | GET | `/api/v1/knowledge/categories` | 知识分类类目集 | 免登录 | | GET | `/api/v1/knowledge/stats` | 统计(四层点数 + 类目分布 + uncategorized 数) | 免登录 | | GET | `/api/v1/documents` | 文档列表(limit/offset 分页) | 免登录 | | GET | `/api/v1/documents/{doc_id}` | 文档详情 | 免登录 | | GET | `/api/v1/documents/{doc_id}/file` | 下载关联的原始文件(免登录) | 免登录 | +| POST | `/api/v1/documents/{doc_id}/reingest` | 重新摘要入库(读原文件→删旧→重跑流水线) | Bearer | | DELETE | `/api/v1/documents/{doc_id}` | 删除文档(幂等) | Bearer | | POST | `/api/v1/auth/login` | 用户名密码登录,签发 session token(TTL 12h) | 免登录 | | POST | `/api/v1/auth/logout` | 退出登录(删除当前 session) | Bearer | @@ -108,13 +111,13 @@ QMDSearch/ | DELETE | `/api/v1/auth/users/{username}` | 删除用户(清除其 session;禁删自己/最后一个 admin) | Bearer + admin | | GET | `/admin` | 管理页面 | 页面登录门禁 | -「Bearer」指请求头 `Authorization: Bearer `,token 经 `/api/v1/auth/login` 获取;变更类文档端点(POST /documents、POST /documents/upload、DELETE /documents/{id})需 Bearer token(admin/user 角色均可),查询类端点免登录。 +「Bearer」指请求头 `Authorization: Bearer `,token 经 `/api/v1/auth/login` 获取;变更类文档端点(POST /documents、POST /documents/upload、POST /documents/upload-batch、POST /documents/{id}/reingest、DELETE /documents/{id})与 GET /documents/tasks 任务列表需 Bearer token(admin/user 角色均可),其余查询类端点免登录。 入库任务状态持久化在 Redis(key: `ingest_task:{task_id}`):进行中与 done 保留 24h,failed 保留 7 天;Redis 不可用时降级为纯内存。 ## 管理页面 -浏览器访问 `/admin`,页面带登录门禁(未登录/token 失效自动回登录卡片;must_change_password 用户先强制改密后方可进入)。单页面含八个区块(admin 视角;user 角色无用户管理区块):概览、个人中心(user 角色默认进入且可见,含账号信息/改密/退出)、文档管理(列表/详情/删除)、文档入库(文本 + 文件上传)、检索测试台、类目列表、API 指南(端点清单 + 在线测试台)、用户管理(仅 admin 角色挂载,含搜索筛选/行内编辑角色与启用状态/弹窗式重置密码与删除)。 +浏览器访问 `/admin`,页面带登录门禁(未登录/token 失效自动回登录卡片;must_change_password 用户先强制改密后方可进入)。单页面含九个区块(admin 视角;user 角色无用户管理区块):概览、个人中心(user 角色默认进入且可见,含账号信息/改密/退出)、文档管理(列表/详情/删除)、文档入库(文本 + 文件上传,文件上传支持批量多选)、入库进度(展示近期任务与状态,2s 轮询自动刷新)、检索测试台、类目列表、API 指南(端点清单 + 在线测试台)、用户管理(仅 admin 角色挂载,含搜索筛选/行内编辑角色与启用状态/弹窗式重置密码与删除)。 ## 认证与用户 diff --git a/app/api/deps.py b/app/api/deps.py index 57ad600..217f9e1 100644 --- a/app/api/deps.py +++ b/app/api/deps.py @@ -8,6 +8,7 @@ UserStore/SessionStore 为模块级懒加载单例:Redis 客户端创建失败 import structlog from fastapi import Depends, Header from redis import asyncio as redis_async +from redis import Redis as RedisSync from app.api.response import ApiError from app.config import settings @@ -22,11 +23,20 @@ _session_store: SessionStore | None = None def _create_redis_client() -> redis_async.Redis | None: - """创建 redis.asyncio 客户端(decode_responses=True);失败返回 None 走内存降级""" + """创建 redis.asyncio 客户端(decode_responses=True);连接失败返回 None 走内存降级 + + 注意:redis.asyncio.from_url() 仅构造客户端对象,不会实际连接 Redis, + 因此需要用同步客户端 ping 一次确认连接可用,否则后续操作才会报错, + 导致 UserStore/SessionStore 无法降级为内存模式。 + """ try: + # 同步 ping 确认 Redis 可达(from_url 不会实际连接) + sync_client = RedisSync.from_url(settings.redis_url, decode_responses=True) + sync_client.ping() + sync_client.close() return redis_async.from_url(settings.redis_url, decode_responses=True) except Exception: - logger.warning("Redis 客户端创建失败,认证存储降级为内存模式", exc_info=True) + logger.warning("Redis 连接失败,认证存储降级为内存模式", exc_info=True) return None diff --git a/app/api/v1/document.py b/app/api/v1/document.py index dba5538..6abb00c 100644 --- a/app/api/v1/document.py +++ b/app/api/v1/document.py @@ -183,6 +183,149 @@ async def upload_document( ) +async def _process_single_upload(file: UploadFile) -> dict[str, Any]: + """处理单个上传文件:校验 → 提取文本 → 落盘 → 提交入库 + + 供批量上传复用:返回成功 {"filename", "task_id"} 或失败 {"filename", "error"}。 + title/source 使用文件名默认值,不接受用户传入的额外表单字段。 + """ + original_filename = file.filename or "unnamed" + ext = Path(original_filename).suffix.lower() + + # 1. 扩展名校验 + allowed = _allowed_extensions() + if ext not in allowed: + return {"filename": original_filename, "error": f"不支持的文件类型: {ext or '(无扩展名)'}"} + + # 2. 读取字节并校验大小 + try: + content = await file.read() + except Exception as exc: + return {"filename": original_filename, "error": f"文件读取失败: {exc}"} + max_bytes = settings.upload_max_size_mb * 1024 * 1024 + if len(content) > max_bytes: + return {"filename": original_filename, "error": f"文件超过大小上限: {settings.upload_max_size_mb}MB"} + + # 3. 提取文本 + try: + text = parse_file(original_filename, content) + except ValueError as exc: + return {"filename": original_filename, "error": str(exc)} + if not text.strip(): + return {"filename": original_filename, "error": "无法从文件提取文本"} + + # 4. 落盘(按 YYYY/MM 日期分片;失败仅 warning,不阻塞入库) + doc_id = uuid.uuid4().hex + metadata_dict: dict[str, str] = {} + try: + upload_dir = Path(settings.upload_dir).resolve() + shard_subdir = datetime.now(UTC).strftime("%Y/%m") + target_dir = upload_dir / shard_subdir + target_dir.mkdir(parents=True, exist_ok=True) + target = target_dir / f"{doc_id}_{original_filename}" + target.write_bytes(content) + metadata_dict.update( + { + "raw_file_path": str(target), + "original_filename": original_filename, + "original_size_bytes": str(len(content)), + } + ) + logger.info( + "上传文件已落盘", doc_id=doc_id, saved_path=str(target), size=len(content) + ) + except Exception: + logger.warning( + "上传文件落盘失败,仅做文本入库", + doc_id=doc_id, + filename=original_filename, + exc_info=True, + ) + + # 5. 默认 title / source + title = Path(original_filename).stem + source = f"file:{original_filename}" + + # 6. 提交入库流水线 + doc_input = DocumentInput(text=text, title=title, source=source, metadata=metadata_dict) + task_id = await _get_task_manager().submit(doc_input) + return {"filename": original_filename, "task_id": task_id} + + +@router.post("/documents/upload-batch") +async def upload_batch( + files: list[UploadFile] = File(...), + user: UserRecord = Depends(get_current_user), +) -> JSONResponse: + """批量文件上传入库:逐文件复用单文件逻辑,收集成功与失败结果 + + 返回 202 + {tasks: [{filename, task_id}], failed: [{filename, error}]}; + 空文件列表返回 1001。 + """ + if not files: + raise ApiError(1001, "未提供任何文件") + tasks: list[dict[str, Any]] = [] + failed: list[dict[str, Any]] = [] + for file in files: + result = await _process_single_upload(file) + if "task_id" in result: + tasks.append({"filename": result["filename"], "task_id": result["task_id"]}) + else: + failed.append({"filename": result["filename"], "error": result["error"]}) + return JSONResponse(status_code=202, content=ok({"tasks": tasks, "failed": failed})) + + +@router.get("/documents/tasks") +async def list_ingest_tasks( + limit: int = Query(default=20, ge=1, le=100), + user: UserRecord = Depends(get_current_user), +) -> dict[str, Any]: + """列出近期入库任务(合并内存与 Redis 镜像,按 updated_at 降序)""" + tasks = await _get_task_manager().list_tasks(limit=limit) + return ok({"items": tasks, "total": len(tasks)}) + + +@router.post("/documents/{doc_id}/reingest") +async def reingest_document( + doc_id: str, user: UserRecord = Depends(get_current_user) +) -> JSONResponse: + """重新入库:读原始文件 → 删旧数据 → 提交新入库任务 + + 仅对有原始文件落盘记录的文档可重新入库;重新入库会生成新 doc_id。 + """ + meta = await _get_qdrant().get_l1_metadata(doc_id) + if meta is None: + raise ApiError(1004, "文档不存在") + raw_file_path = meta.get("raw_file_path", "") + if not raw_file_path: + raise ApiError(1001, "该文档无原始文件,无法重新入库") + path = Path(raw_file_path) + if not path.is_file(): + raise ApiError(1004, "文件不存在") + + # 读原文件并重新提取文本 + try: + text = parse_file(meta.get("original_filename", path.name), path.read_bytes()) + except ValueError as exc: + raise ApiError(1001, str(exc)) from exc + if not text.strip(): + raise ApiError(1001, "无法从文件提取文本") + + # 删旧数据(四层集合按 doc_id 清除,幂等) + await _get_qdrant().delete_by_doc_id(doc_id) + + # 构造新文档输入(保留原 metadata,title/source 从 meta 取或文件名回退) + original_filename = meta.get("original_filename", path.name) + title = meta.get("title") or Path(original_filename).stem + source = meta.get("source") or f"file:{original_filename}" + doc_input = DocumentInput(text=text, title=title, source=source, metadata=dict(meta)) + + task_id = await _get_task_manager().submit(doc_input) + return JSONResponse( + status_code=202, content=ok({"task_id": task_id, "status": "pending"}) + ) + + @router.get("/documents/tasks/{task_id}") async def get_ingest_task(task_id: str) -> dict[str, Any]: """查询入库任务状态:含 task_id/status/created_at/updated_at,done 附 result,failed 附 error""" diff --git a/app/api/v1/settings.py b/app/api/v1/settings.py index cf9728e..467dd87 100644 --- a/app/api/v1/settings.py +++ b/app/api/v1/settings.py @@ -11,7 +11,8 @@ from fastapi import APIRouter, Depends from pydantic import BaseModel from app.api.response import ApiError, ok -from app.core.auth import AuthUser, get_current_user, require_admin +from app.api.deps import get_current_user, require_admin +from app.core.users import UserRecord from app.core.dedup import invalidate_dedup_strategy_cache from app.core.file_parser import ( list_docx_plugins, @@ -46,7 +47,7 @@ class SettingsUpdateRequest(BaseModel): @router.get("/settings") -async def get_settings(user: AuthUser = Depends(get_current_user)) -> dict[str, Any]: +async def get_settings(user: UserRecord = Depends(get_current_user)) -> dict[str, Any]: """返回当前 RuntimeSettings(任何登录用户可读)""" cfg = get_runtime_settings() return ok(cfg.model_dump(mode="json")) @@ -55,7 +56,7 @@ async def get_settings(user: AuthUser = Depends(get_current_user)) -> dict[str, @router.put("/settings") async def update_settings( body: SettingsUpdateRequest, - user: AuthUser = Depends(require_admin), + user: UserRecord = Depends(require_admin), ) -> dict[str, Any]: """部分更新 RuntimeSettings(仅 admin) @@ -82,7 +83,7 @@ async def update_settings( @router.get("/settings/schema") async def get_settings_schema( - user: AuthUser = Depends(get_current_user), + user: UserRecord = Depends(get_current_user), ) -> dict[str, Any]: """返回可选插件与策略列表(前端 Settings 页渲染选项用)""" return ok( @@ -98,7 +99,7 @@ async def get_settings_schema( @router.post("/settings/reset") async def reset_settings( - user: AuthUser = Depends(require_admin), + user: UserRecord = Depends(require_admin), ) -> dict[str, Any]: """重置 RuntimeSettings 为默认值(仅 admin),同时清缓存""" try: diff --git a/app/core/ingest_tasks.py b/app/core/ingest_tasks.py index 15bd1e6..a94b887 100644 --- a/app/core/ingest_tasks.py +++ b/app/core/ingest_tasks.py @@ -20,7 +20,7 @@ from typing import Any import structlog from app.config import Settings -from app.core.dedup import DEDUP_KEY_PREFIX, get_dedup_strategy +from app.core.dedup import get_dedup_strategy from app.core.ingestion import Ingester, IngestionError from app.models.document import DocumentInput from app.services.redis import RedisCache @@ -79,6 +79,7 @@ class IngestTaskManager: """ task_id = uuid.uuid4().hex now = _utc_now_iso() + filename = self._extract_filename(doc) dedup = get_dedup_strategy(self._redis) # 1. 去重命中:直接置 done,复用旧结果,不调 _run @@ -93,6 +94,7 @@ class IngestTaskManager: "updated_at": now, "result": result_dict, "error": None, + "filename": filename, } self._schedule_mirror(task_id) logger.info( @@ -110,6 +112,7 @@ class IngestTaskManager: "updated_at": now, "result": None, "error": None, + "filename": filename, } self._schedule_mirror(task_id) background = asyncio.create_task(self._run(task_id, doc, dedup)) @@ -135,6 +138,72 @@ class IngestTaskManager: ) return None + async def list_tasks(self, limit: int = 20) -> list[dict[str, Any]]: + """合并内存注册表与 Redis 镜像,去重,按 updated_at 降序,截断 limit + + 每项提取 {task_id, status, filename, created_at, updated_at, doc_id}: + doc_id 从 done 任务的 result.document_id 提取,其余状态为 None。 + """ + records: dict[str, dict[str, Any]] = {} + # 1. 内存注册表(主) + for task_id, record in self._tasks.items(): + records[task_id] = record + # 2. Redis 镜像补充内存中没有的(内存未命中的任务,如重启后只存在 Redis 的历史记录) + for task_id, record in await self._scan_redis_tasks(): + records.setdefault(task_id, record) + # 3. 按 updated_at 降序,截断 limit + sorted_records = sorted( + records.values(), + key=lambda r: r.get("updated_at") or "", + reverse=True, + )[:limit] + # 4. 提取展示字段 + return [self._to_list_item(r) for r in sorted_records] + + async def _scan_redis_tasks(self) -> list[tuple[str, dict[str, Any]]]: + """扫描 Redis 中 ingest_task:* 键,返回 (task_id, record) 列表 + + RedisCache 未暴露公开扫描接口,借道底层 _get_client().scan_iter; + 任何异常降级为空列表,不影响 list_tasks 主流程。 + """ + if self._redis is None: + return [] + get_client = getattr(self._redis, "_get_client", None) + if not callable(get_client): + return [] + try: + client = get_client() + tasks: list[tuple[str, dict[str, Any]]] = [] + async for key in client.scan_iter(match=f"{REDIS_KEY_PREFIX}*"): + key_str = key if isinstance(key, str) else key.decode("utf-8", "replace") + task_id = key_str[len(REDIS_KEY_PREFIX) :] + record = await self._redis.get_json(key_str) + if isinstance(record, dict): + tasks.append((task_id, record)) + return tasks + except Exception: + logger.warning("扫描 Redis 任务键失败", exc_info=True) + return [] + + @staticmethod + def _extract_filename(doc: DocumentInput) -> str | None: + """从文档输入提取展示用文件名:优先 metadata.original_filename,其次 title""" + return doc.metadata.get("original_filename") or doc.title or None + + @staticmethod + def _to_list_item(record: dict[str, Any]) -> dict[str, Any]: + """从完整任务记录提取列表展示字段""" + result = record.get("result") + doc_id = result.get("document_id") if isinstance(result, dict) else None + return { + "task_id": record.get("task_id"), + "status": record.get("status"), + "filename": record.get("filename"), + "created_at": record.get("created_at"), + "updated_at": record.get("updated_at"), + "doc_id": doc_id, + } + async def wait_done(self, task_id: str, timeout: float = 30.0) -> dict[str, Any]: """轮询内存注册表直到任务进入终态(done/failed)或超时 diff --git a/app/main.py b/app/main.py index 749c7e7..c9fadd5 100644 --- a/app/main.py +++ b/app/main.py @@ -6,7 +6,6 @@ import structlog from fastapi import FastAPI, Request from fastapi.exceptions import RequestValidationError from fastapi.responses import FileResponse, JSONResponse, RedirectResponse -from redis import asyncio as redis_async from app.api.response import ApiError, error from app.api.v1.auth import router as auth_router @@ -15,7 +14,7 @@ from app.api.v1.knowledge import router as knowledge_router from app.api.v1.search import router as search_router from app.api.v1.settings import router as settings_router from app.config import settings -from app.core.users import UserStore, bootstrap_admin +from app.core.users import bootstrap_admin from app.services.qdrant import QdrantService logger = structlog.get_logger() @@ -33,13 +32,13 @@ async def lifespan(_: FastAPI) -> AsyncIterator[None]: except Exception: logger.error("Qdrant 集合初始化失败,跳过初始化继续启动") try: - # Redis 客户端创建失败(如 URL 非法)时传 None,UserStore 降级为内存模式 - redis_client = redis_async.from_url(settings.redis_url, decode_responses=True) - except Exception: - redis_client = None - try: + # 复用 deps 的 UserStore 单例,确保 lifespan 创建的 admin 与 API 请求用的是同一实例 + # _get_user_store 内部会检测 Redis 可达性,不可达时降级为内存模式 + from app.api.deps import _get_user_store + + user_store = _get_user_store() # 空库引导默认管理员;明文密码由 bootstrap_admin 内部 warning 打印一次 - await bootstrap_admin(UserStore(redis_client), logger) + await bootstrap_admin(user_store, logger) except Exception: logger.warning("默认管理员初始化失败,跳过", exc_info=True) try: @@ -150,8 +149,8 @@ async def admin_spa(rest: str): _AGENT_SKILL_ZIP = _STATIC_DIR / "agent-skill" / "QMDSearch-Agent-Skill.zip" -@app.get("/agent-skill", include_in_schema=False) -async def agent_skill_download() -> FileResponse | JSONResponse: +@app.get("/agent-skill", include_in_schema=False, response_model=None) +async def agent_skill_download(): """下载 AI Agent Skill 压缩包(qmdsearch-agent skill,含 SKILL.md 与示例)""" if not _AGENT_SKILL_ZIP.is_file(): return JSONResponse( diff --git a/app/static/admin.html b/app/static/admin.html index 4e25161..c2926f1 100644 --- a/app/static/admin.html +++ b/app/static/admin.html @@ -259,6 +259,7 @@ + @@ -340,8 +341,8 @@

或上传文件

- - + +
@@ -356,6 +357,21 @@ + +