fix: preserve embedded client uploads

This commit is contained in:
hetaoBackend 2026-08-06 10:18:37 +08:00
parent 09d3fc04cd
commit 0d06a37773
2 changed files with 111 additions and 66 deletions

View File

@ -21,7 +21,6 @@ import copy
import logging
import mimetypes
import os
import shutil
import uuid
from collections.abc import Generator, Mapping, Sequence
from dataclasses import dataclass, field
@ -62,13 +61,15 @@ from deerflow.skills.storage import get_or_new_user_skill_storage
from deerflow.tools.builtins.tool_search import assemble_deferred_tools, build_mcp_routing_middleware, get_mcp_routing_hints_prompt_section
from deerflow.trace_context import DEERFLOW_TRACE_METADATA_KEY, generate_trace_id, get_current_trace_id, reset_current_trace_id, set_current_trace_id
from deerflow.tracing import build_tracing_callbacks, inject_langfuse_metadata
from deerflow.uploads.conversion import convert_uploaded_file_to_markdown
from deerflow.uploads.layout import artifact_url_for_virtual_path, conversion_virtual_path
from deerflow.uploads.manager import (
claim_unique_filename,
delete_file_safe,
enrich_file_listing,
ensure_uploads_dir,
get_uploads_dir,
list_files_in_dir,
publish_upload_copy,
upload_artifact_url,
upload_virtual_path,
)
@ -1495,11 +1496,10 @@ class DeerFlowClient:
ValueError: If any supplied path exists but is not a regular file.
"""
validate_thread_id(thread_id)
from deerflow.utils.file_conversion import CONVERTIBLE_EXTENSIONS, convert_file_to_markdown
from deerflow.utils.file_conversion import CONVERTIBLE_EXTENSIONS
# Validate all files upfront to avoid partial uploads.
resolved_files = []
seen_names: set[str] = set()
resolved_files: list[Path] = []
has_convertible_file = False
for f in files:
p = Path(f)
@ -1507,8 +1507,7 @@ class DeerFlowClient:
raise FileNotFoundError(f"File not found: {f}")
if not p.is_file():
raise ValueError(f"Path is not a file: {f}")
dest_name = claim_unique_filename(p.name, seen_names)
resolved_files.append((p, dest_name))
resolved_files.append(p)
if not has_convertible_file and p.suffix.lower() in CONVERTIBLE_EXTENSIONS:
has_convertible_file = True
@ -1528,13 +1527,13 @@ class DeerFlowClient:
# creating a new ThreadPoolExecutor per converted file.
conversion_pool = concurrent.futures.ThreadPoolExecutor(max_workers=1)
def _convert_in_thread(path: Path, output_path: Path | None = None):
return asyncio.run(convert_file_to_markdown(path, output_path=output_path))
def _convert_in_thread(path: Path):
return asyncio.run(convert_uploaded_file_to_markdown(path))
try:
for src_path, dest_name in resolved_files:
dest = uploads_dir / dest_name
shutil.copy2(src_path, dest)
for src_path in resolved_files:
dest = publish_upload_copy(uploads_dir, src_path.name, src_path)
dest_name = dest.name
info: dict[str, Any] = {
"filename": dest_name,
@ -1547,17 +1546,11 @@ class DeerFlowClient:
info["original_filename"] = src_path.name
if src_path.suffix.lower() in CONVERTIBLE_EXTENSIONS:
# Reserve companion .md name before convert so two stems
# that collapse to the same .md (or a prior .md upload)
# cannot silently overwrite each other.
provisional_md_name = Path(dest_name).with_suffix(".md").name
unique_md_name = claim_unique_filename(provisional_md_name, seen_names)
md_output = dest.with_name(unique_md_name)
try:
if conversion_pool is not None:
md_path = conversion_pool.submit(_convert_in_thread, dest, md_output).result()
md_path = conversion_pool.submit(_convert_in_thread, dest).result()
else:
md_path = asyncio.run(convert_file_to_markdown(dest, output_path=md_output))
md_path = asyncio.run(convert_uploaded_file_to_markdown(dest))
except Exception:
logger.warning(
"Failed to convert %s to markdown",
@ -1567,15 +1560,11 @@ class DeerFlowClient:
md_path = None
if md_path is not None:
md_virtual_path = conversion_virtual_path(dest_name)
info["markdown_file"] = md_path.name
info["markdown_path"] = str(uploads_dir / md_path.name)
info["markdown_virtual_path"] = upload_virtual_path(md_path.name)
info["markdown_artifact_url"] = upload_artifact_url(thread_id, md_path.name)
else:
# Conversion failed and wrote nothing, so release the
# claim; holding it would rename a later same-stem
# upload against a name nothing occupies.
seen_names.discard(unique_md_name)
info["markdown_path"] = str(md_path)
info["markdown_virtual_path"] = md_virtual_path
info["markdown_artifact_url"] = artifact_url_for_virtual_path(thread_id, md_virtual_path)
uploaded_files.append(info)
finally:
@ -1619,10 +1608,8 @@ class DeerFlowClient:
PermissionError: If path traversal is detected.
"""
validate_thread_id(thread_id)
from deerflow.utils.file_conversion import CONVERTIBLE_EXTENSIONS
uploads_dir = get_uploads_dir(thread_id)
return delete_file_safe(uploads_dir, filename, convertible_extensions=CONVERTIBLE_EXTENSIONS)
return delete_file_safe(uploads_dir, filename)
# ------------------------------------------------------------------
# Public API — artifacts

