Onyx(Craft)用户文件库同步到沙箱的完整架构与实现解析
【免费下载链接】danswerOpen Source AI Platform - AI Chat with advanced features that works with every LLM项目地址: https://gitcode.com/GitHub_Trending/da/danswer
导读
用户上传的 PDF、表格、幻灯片等文件,如何从"用户文件库"进入 AI 智能体(Agent)的沙箱工作区,是 Craft 构建型会话(Build Session)中文件可见性的核心问题。本文以仓库文档 docs/craft/features/user-library-sync.md 为骨架,结合 Onyx 后端源码,系统讲解 User Library 到沙箱同步的分层架构、存储模型、同步流水线、DB 层 API、所有权模型与测试策略。读完你将掌握:文件如何以CRAFT_FILE文档形式入库、如何构建 FileSet 并经 push daemon 原子写入沙箱的/workspace/managed/user_library/挂载点,以及各同步触发点与配额校验的实现细节。
该机制替代了此前在 PR #11042 中被移除的 symbiotic-S3 同步方案,与 skills 推送共用同一条托管内容(managed content)流水线,是理解 Craft 会话文件系统的重要组成部分。
1. 三层架构:镜像 skills 流水线
User Library 同步采用与 skills 推送一致的三层结构,职责边界清晰:
HTTP layer backend/onyx/server/features/build/user_library/api.py │ thin: validate input, call db helpers, return response ▼ DB layer backend/onyx/server/features/build/db/user_library.py │ owns: file_store I/O, Document upserts, ownership checks, quota ▼ Sandbox sync backend/onyx/server/features/build/sandbox/user_library.py builds FileSet from DB + file_store, pushes via push daemon- HTTP 层(api.py):只做输入校验(文件大小、数量、配额、zip 炸弹防护、路径清洗)、调用 DB 层 helper、返回响应,不包含任何业务逻辑;
- DB 层(db/user_library.py):负责 file_store 读写、Document 记录的 upsert、所有权校验与存储配额聚合;
- 沙箱同步层(sandbox/user_library.py):从 DB 与 file_store 构建 FileSet,并通过 push daemon 推送到所有活跃沙箱。
文件统一存放在默认文件存储(get_default_file_store())中——与 skills 使用的后端完全一致,Craft 并没有为用户文件单独准备 S3 bucket。这样文件字节始终落在统一的 file_store 抽象之下,无论底层是本地磁盘还是 S3。
2. 存储模型:一个确定性 doc_id 串联全部信息
User Library 复用 Onyx 既有的 Document 模型,不新增任何数据库表。各字段的分工如下表:
| 字段 | 存放位置 | 用途 |
|---|---|---|
| 文件内容 | file_store(S3/本地) | 原始字节 |
Document.id | CRAFT_FILE__{user_id}__{sha256(path)[:16]} | 确定性 ID,前缀编码所有权 |
Document.link | file_store 的file_id | 检索时的指针 |
Document.file_id | file_store 的file_id(镜像) | 模式层指针 |
Document.doc_metadata | {file_path, file_size, mime_type, is_directory, sync_disabled, file_store_id} | 展示与过滤 |
Connector | 共享的 "User Library" 连接器,DocumentSource.CRAFT_FILE | 复用既有索引基础设施 |
Credential | 每个用户一个,credential_json={} | 与 Connector 配对形成所有权 |
FileOrigin.USER_FILE是 file_store 侧的文件来源标签,用于标记用户上传文件。
2.1 确定性文档 ID 的生成
build_document_id(user_id, path)(db/user_library.py)实现如下:
def build_document_id(user_id: UUID, path: str) -> str: """Deterministic document ID for a user library path.""" path_hash = hashlib.sha256(path.encode()).hexdigest()[:16] return f"CRAFT_FILE__{user_id}__{path_hash}"对path取 SHA-256 摘要的前 16 个十六进制字符,拼接为CRAFT_FILE__{user_id}__{hash}。这意味着:
- 同一用户上传同一路径的文件,doc_id 完全确定,天然幂等(重复上传即覆盖更新);
- 用户 ID 作为前缀,所有权信息直接编码在 ID 中,为第 5 节的所有权检查提供了前置条件。
2.2 元数据字段
doc_metadata中的is_directory用于区分真实文件与虚拟目录(zip 解压时创建的目录记录);sync_disabled用于控制某个文件/目录是否参与同步;file_store_id与link互为镜像,均在检索时指向同一个 file_store blob。
semantic_identifier被设置为user_library{file_path}的形式(USER_LIBRARY_SOURCE_DIR = "user_library",见 configs.py),文件树接口正是从它解析出路径与名称。
3. 同步流水线(Sync Pipeline)
3.1 挂载路径
沙箱 pod 内的挂载路径为/workspace/managed/user_library/。与 skills 相同,这里使用原子符号链接交换(atomic symlink swap):智能体始终看到一条稳定路径,而 push daemon 在底层轮换带版本号的目录。同步模块中该常量定义在 sandbox/user_library.py:
USER_LIBRARY_MOUNT_PATH = "/workspace/managed/user_library"3.2 触发点
同步由以下事件驱动,全部为同步执行(无 Celery 异步任务),与 skills 模式保持一致:
| 触发点 | 调用 |
|---|---|
| 会话创建(冷启动) | hydrate_user_library(sandbox_id, user_id, db_session) |
| 恢复既有会话 | hydrate_user_library(...)(再次水合以获取最新状态) |
| 文件上传(单文件或 zip) | sync_user_library_to_active_sandboxes(user_id, db_session) |
| 文件删除 | sync_user_library_to_active_sandboxes(...) |
| 切换 sync_disabled | sync_user_library_to_active_sandboxes(...) |
冷启动路径在 session/manager.py 的create_session__no_commit与get_or_create_empty_session中都会调用_hydrate_user_library,确保新建会话的沙箱在报告就绪前就已完成托管内容(skills + 用户文件库)的推送。这一点在 sandbox_lifecycle.py 的build_managed_content_payload中体现:library_files(由build_user_library_fileset构建)与skills_files一同打包进ManagedContentPayload,再由push_managed_content在无打开事务的情况下执行外部推送(sandbox_lifecycle.py)。
3.3 数据流
build_user_library_fileset(user_id, db_session) 1. list_user_files() → CRAFT_FILE Document rows for user 2. Filter: skip is_directory=True, skip sync_disabled=True 3. For each remaining doc: file_store.read_file(doc.link) 4. Strip leading "/" from file_path (push daemon rejects absolute paths) 5. Return FileSet { relative_path: bytes } │ ▼ sandbox_manager.push_to_sandbox / push_to_sandboxes mount_path="/workspace/managed/user_library" → tar.gz, Ed25519-sign, POST to push daemon, atomic swapbuild_user_library_fileset(sandbox/user_library.py)的关键细节:
- 遍历用户的全部
CRAFT_FILE文档,跳过is_directory=True与sync_disabled=True的记录; - 通过
doc.link从 file_store 读取原始字节; - 剥离路径前导
/——push daemon 拒绝绝对路径,file_path.lstrip("/")保证 FileSet 中的 key 始终是相对路径; - 单文件读取失败仅记录 warning 并跳过,不阻塞整个 fileset 构建。
推送端由 sandbox/base.py 的push_to_sandbox/push_to_sandboxes提供:单沙箱推送带重试与超时(默认 30s),FatalWriteError立即失败,瞬时错误按退避重试;多沙箱推送则并行执行并汇总PushResult(含失败明细,逐条记录日志)。底层 K8s 实现通过 sidecar push daemon 接收Ed25519 签名的 tar.gz 包并做原子交换(签名公钥通过ONYX_SANDBOX_PUSH_PUBLIC_KEY环境变量注入,见 contract.py);Docker 实现则退化为docker exec tar -x,同样保证原子落盘语义(docker_sandbox_manager.py)。
空 fileset 也会照常推送——这正是清理机制的核心:通过原子交换,旧目录会被整体替换为空目录,从而清掉沙箱中已被删除的陈旧文件。
3.4 上传配额与安全校验
上传接口(api.py)在写入前会做多层校验,全部使用 configs.py 中定义的配置:
| 配置项 | 环境变量 | 默认值 | 说明 |
|---|---|---|---|
USER_LIBRARY_MAX_FILE_SIZE_BYTES | USER_LIBRARY_MAX_FILE_SIZE_MB | 500 MB | 单文件大小上限 |
USER_LIBRARY_MAX_TOTAL_SIZE_BYTES | USER_LIBRARY_MAX_TOTAL_SIZE_GB | 10 GB | 单用户存储总量上限 |
USER_LIBRARY_MAX_FILES_PER_UPLOAD | USER_LIBRARY_MAX_FILES_PER_UPLOAD | 100 | 单次上传文件数量上限 |
USER_LIBRARY_CONNECTOR_NAME | — | "User Library" | 共享连接器名称 |
USER_LIBRARY_CREDENTIAL_NAME | — | "User Library Credential" | 每用户凭据名称 |
USER_LIBRARY_SOURCE_DIR | — | "user_library" | 语义标识前缀 |
上传接口还包含:
- PDF 内嵌图片数量限制:调用
count_pdf_embedded_images统计单文件与单批次内嵌图片数,分别对照MAX_EMBEDDED_IMAGES_PER_FILE与MAX_EMBEDDED_IMAGES_PER_UPLOAD(api.py); - zip 炸弹防护:
_validate_zip_contents在解压前就依据 zip 目录项声明的file_size总和判断是否超配额(api.py); - 路径清洗:
_sanitize_path剥离..、.与非白名单字符([^a-zA-Z0-9\-_. ]),杜绝路径穿越(api.py); - zip 解压时跳过
__MACOSX与隐藏文件条目。
4. DB 层 API
db/user_library.py 是全部文件存储与文档管理的归属层:
| 函数 | 用途 |
|---|---|
get_or_create_craft_connector(db, user) | 返回(connector_id, credential_id),幂等 |
get_user_storage_bytes(db, user_id) | SQL 聚合存储用量(配额校验) |
build_document_id(user_id, path) | 确定性 doc_id |
list_user_files(db, user_id) | 该用户的全部 CRAFT_FILE 文档 |
fetch_user_file_for_user(db, doc_id, user_id) | 查询 + 所有权检查;未命中抛OnyxError(NOT_FOUND) |
store_user_file(db, ..., file_path, content, mime_type) | 写入 file_store + upsert Document,返回(doc_id, file_id, old_blob_id_to_delete) |
cleanup_old_blobs(blob_ids) | 删除被取代的旧 blob |
create_directory_record(db, user_id, connector_id, credential_id, dir_path) | 虚拟目录文档(无 file_store 对象) |
set_sync_disabled(db, user_id, doc, sync_disabled) | 切换文件/目录同步开关(目录递归到子项) |
delete_user_file(db, doc) | 删除 file_store blob + Document 行 |
HTTP 层直接调用这些 helper,端点本身不含业务逻辑。
4.1 幂等的连接器/凭据配对
get_or_create_craft_connector(db/user_library.py)保证每个用户至多一条 User Library cc_pair:
- 先查用户既有 cc_pair 中是否存在
CRAFT_FILE来源且属于该用户的记录,命中直接返回; - 否则查找共享的 "User Library" 连接器(按名称匹配,不存在则创建,
InputType.LOAD_STATE、connector_specific_config={"disabled_paths": []}); - 复用该用户的 "User Library Credential"(不存在则以
credential_json={}创建); - 以
ProcessingMode.RAW_BINARY(原始二进制,不做文本抽取)和AccessType.PRIVATE关联二者,然后提交事务。
4.2 存储、覆盖与旧 blob 清理的顺序契约
store_user_file(db/user_library.py)返回三个值:doc_id、新的file_id、以及被取代的old_blob_id_to_delete。这里存在一个严格的调用顺序约束:
- 新 blob 先落盘:先调用
file_store.save_file()写入新内容,再 upsert Document; - 旧 blob 后删除:调用方必须在最终 DB 提交之后再调用
cleanup_old_blobs。如果在提交前删除,一旦后续提交失败,事务回滚会让文档重新指向已被删除的 blob,造成永久数据丢失。
API 层的实际用法(api.py):
db_session.commit() cleanup_old_blobs(stale_blobs)cleanup_old_blobs对每个 blob 调用file_store.delete_file(error_on_missing=False),删除失败仅记录 warning,不阻断主流程。
4.3 配额聚合查询
get_user_storage_bytes(db/user_library.py)通过DocumentByConnectorCredentialPair → ConnectorCredentialPair → Connector三级 join,过滤Connector.source == DocumentSource.CRAFT_FILE且creator_id == user_id,再对doc_metadata["file_size"]求和(排除is_directory=True的记录)。doc_metadata中的 JSON 字段通过 PostgreSQL 的as_string()+cast(... Integer)完成类型转换。
5. 所有权模型:ID 前缀即权限边界
所有权被直接编码在文档 ID 前缀中:CRAFT_FILE__{user_id}__{hash}。fetch_user_file_for_user(db/user_library.py)在任何 DB 查询之前先校验前缀:
if not doc_id.startswith(f"CRAFT_FILE__{user_id}__"): raise OnyxError(OnyxErrorCode.NOT_FOUND, "File not found") doc = get_document(doc_id, db_session) if doc is None: raise OnyxError(OnyxErrorCode.NOT_FOUND, "File not found") return doc关键设计决策:其他用户的 doc_id 返回NOT_FOUND(404)而非403 Forbidden——因为 doc_id 本身就是可猜测的确定性值,若返回 403 会泄露"该文件存在"的信息(存在性侧信道)。所有跨用户访问一律表现为"文件不存在",从根上杜绝了存在性泄露。上传、删除、目录创建等接口均通过require_permission(Permission.BASIC_ACCESS)鉴权,并最终落到上述前缀校验。
6. 相关文件清单
新增文件:
- backend/onyx/server/features/build/sandbox/user_library.py — 同步模块(FileSet 构建 + 推送编排)
- backend/tests/integration/tests/craft/k8s/test_user_library_sync.py — 面向 Helm 安装的 kind 集群的 API 驱动 k8s 集成测试(真实 API、web_server、Celery、支撑服务、sandbox-proxy 与 sandbox pod)
- backend/tests/external_dependency_unit/craft/test_user_library_fileset.py — fileset/同步 helper 的直接单测
修改文件:
- backend/onyx/server/features/build/db/user_library.py — 新增全部 CRUD/存储 helper
- backend/onyx/server/features/build/user_library/api.py — 精简为调用 DB helper;移除
PersistentDocumentWriter用法;HTTPException→OnyxError - backend/onyx/server/features/build/session/manager.py — 在
create_session__no_commit与get_or_create_empty_session中调用_hydrate_user_library - backend/onyx/skills/push.py — 逐失败日志(与 user_library 推送日志一致)
- backend/tests/integration/tests/craft/k8s/k8s_fixtures.py —
SandboxHandle.provision_api_user返回WorkspaceProxy,消除重复的_provision_with_statushelper
删除文件:
- backend/onyx/server/features/build/indexing/persistent_document_writer.py
- backend/tests/external_dependency_unit/craft/test_persistent_document_writer.py
7. 测试策略:分层覆盖
测试遵循 Craft 的分层测试策略:需要真实沙箱行为的场景走 API 驱动集成测试,纯文件集逻辑则用更窄的单元测试。
K8s 集成测试(test_user_library_sync.py)
使用 Helm 安装的 kind chart 提供的真实 Postgres、Redis、MinIO,以及真实的 api_server、web_server、Celery worker、sandbox-proxy 与 sandbox pod。测试通过部署后的/build/user-libraryAPI 播种数据,然后在沙箱 pod 内断言文件确实落盘:
test_upload_api_syncs_file_to_running_sandbox— 上传 API 把文件推入运行中的沙箱;test_upload_zip_api_syncs_nested_file— zip 上传推送嵌套路径(断言路径为bundle/folder-xxx/real_file.csv);test_session_workspace_links_user_library_after_api_upload— 会话符号链接暴露已同步文件(通过workspace / "managed" / "user_library" / file_path断言);test_delete_api_removes_file_from_running_sandbox— 删除 API 从运行中沙箱移除文件。
这些测试在SANDBOX_BACKEND != kubernetes时通过pytest.mark.skipif跳过,仅在专用 K8s CI 任务中运行。
外部依赖单元测试(test_user_library_fileset.py)
用StubSandboxManager直测不需要 Kubernetes 的纯逻辑:
test_sync_disabled_files_excluded—build_user_library_fileset排除 sync-disabled 文件;test_directories_excluded_from_fileset— 排除目录记录;test_managed_content_push_delivers_current_library_fileset— 冷启动水合路径:build_managed_content_payload携带当前库文件集,push_managed_content推送到正确挂载点;test_sync_user_library_pushes_to_running_sandbox— 仅向 RUNNING 状态的沙箱行推送。
既有 HTTP 集成测试(test_user_library_api.py)
覆盖跨用户 404、sync 开关、删除、上传上限等 HTTP 层行为,属于测试重构分支上既有的测试。
底层推送契约
test_user_library_sync.py覆盖 API 触发的同步;更底层的 push daemon 契约(签名 tarball、原子交换)由更广泛的 Craft k8s 集成套件覆盖,User Library 同步复用的是未改动的write_files_to_sandbox(),因此无需重复测试。
8. 设计取舍与明确不在范围内的事项
以下是当前设计刻意划定的边界(来自原文档,均可在源码中得到印证):
- 无按会话的文件选择:用户所有未禁用同步的文件会进入其全部活跃沙箱,粒度是"用户级"而非"会话级";
- 不做大文件流式传输:Office 文档通常为 MB 级,100 MiB 的包体上限足够;
- 不新增数据库表或迁移:完全复用 Document + doc_metadata 模型;
- 不改动 push daemon 与解压逻辑:仅复用既有的签名 tarball + 原子交换通道。
同步全程为同步调用,不引入 Celery,与 skills 模式保持一致;推送失败仅记录日志,不影响上传/删除接口的成功返回,属于尽力而为(best-effort)的同步语义。理解这套"确定性文档 ID + 共享连接器 + FileSet 推送 + 原子交换"的组合,也就掌握了 Craft 托管内容同步(skills、MCP 配置、用户文件库)的统一心智模型。
【免费下载链接】danswerOpen Source AI Platform - AI Chat with advanced features that works with every LLM项目地址: https://gitcode.com/GitHub_Trending/da/danswer
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考