mirror of
https://github.com/bytedance/deer-flow.git
synced 2026-09-20 19:46:16 +00:00
* feat(projects): Projects MVP Phase 2 — instructions, document shelf, promotion, trash Implements docs/superpowers/specs/2026-09-12-projects-mvp-phase2-design.md (issue #5160, tracker #5129) in the slice order of the spec's §16. Slices: - A: ProjectsConfig + write-time 422 UTF-8 byte cap; PROJECT_CONTEXT_KEY admission pinning (both server-owned sets + worker hoist); latest-only request-scoped <project> block via DynamicContextMiddleware wrap_model_call/awrap_model_call (idempotent reassembly, reserved ID prefix + marker + provenance, never persisted); journal audit fingerprints; Instructions tab. - B: ProjectDocumentRow + migration 0023; ProjectDocumentRepository with locked check-and-set; hash-qualified immutable shelf storage with Paths helpers; upload/list/content/delete-to-trash routes; project delete trashes the shelf in-transaction; request-scoped bounded <documents> index with honest count/shown + actionable overflow note; list_project_documents/read_project_document tools registered only on pinned runs; PAT allowlist + drift guards; blocking-IO anchors. - C: shared thread-upload ingestion service (uploads router refactored to parity); POST from-thread with provenance; attach-to-thread with lock-staged copy (archived source allowed); read-only thread-files view with per-group truncation reporting. - D: restore (restored/merged/not_found/no_target/content_missing; no file moves), purge (continuous row lock across unlink/delete/commit, retryable on FS errors), retention sweep (lazy + startup, 24h orphan guard, row-side reconciliation never deletes). - E: Documents tab (shelf + conversation-files browser, provenance, archived banner, content-missing rows), /workspace/trash route, sidebar entry, composer attach handoff, i18n (en-US/zh-CN), e2e mocks + specs. Review hardening folded in (10 rounds, all with tests): - force active shelf content (HTML/XML family) to download; nosniff on artifact + content responses; unified unsandboxed-iframe PDF preview (fixes the pre-existing Chromium sandbox blank in the artifact viewer) - scope document trash to the URL project under the document lock - atomic no-overwrite filename reservation for ALL ingestion (seeded claims + os.link commit with suffix retry; same-name re-upload now unique-names instead of replacing); hidden staging only, no visible placeholders; lease cleanup on setup failure - serialize conversion under the document lock with post-lock active revalidation; drain locked filesystem work on cancellation; preserve bytes when an insert's commit state is uncertain (including trashed rows) - original-integrity checks before serving text or cached conversions; content_missing surfaced in list responses (UI reads the flag, no 409-probe); downloads always serve original bytes - bounded streaming document reads with cached char counts; shelf limits declared in middleware release identity - thread-root confinement for from-thread sources; config fallback rejects fractional/infinite values; composer counts staged attachments; pending attachments persist until submission or removal; in-flight instruction/rename edits survive save refetches; shelf and trash pagination; conversation-file and thread-files pages stay subscribed to refetches Docs: README/README_zh, backend API.md/ARCHITECTURE.md, AGENTS.md contracts, config.example.yaml projects block. Review follow-ups (head b4807477 → this revision): - The trash retention sweep is split so repeated lazy triggers stay bounded: the indexed expiry purge still runs on every trigger (GET /api/trash/documents, POST /api/trash/purge) while the O(all rows + all files) reconciliation is throttled to one run per user per 15 minutes (process-local, per-user window). The startup sweep now runs as a background task instead of blocking gateway readiness, and shutdown awaits it (bounded). - The export scrub (stripInternalMarkers) is fence- and indentation-aware like the render path, so a pasted, fenced <project>/<documents> snippet survives markdown export while real injected blocks (never fenced) are still removed. Fence regexes moved to a dependency-free leaf module to avoid the messages↔streamdown import cycle. - The artifact viewer's PDF iframe no longer carries an added title attribute (the upstream e2e contract locates it via :not([title])), and the upstream artifact-preview spec now pins the new contract: PDFs render unsandboxed, images keep sandbox="". * fix(projects): round-2 review — cancel an overrun trash sweep, restore the PDF frame title - Shutdown cancelled only the shield around the background startup sweep, so an all-users reconciliation that outlived the 5s budget kept walking rows and files while the document repo and DB engine were disposed underneath it. The wait now lives in `_shutdown_startup_trash_sweep`, which cancels the task and drains it before worker exit: the shield keeps the wait bounded, the cancel makes it final (CancelledError lands at the sweep's next await, and `_run_startup_trash_sweep` only catches `Exception`, so nothing swallows it). - The browser-preview iframe lost `title={getFileName(filepath)}` in the previous fix round, leaving the PDF frame without an accessible name while its siblings keep theirs. Restore it (WCAG frame titles), assert it in the DOM test, and anchor the e2e on `iframe[title="report.pdf"]` instead of `iframe:not([title])`. * fix(projects): round-3 review — report the sweep's late finish, not a phantom cancel `Task.cancel()` returns False when the sweep already finished inside the window between the deadline firing and the cancel, so the shutdown log claimed a cancellation that never happened. Branch on that outcome: the warning stays for a real cancel, a late finish is logged at info, and both paths still reap the task before worker exit. * fix(projects): round-4 review — make Empty trash delete what it confirms `POST /api/trash/purge` only ran the retention sweep, and the sweep's candidate selection is age-gated, so a freshly trashed document survived "Empty trash" even though the confirmation promises that every listed document is permanently deleted. With one trashed row the route answered `{"purged": 0}` and left it in place; `GET /api/trash/documents` sweeps expired rows before listing, so the visible rows were normally ineligible for the action by construction. Empty trash now drives `purge_all_trashed`: the caller's trashed rows (`list_all_trashed`, no age filter) each go through the same guarded, row-locked `purge` as the single-document delete — bytes first, then the row, in one transaction — so a row restored mid-flight is skipped instead of force-deleted, and an unlink failure rolls that row back and answers 500 with a retryable message. Retention expiry stays where it was: the sweep's `purge_candidates` is now the only age-gated selection, and the lazy retention sweep still runs on the listing and at startup. Tests: the router suite replaces the retention-gated expectation with the reviewer's repro (fresh row purged, bytes unlinked, shelf and other users' trash untouched, a failing unlink stays retryable and 500); a blocking-I/O anchor drives the new entry point through the offload; the mocked e2e covers the action end to end; a new real-backend spec performs it against the real gateway and re-reads `GET /api/trash/documents`. README, API, ARCHITECTURE and the phase-2 design docs (en+zh) state the age-independent contract.
391 lines
14 KiB
Python
391 lines
14 KiB
Python
"""Shared upload management logic.
|
|
|
|
Pure business logic — no FastAPI/HTTP dependencies.
|
|
Both Gateway and Client delegate to these functions.
|
|
"""
|
|
|
|
import errno
|
|
import logging
|
|
import os
|
|
import stat
|
|
from pathlib import Path
|
|
from urllib.parse import quote
|
|
|
|
from deerflow.config.paths import VIRTUAL_PATH_PREFIX, get_paths
|
|
from deerflow.runtime.user_context import get_effective_user_id
|
|
from deerflow.utils.thread_id import validate_thread_id
|
|
|
|
|
|
class PathTraversalError(ValueError):
|
|
"""Raised when a path escapes its allowed base directory."""
|
|
|
|
|
|
class UnsafeUploadPathError(ValueError):
|
|
"""Raised when an upload destination is not a safe regular file path."""
|
|
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
UPLOAD_STAGING_PREFIX = ".upload-"
|
|
UPLOAD_STAGING_SUFFIX = ".part"
|
|
|
|
_MAX_FILENAME_BYTES = 255
|
|
|
|
|
|
def get_uploads_dir(thread_id: str, *, user_id: str | None = None) -> Path:
|
|
"""Return the uploads directory path for a thread (no side effects)."""
|
|
validate_thread_id(thread_id)
|
|
return get_paths().sandbox_uploads_dir(thread_id, user_id=user_id or get_effective_user_id())
|
|
|
|
|
|
def ensure_uploads_dir(thread_id: str, *, user_id: str | None = None) -> Path:
|
|
"""Return the uploads directory for a thread, creating it if needed."""
|
|
base = get_uploads_dir(thread_id, user_id=user_id)
|
|
base.mkdir(parents=True, exist_ok=True)
|
|
return base
|
|
|
|
|
|
def normalize_filename(filename: str) -> str:
|
|
"""Sanitize a filename by extracting its basename.
|
|
|
|
Strips any directory components and rejects traversal patterns.
|
|
|
|
Args:
|
|
filename: Raw filename from user input (may contain path components).
|
|
|
|
Returns:
|
|
Safe filename (basename only).
|
|
|
|
Raises:
|
|
ValueError: If filename is empty or resolves to a traversal pattern.
|
|
"""
|
|
if not filename:
|
|
raise ValueError("Filename is empty")
|
|
safe = Path(filename).name
|
|
if not safe or safe in {".", ".."}:
|
|
raise ValueError(f"Filename is unsafe: {filename!r}")
|
|
# Reject backslashes — on Linux Path.name keeps them as literal chars,
|
|
# but they indicate a Windows-style path that should be stripped or rejected.
|
|
if "\\" in safe:
|
|
raise ValueError(f"Filename contains backslash: {filename!r}")
|
|
if len(safe.encode("utf-8")) > _MAX_FILENAME_BYTES:
|
|
raise ValueError(f"Filename too long: {len(safe)} chars")
|
|
return safe
|
|
|
|
|
|
def _fit_utf8_bytes(text: str, budget: int) -> str:
|
|
"""Truncate *text* to at most *budget* UTF-8 bytes without splitting a code point."""
|
|
encoded = text.encode("utf-8")
|
|
if len(encoded) <= budget:
|
|
return text
|
|
return encoded[:budget].decode("utf-8", errors="ignore")
|
|
|
|
|
|
def claim_unique_filename(name: str, seen: set[str]) -> str:
|
|
"""Generate a unique filename by appending ``_N`` suffix on collision.
|
|
|
|
Automatically adds the returned name to *seen* so callers don't need to.
|
|
|
|
The deduplicated name stays within the 255-byte filename limit that
|
|
:func:`normalize_filename` enforces: when appending ``_N`` (plus the
|
|
preserved extension) would exceed it, the stem is truncated on a UTF-8
|
|
boundary to make room. Otherwise a maximum-length upload that collides
|
|
would produce a name the filesystem (and a later ``normalize_filename``
|
|
call on the write path) rejects.
|
|
|
|
Args:
|
|
name: Candidate filename.
|
|
seen: Set of filenames already claimed (mutated in place).
|
|
|
|
Returns:
|
|
A filename not present in *seen* (already added to *seen*).
|
|
"""
|
|
if name not in seen:
|
|
seen.add(name)
|
|
return name
|
|
stem, suffix = Path(name).stem, Path(name).suffix
|
|
counter = 1
|
|
while True:
|
|
tag = f"_{counter}"
|
|
budget = _MAX_FILENAME_BYTES - len(tag.encode("utf-8")) - len(suffix.encode("utf-8"))
|
|
if budget < 1:
|
|
# Pathological suffix that leaves no room for a stem; keep the
|
|
# unique tag and fit the rest (stem + suffix tail) around it.
|
|
candidate = _fit_utf8_bytes(stem + suffix, _MAX_FILENAME_BYTES - len(tag.encode("utf-8"))) + tag
|
|
else:
|
|
candidate = f"{_fit_utf8_bytes(stem, budget)}{tag}{suffix}"
|
|
if candidate not in seen:
|
|
break
|
|
counter += 1
|
|
seen.add(candidate)
|
|
return candidate
|
|
|
|
|
|
def is_upload_staging_file(filename: str) -> bool:
|
|
"""Return whether *filename* is a transient Gateway upload staging file."""
|
|
return filename.startswith(UPLOAD_STAGING_PREFIX) and filename.endswith(UPLOAD_STAGING_SUFFIX)
|
|
|
|
|
|
def validate_path_traversal(path: Path, base: Path) -> None:
|
|
"""Verify that *path* is inside *base*.
|
|
|
|
Raises:
|
|
PathTraversalError: If a path traversal is detected.
|
|
"""
|
|
try:
|
|
path.resolve().relative_to(base.resolve())
|
|
except ValueError:
|
|
raise PathTraversalError("Path traversal detected") from None
|
|
|
|
|
|
def validate_upload_destination(base_dir: Path, filename: str) -> Path:
|
|
"""Validate an upload destination without mutating an existing file."""
|
|
safe_name = normalize_filename(filename)
|
|
dest = base_dir / safe_name
|
|
|
|
try:
|
|
st = os.lstat(dest)
|
|
except FileNotFoundError:
|
|
st = None
|
|
|
|
if st is not None and not stat.S_ISREG(st.st_mode):
|
|
raise UnsafeUploadPathError(f"Upload destination is not a regular file: {safe_name}")
|
|
if st is not None and st.st_nlink > 1:
|
|
raise UnsafeUploadPathError(f"Upload destination has multiple links: {safe_name}")
|
|
|
|
validate_path_traversal(dest, base_dir)
|
|
return dest
|
|
|
|
|
|
def _iter_upload_dirs(base_dir: Path):
|
|
yield from base_dir.glob("threads/*/user-data/uploads")
|
|
yield from base_dir.glob("users/*/threads/*/user-data/uploads")
|
|
|
|
|
|
def cleanup_stale_upload_staging_files(base_dir: Path | str | None = None) -> int:
|
|
"""Remove orphaned Gateway upload staging files left by a hard crash."""
|
|
root = Path(base_dir) if base_dir is not None else get_paths().base_dir
|
|
removed = 0
|
|
for uploads_dir in _iter_upload_dirs(root):
|
|
if not uploads_dir.is_dir():
|
|
continue
|
|
try:
|
|
with os.scandir(uploads_dir) as entries:
|
|
for entry in entries:
|
|
if not is_upload_staging_file(entry.name) or not entry.is_file(follow_symlinks=False):
|
|
continue
|
|
try:
|
|
os.unlink(entry.path)
|
|
removed += 1
|
|
except FileNotFoundError:
|
|
pass
|
|
except OSError:
|
|
logger.warning("Failed to remove stale upload staging file: %s", entry.path, exc_info=True)
|
|
except FileNotFoundError:
|
|
continue
|
|
except OSError:
|
|
logger.warning("Failed to scan uploads directory for stale staging files: %s", uploads_dir, exc_info=True)
|
|
return removed
|
|
|
|
|
|
def open_upload_file_no_symlink(base_dir: Path, filename: str) -> tuple[Path, object]:
|
|
"""Open an upload destination for safe streaming writes.
|
|
|
|
Upload directories may be mounted into local sandboxes. A sandbox process can
|
|
therefore leave a symlink at a future upload filename. Normal ``Path.write_bytes``
|
|
follows that link and can overwrite files outside the uploads directory with
|
|
gateway privileges. This helper rejects symlink destinations using ``O_NOFOLLOW``
|
|
on POSIX. On Windows (which lacks ``O_NOFOLLOW``), it uses dual ``lstat`` checks
|
|
and ``fstat`` validation after ``open()`` to reduce the TOCTOU window; this does
|
|
not eliminate all races but makes exploitation significantly harder. Path-traversal
|
|
validation prevents escapes from *base_dir* in both cases.
|
|
"""
|
|
safe_name = normalize_filename(filename)
|
|
dest = validate_upload_destination(base_dir, safe_name)
|
|
try:
|
|
st = os.lstat(dest)
|
|
except FileNotFoundError:
|
|
st = None
|
|
|
|
has_nofollow = hasattr(os, "O_NOFOLLOW")
|
|
|
|
if has_nofollow:
|
|
# POSIX: O_NOFOLLOW makes open() fail with ELOOP if dest is a symlink.
|
|
flags = os.O_WRONLY | os.O_CREAT | os.O_NOFOLLOW
|
|
if hasattr(os, "O_NONBLOCK"):
|
|
flags |= os.O_NONBLOCK
|
|
|
|
try:
|
|
fd = os.open(dest, flags, 0o600)
|
|
except OSError as exc:
|
|
if exc.errno in {errno.ELOOP, errno.EISDIR, errno.ENOTDIR, errno.ENXIO, errno.EAGAIN}:
|
|
raise UnsafeUploadPathError(f"Unsafe upload destination: {safe_name}") from exc
|
|
raise
|
|
|
|
try:
|
|
opened_stat = os.fstat(fd)
|
|
if not stat.S_ISREG(opened_stat.st_mode) or opened_stat.st_nlink != 1:
|
|
raise UnsafeUploadPathError(f"Upload destination is not an exclusive regular file: {safe_name}")
|
|
os.ftruncate(fd, 0)
|
|
fh = os.fdopen(fd, "wb")
|
|
fd = -1
|
|
finally:
|
|
if fd >= 0:
|
|
os.close(fd)
|
|
return dest, fh
|
|
|
|
# Windows: no O_NOFOLLOW available. Uses a second lstat immediately before open()
|
|
# to narrow the TOCTOU window, then fstat after open() as a further defence.
|
|
# Note: a narrow race window remains between the pre-open lstat and open(); the
|
|
# path-traversal check mitigates escapes from base_dir but cannot prevent an
|
|
# attacker who can atomically replace dest with a symlink after the check.
|
|
if st is not None and st.st_nlink > 1:
|
|
raise UnsafeUploadPathError(f"Upload destination has multiple links: {safe_name}")
|
|
|
|
flags = os.O_WRONLY | os.O_CREAT
|
|
if hasattr(os, "O_BINARY"):
|
|
flags |= os.O_BINARY
|
|
|
|
try:
|
|
pre_open_st = os.lstat(dest)
|
|
except FileNotFoundError:
|
|
pre_open_st = None
|
|
|
|
if pre_open_st is not None and not stat.S_ISREG(pre_open_st.st_mode):
|
|
raise UnsafeUploadPathError(f"Upload destination is not a regular file: {safe_name}")
|
|
if pre_open_st is not None and pre_open_st.st_nlink > 1:
|
|
raise UnsafeUploadPathError(f"Upload destination has multiple links: {safe_name}")
|
|
|
|
try:
|
|
fd = os.open(dest, flags, 0o600)
|
|
except OSError as exc:
|
|
if exc.errno in {errno.EISDIR, errno.ENOTDIR, errno.ENXIO, errno.EAGAIN}:
|
|
raise UnsafeUploadPathError(f"Unsafe upload destination: {safe_name}") from exc
|
|
raise
|
|
|
|
try:
|
|
opened_stat = os.fstat(fd)
|
|
if not stat.S_ISREG(opened_stat.st_mode) or opened_stat.st_nlink > 1:
|
|
raise UnsafeUploadPathError(f"Upload destination is not an exclusive regular file: {safe_name}")
|
|
os.ftruncate(fd, 0)
|
|
fh = os.fdopen(fd, "wb")
|
|
fd = -1
|
|
finally:
|
|
if fd >= 0:
|
|
os.close(fd)
|
|
return dest, fh
|
|
|
|
|
|
def write_upload_file_no_symlink(base_dir: Path, filename: str, data: bytes) -> Path:
|
|
"""Write upload bytes without following a pre-existing destination symlink."""
|
|
dest, fh = open_upload_file_no_symlink(base_dir, filename)
|
|
with fh:
|
|
fh.write(data)
|
|
return dest
|
|
|
|
|
|
def list_files_in_dir(directory: Path) -> dict:
|
|
"""List files (not directories) in *directory*.
|
|
|
|
Args:
|
|
directory: Directory to scan.
|
|
|
|
Returns:
|
|
Dict with "files" list (sorted by name) and "count".
|
|
Each file entry has ``size`` as *int* (bytes). Call
|
|
:func:`enrich_file_listing` to add virtual / artifact URLs.
|
|
"""
|
|
if not directory.is_dir():
|
|
return {"files": [], "count": 0}
|
|
|
|
files = []
|
|
with os.scandir(directory) as entries:
|
|
for entry in sorted(entries, key=lambda e: e.name):
|
|
if is_upload_staging_file(entry.name):
|
|
continue
|
|
if not entry.is_file(follow_symlinks=False):
|
|
continue
|
|
st = entry.stat(follow_symlinks=False)
|
|
files.append(
|
|
{
|
|
"filename": entry.name,
|
|
"size": st.st_size,
|
|
"path": entry.path,
|
|
"extension": Path(entry.name).suffix,
|
|
"modified": st.st_mtime,
|
|
}
|
|
)
|
|
return {"files": files, "count": len(files)}
|
|
|
|
|
|
def delete_file_safe(base_dir: Path, filename: str, *, convertible_extensions: set[str] | None = None) -> dict:
|
|
"""Delete a file inside *base_dir* after path-traversal validation.
|
|
|
|
If *convertible_extensions* is provided and the file's extension matches,
|
|
the companion ``.md`` file is also removed (if it exists).
|
|
|
|
Args:
|
|
base_dir: Directory containing the file.
|
|
filename: Name of file to delete.
|
|
convertible_extensions: Lowercase extensions (e.g. ``{".pdf", ".docx"}``)
|
|
whose companion markdown should be cleaned up.
|
|
|
|
Returns:
|
|
Dict with success and message.
|
|
|
|
Raises:
|
|
FileNotFoundError: If the file does not exist.
|
|
PathTraversalError: If path traversal is detected.
|
|
"""
|
|
file_path = (base_dir / filename).resolve()
|
|
validate_path_traversal(file_path, base_dir)
|
|
|
|
if not file_path.is_file():
|
|
raise FileNotFoundError(f"File not found: {filename}")
|
|
|
|
file_path.unlink()
|
|
|
|
# Clean up companion markdown generated during upload conversion.
|
|
if convertible_extensions and file_path.suffix.lower() in convertible_extensions:
|
|
file_path.with_suffix(".md").unlink(missing_ok=True)
|
|
|
|
return {"success": True, "message": f"Deleted {filename}"}
|
|
|
|
|
|
def upload_artifact_url(thread_id: str, filename: str) -> str:
|
|
"""Build the artifact URL for a file in a thread's uploads directory.
|
|
|
|
*filename* is percent-encoded so that spaces, ``#``, ``?`` etc. are safe.
|
|
"""
|
|
return f"/api/threads/{thread_id}/artifacts{VIRTUAL_PATH_PREFIX}/uploads/{quote(filename, safe='')}"
|
|
|
|
|
|
def upload_virtual_path(filename: str) -> str:
|
|
"""Build the virtual path for a file in the uploads directory."""
|
|
return f"{VIRTUAL_PATH_PREFIX}/uploads/{filename}"
|
|
|
|
|
|
def output_artifact_url(thread_id: str, filename: str) -> str:
|
|
"""Build the artifact URL for a file in a thread's outputs directory.
|
|
|
|
*filename* is percent-encoded so that spaces, ``#``, ``?`` etc. are safe.
|
|
"""
|
|
return f"/api/threads/{thread_id}/artifacts{VIRTUAL_PATH_PREFIX}/outputs/{quote(filename, safe='')}"
|
|
|
|
|
|
def output_virtual_path(filename: str) -> str:
|
|
"""Build the virtual path for a file in the outputs directory."""
|
|
return f"{VIRTUAL_PATH_PREFIX}/outputs/{filename}"
|
|
|
|
|
|
def enrich_file_listing(result: dict, thread_id: str) -> dict:
|
|
"""Add virtual paths and artifact URLs on a listing result.
|
|
|
|
Mutates *result* in place and returns it for convenience.
|
|
"""
|
|
for f in result["files"]:
|
|
filename = f["filename"]
|
|
f["virtual_path"] = upload_virtual_path(filename)
|
|
f["artifact_url"] = upload_artifact_url(thread_id, filename)
|
|
return result
|