View File

@ -28,7 +28,8 @@ from deerflow.config.extensions_config import ExtensionsConfig, McpServerConfig
from deerflow.config.paths import Paths
from deerflow.skills.types import SkillCategory
from deerflow.tools.mcp_metadata import tag_mcp_tool
from deerflow.uploads.manager import PathTraversalError
from deerflow.uploads.layout import conversion_path_for_upload
from deerflow.uploads.manager import PathTraversalError, UnsafeUploadPathError
# ---------------------------------------------------------------------------
# Fixtures
@ -2143,6 +2144,52 @@ class TestUploads:
assert "message" in result
assert (uploads_dir / "test.txt").exists()
def test_upload_files_across_calls_never_overwrite(self, client, tmp_path):
uploads_dir = tmp_path / "user-data" / "uploads"
uploads_dir.mkdir(parents=True)
first_dir = tmp_path / "first"
second_dir = tmp_path / "second"
first_dir.mkdir()
second_dir.mkdir()
first = first_dir / "same.txt"
second = second_dir / "same.txt"
first.write_bytes(b"first")
second.write_bytes(b"second")
with (
patch("deerflow.client.ensure_uploads_dir", return_value=uploads_dir),
patch("deerflow.client.get_uploads_dir", return_value=uploads_dir),
):
first_result = client.upload_files("thread-1", [first])
second_result = client.upload_files("thread-1", [second])
assert first_result["files"][0]["filename"] == "same.txt"
assert second_result["files"][0]["filename"] == "same_1.txt"
assert (uploads_dir / "same.txt").read_bytes() == b"first"
assert (uploads_dir / "same_1.txt").read_bytes() == b"second"
def test_concurrent_client_uploads_preserve_all_payloads(self, client, tmp_path):
uploads_dir = tmp_path / "user-data" / "uploads"
uploads_dir.mkdir(parents=True)
sources = []
for index in range(8):
source_dir = tmp_path / f"source-{index}"
source_dir.mkdir()
source = source_dir / "same.bin"
source.write_bytes(f"payload-{index}".encode())
sources.append(source)
with (
patch("deerflow.client.ensure_uploads_dir", return_value=uploads_dir),
patch("deerflow.client.get_uploads_dir", return_value=uploads_dir),
concurrent.futures.ThreadPoolExecutor(max_workers=len(sources)) as pool,
):
results = list(pool.map(lambda source: client.upload_files("thread-1", [source]), sources))
names = [result["files"][0]["filename"] for result in results]
assert len(set(names)) == len(sources)
assert {(uploads_dir / name).read_bytes() for name in names} == {source.read_bytes() for source in sources}
def test_upload_files_not_found(self, client):
with pytest.raises(FileNotFoundError):
client.upload_files("thread-1", ["/nonexistent/file.txt"])
@ -2166,8 +2213,9 @@ class TestUploads:
created_executors = []
real_executor_cls = concurrent.futures.ThreadPoolExecutor
async def fake_convert(path: Path, output_path: Path | None = None) -> Path:
md_path = output_path if output_path is not None else path.with_suffix(".md")
async def fake_convert(path: Path) -> Path:
md_path = conversion_path_for_upload(path)
md_path.parent.mkdir(parents=True, exist_ok=True)
md_path.write_text(f"converted {path.name}")
return md_path
@ -2192,7 +2240,7 @@ class TestUploads:
patch("deerflow.client.get_uploads_dir", return_value=uploads_dir),
patch("deerflow.client.ensure_uploads_dir", return_value=uploads_dir),
patch("deerflow.utils.file_conversion.CONVERTIBLE_EXTENSIONS", {".pdf"}),
patch("deerflow.utils.file_conversion.convert_file_to_markdown", side_effect=fake_convert),
patch("deerflow.client.convert_uploaded_file_to_markdown", side_effect=fake_convert),
patch("concurrent.futures.ThreadPoolExecutor", FakeExecutor),
):
result = asyncio.run(call_upload())
@ -2202,8 +2250,8 @@ class TestUploads:
assert len(created_executors) == 1
assert created_executors[0].max_workers == 1
assert created_executors[0].shutdown_calls == [True]
assert result["files"][0]["markdown_file"] == "first.md"
assert result["files"][1]["markdown_file"] == "second.md"
assert result["files"][0]["markdown_file"] == "first.pdf.md"
assert result["files"][1]["markdown_file"] == "second.pdf.md"
def test_upload_files_converted_markdown_uses_unique_names_on_stem_collision(self, client):
"""Companion .md from convert must not clobber another same-stem companion."""
@ -2217,8 +2265,12 @@ class TestUploads:
docx.write_bytes(b"DOCX")
pdf.write_bytes(b"PDF")
async def fake_convert(path: Path, output_path: Path | None = None) -> Path:
md_path = output_path if output_path is not None else path.with_suffix(".md")
markdown = tmp_path / "a.md"
markdown.write_bytes(b"USER")
async def fake_convert(path: Path) -> Path:
md_path = conversion_path_for_upload(path)
md_path.parent.mkdir(parents=True, exist_ok=True)
md_path.write_text(f"FROM:{path.name}", encoding="utf-8")
return md_path
@ -2226,23 +2278,22 @@ class TestUploads:
patch("deerflow.client.get_uploads_dir", return_value=uploads_dir),
patch("deerflow.client.ensure_uploads_dir", return_value=uploads_dir),
patch("deerflow.utils.file_conversion.CONVERTIBLE_EXTENSIONS", {".docx", ".pdf"}),
patch("deerflow.utils.file_conversion.convert_file_to_markdown", side_effect=fake_convert),
patch("deerflow.client.convert_uploaded_file_to_markdown", side_effect=fake_convert),
):
result = client.upload_files("thread-1", [docx, pdf])
result = client.upload_files("thread-1", [docx, pdf, markdown])
assert result["success"] is True
assert result["files"][0]["markdown_file"] == "a.md"
assert result["files"][1]["markdown_file"] == "a_1.md"
assert (uploads_dir / "a.md").read_text(encoding="utf-8") == "FROM:a.docx"
assert (uploads_dir / "a_1.md").read_text(encoding="utf-8") == "FROM:a.pdf"
assert result["files"][0]["markdown_file"] == "a.docx.md"
assert result["files"][0]["markdown_virtual_path"] == "/mnt/user-data/.upload-conversions/a.docx.md"
assert result["files"][1]["markdown_file"] == "a.pdf.md"
assert result["files"][1]["markdown_artifact_url"] == ("/api/threads/thread-1/artifacts/mnt/user-data/.upload-conversions/a.pdf.md")
assert result["files"][2]["filename"] == "a.md"
assert (uploads_dir / "a.md").read_bytes() == b"USER"
assert conversion_path_for_upload(uploads_dir / "a.docx").read_text(encoding="utf-8") == "FROM:a.docx"
assert conversion_path_for_upload(uploads_dir / "a.pdf").read_text(encoding="utf-8") == "FROM:a.pdf"
def test_upload_files_failed_conversion_releases_the_claimed_markdown_name(self, client):
"""A conversion that writes nothing must not reserve stem.md against a later companion.
Destination names are claimed upfront, so a same-stem ``.md`` upload
always wins ``a.md``; the only reachable victim of a stale claim is the
next convertible's companion.
"""
def test_upload_files_failed_conversion_does_not_block_other_conversion(self, client):
"""A failed conversion does not affect another primary's owned asset."""
with tempfile.TemporaryDirectory() as tmp:
tmp_path = Path(tmp)
uploads_dir = tmp_path / "uploads"
@ -2253,10 +2304,11 @@ class TestUploads:
docx.write_bytes(b"DOCX")
pdf.write_bytes(b"PDF")
async def convert_failing_on_docx(path: Path, output_path: Path | None = None) -> Path | None:
async def convert_failing_on_docx(path: Path) -> Path | None:
if path.suffix.lower() == ".docx":
return None
md_path = output_path if output_path is not None else path.with_suffix(".md")
md_path = conversion_path_for_upload(path)
md_path.parent.mkdir(parents=True, exist_ok=True)
md_path.write_text(f"FROM:{path.name}", encoding="utf-8")
return md_path
@ -2264,15 +2316,14 @@ class TestUploads:
patch("deerflow.client.get_uploads_dir", return_value=uploads_dir),
patch("deerflow.client.ensure_uploads_dir", return_value=uploads_dir),
patch("deerflow.utils.file_conversion.CONVERTIBLE_EXTENSIONS", {".docx", ".pdf"}),
patch("deerflow.utils.file_conversion.convert_file_to_markdown", side_effect=convert_failing_on_docx),
patch("deerflow.client.convert_uploaded_file_to_markdown", side_effect=convert_failing_on_docx),
):
result = client.upload_files("thread-1", [docx, pdf])
assert result["success"] is True
assert result["files"][0].get("markdown_file") is None
assert result["files"][1]["markdown_file"] == "a.md"
assert (uploads_dir / "a.md").read_text(encoding="utf-8") == "FROM:a.pdf"
assert not (uploads_dir / "a_1.md").exists()
assert result["files"][1]["markdown_file"] == "a.pdf.md"
assert conversion_path_for_upload(uploads_dir / "a.pdf").read_text(encoding="utf-8") == "FROM:a.pdf"
def test_list_uploads(self, client):
with tempfile.TemporaryDirectory() as tmp:
@ -2295,15 +2346,24 @@ class TestUploads:
def test_delete_upload(self, client):
with tempfile.TemporaryDirectory() as tmp:
uploads_dir = Path(tmp)
(uploads_dir / "delete-me.txt").write_text("gone")
uploads_dir = Path(tmp) / "user-data" / "uploads"
uploads_dir.mkdir(parents=True)
primary = uploads_dir / "delete-me.txt"
primary.write_text("gone")
user_markdown = uploads_dir / "delete-me.md"
user_markdown.write_text("user")
conversion = conversion_path_for_upload(primary)
conversion.parent.mkdir()
conversion.write_text("generated")
with patch("deerflow.client.get_uploads_dir", return_value=uploads_dir), patch("deerflow.client.ensure_uploads_dir", return_value=uploads_dir):
result = client.delete_upload("thread-1", "delete-me.txt")
assert result["success"] is True
assert "delete-me.txt" in result["message"]
assert not (uploads_dir / "delete-me.txt").exists()
assert not primary.exists()
assert not conversion.exists()
assert user_markdown.read_text() == "user"
def test_delete_upload_not_found(self, client):
with tempfile.TemporaryDirectory() as tmp:
@ -3070,7 +3130,7 @@ class TestScenarioEdgeCases:
patch("deerflow.client.get_uploads_dir", return_value=uploads_dir),
patch("deerflow.client.ensure_uploads_dir", return_value=uploads_dir),
patch("deerflow.utils.file_conversion.CONVERTIBLE_EXTENSIONS", {".pdf"}),
patch("deerflow.utils.file_conversion.convert_file_to_markdown", side_effect=Exception("conversion failed")),
patch("deerflow.client.convert_uploaded_file_to_markdown", side_effect=Exception("conversion failed")),
):
result = client.upload_files("t-pdf-fail", [pdf_file])
@ -3784,7 +3844,7 @@ class TestSerializeMessage:
class TestUploadDeleteSymlink:
def test_delete_upload_symlink_outside_dir(self, client):
"""A symlink in uploads dir pointing outside is caught by path traversal check."""
"""A symlink in uploads dir is rejected without touching its target."""
with tempfile.TemporaryDirectory() as tmp:
uploads_dir = Path(tmp) / "uploads"
uploads_dir.mkdir()
@ -3803,9 +3863,7 @@ class TestUploadDeleteSymlink:
raise
with patch("deerflow.client.get_uploads_dir", return_value=uploads_dir), patch("deerflow.client.ensure_uploads_dir", return_value=uploads_dir):
# The resolved path of the symlink escapes uploads_dir,
# so path traversal check should catch it.
with pytest.raises(PathTraversalError):
with pytest.raises(UnsafeUploadPathError):
client.delete_upload("thread-1", "harmless.txt")
# The outside file must NOT have been deleted.