mirror of
https://github.com/pewdiepie-archdaemon/odysseus.git
synced 2026-10-06 06:52:20 +02:00
"Is this path inside that root" is asked in twenty places in this tree and answered twenty times by a locally written realpath/commonpath pair. Nine test files exist because nine call sites each needed their own proof. Each one is defensible alone; together they are the defect, because the boundary has no single definition and a site that gets a detail wrong is wrong by itself. src/path_confinement.py is that definition, and it settles the details the copies disagreed on. Both sides get canonicalized: comparing a realpath-ed candidate against a root that was only abspath-ed is the macOS /tmp -> /private/tmp mismatch that has already produced a false failure here, and canonicalizing one side is worse than canonicalizing neither. commonpath rather than startswith, because /a/bc begins with /a/b and is not inside it. A relative candidate joins the root rather than os.getcwd(), which is whatever directory the server happens to be running in. NUL and newline are refused with a reason instead of caught by a bare `except Exception` and reported as an ordinary escape. Eighteen call sites go through it now. It deliberately does not decide whether a path is sensitive -- that deny list answers "allowed" rather than "inside", and it stays with src/tool_execution, which owns it. The one commonpath left in the tree, in src/workspace_paths.py, stays: that function translates a host path into a container path, so canonicalizing either side would change the relative path it computes and break the mapping. It is not a confinement check. Two of those sites were weaker than the rest and are fixed rather than moved. The email attachment check used abspath, which folds `..` but does not resolve symlinks, so a symlink written into the extraction directory passed it and was then read through. The skill-reference guard compared a realpath-ed target against a raw dirname, so on a host where the skills tree is reached through a symlink the two sides never matched and the guard could not fire. The execution boundary had two separate holes. The workspace namespace bound /home and /mnt read-write. On the one platform where that namespace engages at all, a command inside it reaches outside the workspace and writes to the user's home directory -- measured by running this argv on a Linux host with working bubblewrap, not inferred from the source. Binding the user's whole home directory into a workspace-confinement namespace gives back most of what the namespace was for. Both are read-only now. The workspace is also bound writable at its real host path, not only at /workspace: BashTool's own /tmp redirect rewrites `/tmp/` to `<agent_cwd()>/.tmp/` before the namespace is built, so the command bwrap receives already names the real path, and those writes previously landed only because the workspace happened to sit under the writable /home. `namespaced or _replace_workspace_alias(...)` chose between a mount namespace and a regex with nothing in the result saying which one ran. The fallback rewrites the literal token /workspace in the command string, so a command that never mentions /workspace is untouched by it and runs on the host unrestricted -- which is every agent shell command on macOS. Both tools now ask containment.probe() instead of each deciding for itself, and every bash and python result carries a containment block naming the mechanism and stating whether the filesystem dimension actually held. Under enforcing mode the command is not run and the result says so. That block reports the filesystem dimension only, and says so in a reported_dimensions field. The probe knows this host could also give a process group and a real wall clock, but these two tools still assemble their own create_subprocess_* call and pass neither, so listing those dimensions would be exactly the false claim src/containment.py calls worse than an honest absence. probe() is new on src/containment.py: the same mechanism table and the same arithmetic as acquire(), stopping before the side effects. acquire() is the wrong shape for a decision -- it writes a durable grant record, and a record whose pid is never filled in and whose release() never runs is an entry a restart reaper keeps finding. CONTAINMENT_MODE stays report_only. Flipping it refuses every agent shell command on macOS and on any Linux host without bubblewrap, which is a product decision rather than a code one. Smaller things in the same area: the /tmp redirect's makedirs was unguarded, so a read-only workspace turned a command that merely mentioned `/tmp/` into an OSError traceback instead of a tool error; it degrades now. WORKSPACE_MOUNT moved to src/constants.py so the namespace and the path resolvers read one definition of the contract rather than two. The ".tmp" dirname got a constant, since it appeared in both tool paths. One generated artifact moved with it: website/configuration-reference.md pins the source line where each ODYSSEUS_* variable is read, and three of those shifted. Regenerated with scripts/generate_env_reference.py; the diff is line numbers only. Three existing tests changed. test_workspace_artifact_tool_floor asserted that an unsafe interpreter prefix produces no `--ro-bind <prefix> <prefix>`, which now fires on /home because /home is legitimately a read-only base mount. Asserting the absence of a literal flag string cannot distinguish "the prefix was rejected" from "the argv mounted that root itself", so it compares the argv against the no-prefix baseline instead: an unsafe prefix must add nothing. The Windows bash test asserted dict equality on the whole result, which makes adding a field to every bash result impossible without touching a test about tmux; it asserts the shape now. The personal-dir symlink test grepped the resolver's source for the literal "os.path.realpath", which is gone because the resolution moved into the shared boundary -- it keeps the negative assertion that the closure must not grow its own abspath check again, and the behavioural half now runs against the boundary, where it covers every call site instead of one closure. Not verified: the bubblewrap argv is asserted, not executed. There is no bwrap on macOS, and in Docker it needs --privileged to work at all -- default and seccomp=unconfined both fail with "Creating new namespace failed", and --cap-add=SYS_ADMIN fails at pivot_root. The Python tool's needs_virtual_namespace gate means ordinary Python code gets no namespace even on a Linux host that could provide one; that is reported now but deliberately not changed, because it alters the Linux Python path on every call and cannot be checked from here.
452 lines
19 KiB
Python
452 lines
19 KiB
Python
# routes/personal_routes.py
|
|
"""Routes for personal documents management."""
|
|
import asyncio
|
|
import os
|
|
import logging
|
|
import shutil
|
|
import uuid
|
|
from typing import Any, Dict, List, Tuple
|
|
from fastapi import APIRouter, HTTPException, Query, Request, UploadFile, File, Depends
|
|
from fastapi.concurrency import run_in_threadpool
|
|
from src.request_models import DirectoryRequest
|
|
from core.constants import BASE_DIR, PERSONAL_DIR, PERSONAL_UPLOADS_DIR
|
|
from src.rag_singleton import get_rag_manager
|
|
from src.auth_helpers import require_privilege, require_user
|
|
from core.middleware import require_admin
|
|
from src.path_confinement import confine
|
|
from src.upload_handler import secure_filename
|
|
from src.upload_limits import PERSONAL_UPLOAD_MAX_BYTES
|
|
|
|
UPLOADS_DIR = PERSONAL_UPLOADS_DIR
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
def _personal_upload_dir_for_owner(owner: str | None, *, create: bool = True) -> str:
|
|
"""Return the per-owner upload directory used for direct RAG uploads."""
|
|
owner_segment = secure_filename((owner or "local").strip())[:80] or "local"
|
|
upload_dir = confine(UPLOADS_DIR, owner_segment, allow_root=False)
|
|
if create:
|
|
os.makedirs(upload_dir, exist_ok=True)
|
|
return upload_dir
|
|
|
|
|
|
def _unique_personal_upload_path(upload_dir: str, original_name: str | None) -> Tuple[str, str, str]:
|
|
"""Build a collision-resistant upload path while preserving a display name."""
|
|
safe_name = secure_filename(os.path.basename(original_name or "upload"))
|
|
if not safe_name or safe_name.startswith("."):
|
|
safe_name = "upload"
|
|
|
|
stem, ext = os.path.splitext(safe_name)
|
|
stem = (stem or "upload")[:80]
|
|
filename = f"{stem}-{uuid.uuid4().hex[:10]}{ext.lower()}"
|
|
file_path = confine(upload_dir, filename, allow_root=False)
|
|
return file_path, filename, safe_name
|
|
|
|
|
|
def _unique_existing_target(path: str) -> str:
|
|
"""Return a non-existing sibling path for rename collision handling."""
|
|
if not os.path.exists(path):
|
|
return path
|
|
stem, ext = os.path.splitext(path)
|
|
while True:
|
|
candidate = f"{stem}-{uuid.uuid4().hex[:10]}{ext}"
|
|
if not os.path.exists(candidate):
|
|
return candidate
|
|
|
|
|
|
def _remove_empty_tree(path: str) -> None:
|
|
"""Best-effort removal of empty directories under ``path``."""
|
|
if not os.path.isdir(path):
|
|
return
|
|
for root, dirs, _files in os.walk(path, topdown=False):
|
|
for dirname in dirs:
|
|
candidate = os.path.join(root, dirname)
|
|
try:
|
|
os.rmdir(candidate)
|
|
except OSError:
|
|
pass
|
|
try:
|
|
os.rmdir(path)
|
|
except OSError:
|
|
pass
|
|
|
|
|
|
def rename_personal_upload_owner(
|
|
old_owner: str,
|
|
new_owner: str,
|
|
*,
|
|
personal_docs_manager: Any = None,
|
|
rag_manager: Any = None,
|
|
) -> Dict[str, Any]:
|
|
"""Move direct personal uploads and rewrite RAG owner metadata on user rename."""
|
|
old_dir = _personal_upload_dir_for_owner(old_owner, create=False)
|
|
new_dir = _personal_upload_dir_for_owner(new_owner, create=False)
|
|
path_map: Dict[str, str] = {}
|
|
moved_files = 0
|
|
|
|
if os.path.isdir(old_dir) and old_dir != new_dir:
|
|
os.makedirs(new_dir, exist_ok=True)
|
|
for root, _dirs, files in os.walk(old_dir):
|
|
rel_root = os.path.relpath(root, old_dir)
|
|
target_root = new_dir if rel_root == "." else os.path.join(new_dir, rel_root)
|
|
os.makedirs(target_root, exist_ok=True)
|
|
for filename in files:
|
|
source = os.path.abspath(os.path.join(root, filename))
|
|
target = _unique_existing_target(os.path.abspath(os.path.join(target_root, filename)))
|
|
shutil.move(source, target)
|
|
path_map[source] = target
|
|
moved_files += 1
|
|
_remove_empty_tree(old_dir)
|
|
|
|
if personal_docs_manager is not None:
|
|
rename_directory = getattr(personal_docs_manager, "rename_directory", None)
|
|
if callable(rename_directory):
|
|
rename_directory(old_dir, new_dir, path_map=path_map)
|
|
|
|
rag_result = None
|
|
if rag_manager is not None:
|
|
rename_owner = getattr(rag_manager, "rename_owner", None)
|
|
if callable(rename_owner):
|
|
rag_result = rename_owner(
|
|
old_owner,
|
|
new_owner,
|
|
path_map=path_map,
|
|
path_prefixes=[(old_dir, new_dir)],
|
|
)
|
|
|
|
return {
|
|
"old_dir": old_dir,
|
|
"new_dir": new_dir,
|
|
"moved_files": moved_files,
|
|
"path_map": path_map,
|
|
"rag_result": rag_result,
|
|
}
|
|
|
|
|
|
def setup_personal_routes(personal_docs_manager, rag_manager, rag_available):
|
|
"""
|
|
Setup personal documents related routes.
|
|
|
|
Args:
|
|
personal_docs_manager: PersonalDocsManager instance
|
|
rag_manager: RAG manager instance (may be None)
|
|
rag_available: Boolean indicating if RAG is available
|
|
|
|
Returns:
|
|
APIRouter instance with personal docs routes
|
|
"""
|
|
router = APIRouter(prefix="/api/personal")
|
|
|
|
# Serializes directory index jobs across requests. Indexing runs in the
|
|
# threadpool (#5558), so concurrent requests would otherwise run in parallel
|
|
# and race PersonalDocsManager's unsynchronized list mutations and file
|
|
# writes; before the threadpool move they serialized on the blocked event
|
|
# loop, so one-at-a-time is behavior parity.
|
|
#
|
|
# An asyncio.Lock acquired in the async handler BEFORE offloading: a waiting
|
|
# request parks on the event loop instead of pinning a threadpool worker (an
|
|
# earlier threading.Lock taken INSIDE the worker meant queued jobs held pool
|
|
# tokens while blocked, starving every other run_in_threadpool caller).
|
|
# add/remove/reload all take this lock, so their mutations never interleave.
|
|
# Per-router (not module-global) so each app binds it to its own event loop.
|
|
# Scope is the single process: multi-worker deployments would need a shared
|
|
# lock (out of scope for #5558).
|
|
_index_job_lock = asyncio.Lock()
|
|
|
|
def _rag():
|
|
"""Get the current RAG manager, retrying init if needed."""
|
|
return get_rag_manager()
|
|
|
|
def _resolve_allowed_personal_dir(directory: str) -> str:
|
|
"""Resolve a user-supplied personal-docs path under the allowed root."""
|
|
if not directory:
|
|
raise HTTPException(400, "Directory path is required")
|
|
|
|
try:
|
|
return confine(PERSONAL_DIR, directory)
|
|
except (ValueError, OSError):
|
|
raise HTTPException(403, "Directory must be inside personal documents")
|
|
|
|
@router.get("")
|
|
def api_personal_list(owner: str = Depends(require_user), _admin: None = Depends(require_admin)):
|
|
"""Enhanced version that includes directories"""
|
|
files = [{"name": f["name"], "size": f["size"], "path": f.get("path", "")} for f in personal_docs_manager.index]
|
|
directories = personal_docs_manager.get_indexed_directories() if hasattr(personal_docs_manager, "get_indexed_directories") else []
|
|
return {"files": files, "directories": directories}
|
|
|
|
@router.post("/reload")
|
|
async def api_personal_reload(owner: str = Depends(require_user), _admin: None = Depends(require_admin)):
|
|
# refresh_index() re-extracts text across every tracked directory —
|
|
# blocking work. Take the shared job lock (so it cannot race an add /
|
|
# remove) and run it off the event loop.
|
|
async with _index_job_lock:
|
|
await run_in_threadpool(personal_docs_manager.refresh_index)
|
|
return {"ok": True, "count": len(personal_docs_manager.index)}
|
|
|
|
@router.post("/add_directory")
|
|
async def add_directory_to_rag(
|
|
request: Request,
|
|
directory_request: DirectoryRequest,
|
|
owner: str = Depends(require_user), _admin: None = Depends(require_admin),
|
|
):
|
|
"""
|
|
Add a directory and all its subdirectories/files to the RAG index.
|
|
|
|
Args:
|
|
directory_request: Directory request model containing the directory path
|
|
|
|
Returns:
|
|
JSON response with indexing results
|
|
"""
|
|
directory = directory_request.directory
|
|
try:
|
|
directory = _resolve_allowed_personal_dir(directory)
|
|
|
|
# Security check - ensure directory exists and is accessible
|
|
if not os.path.exists(directory):
|
|
raise HTTPException(404, f"Directory not found: {directory}")
|
|
|
|
if not os.path.isdir(directory):
|
|
raise HTTPException(400, f"Path is not a directory: {directory}")
|
|
|
|
logger.info(f"Adding directory to RAG: {directory}")
|
|
|
|
# Use the RAGManager to index the directory
|
|
rag = _rag()
|
|
if rag:
|
|
def _index_directory():
|
|
result = rag.index_personal_documents(directory, owner=owner)
|
|
if result["success"]:
|
|
# Also update the personal_docs_manager to track this
|
|
# directory. Kept inside the offloaded call: it triggers
|
|
# refresh_index(), which re-extracts text across tracked
|
|
# directories.
|
|
personal_docs_manager.add_directory(directory, index=False)
|
|
return result
|
|
|
|
# Indexing walks, embeds, and stores the whole tree — minutes
|
|
# on a real directory. The handler is async, so calling it
|
|
# inline runs it on the event loop and every other request
|
|
# queues behind it until it finishes (#5558). Serialize on the
|
|
# async job lock BEFORE offloading so a queued request parks on
|
|
# the loop instead of pinning a threadpool worker.
|
|
async with _index_job_lock:
|
|
result = await run_in_threadpool(_index_directory)
|
|
|
|
if result["success"]:
|
|
return {
|
|
"success": True,
|
|
"message": f"Successfully indexed {result['indexed_count']} chunks from {directory}",
|
|
"indexed_count": result["indexed_count"],
|
|
"failed_count": result.get("failed_count", 0),
|
|
"directory": directory
|
|
}
|
|
else:
|
|
raise HTTPException(500, result.get("message", "Failed to index directory"))
|
|
else:
|
|
raise HTTPException(503, "RAG system is not available")
|
|
|
|
except HTTPException:
|
|
raise
|
|
except Exception as e:
|
|
logger.error(f"Error adding directory to RAG: {e}")
|
|
raise HTTPException(500, f"Failed to add directory: {str(e)}")
|
|
|
|
@router.delete("/remove_directory")
|
|
async def remove_directory_from_rag(directory: str = Query(...), owner: str = Depends(require_user), _admin: None = Depends(require_admin)):
|
|
"""
|
|
Remove a directory from the RAG index.
|
|
|
|
Args:
|
|
directory: Path to the directory to remove
|
|
|
|
Returns:
|
|
JSON response confirming removal
|
|
"""
|
|
try:
|
|
# Confine to PERSONAL_DIR — parity with add_directory_to_rag (which
|
|
# resolves the path the same way). Without this, an arbitrary or
|
|
# `..`-escaping path is passed straight to
|
|
# personal_docs_manager.remove_directory / rag.remove_directory.
|
|
directory = _resolve_allowed_personal_dir(directory)
|
|
|
|
logger.info(f"Removing directory from RAG: {directory}")
|
|
|
|
rag = _rag()
|
|
|
|
def _remove_directory():
|
|
# Always remove from personal_docs_manager tracking. This
|
|
# mutates the same unsynchronized list/index an add job touches
|
|
# and re-extracts text (refresh_index), so it is blocking work.
|
|
if hasattr(personal_docs_manager, 'remove_directory'):
|
|
personal_docs_manager.remove_directory(directory)
|
|
# Remove from RAG vector store (best-effort).
|
|
if rag:
|
|
try:
|
|
rag.remove_directory(directory)
|
|
except Exception as e:
|
|
logger.warning(f"RAG removal failed for directory {directory}: {e}")
|
|
|
|
# Same job lock as add/reload so remove cannot interleave with an
|
|
# in-flight add; offloaded off the event loop.
|
|
async with _index_job_lock:
|
|
await run_in_threadpool(_remove_directory)
|
|
|
|
return {
|
|
"success": True,
|
|
"message": f"Successfully removed {directory} from RAG index",
|
|
"directory": directory
|
|
}
|
|
|
|
except HTTPException:
|
|
raise
|
|
except Exception as e:
|
|
logger.error(f"Error removing directory from RAG: {e}")
|
|
raise HTTPException(500, f"Failed to remove directory: {str(e)}")
|
|
|
|
@router.post("/upload")
|
|
async def upload_files_to_rag(request: Request, files: List[UploadFile] = File(...)):
|
|
"""Upload files directly into RAG. Supports text and PDF."""
|
|
user = require_privilege(request, "can_use_documents")
|
|
rag = _rag()
|
|
if not rag:
|
|
raise HTTPException(503, "RAG system is not available — is the embedding service running?")
|
|
|
|
upload_dir = _personal_upload_dir_for_owner(user)
|
|
|
|
total_indexed = 0
|
|
total_failed = 0
|
|
uploaded_files = []
|
|
|
|
# Chunking, embedding and the tracking update are blocking work over the
|
|
# same vector/tracking state add_directory mutates (#5634). Take the
|
|
# shared job lock BEFORE offloading so a queued request parks on the loop
|
|
# instead of pinning a threadpool worker, matching add_directory.
|
|
# Read and process one capped payload at a time so a multi-file request
|
|
# cannot retain len(files) * PERSONAL_UPLOAD_MAX_BYTES in memory.
|
|
async with _index_job_lock:
|
|
for upload in files:
|
|
try:
|
|
file_path, stored_name, safe_name = _unique_personal_upload_path(
|
|
upload_dir, upload.filename
|
|
)
|
|
content_bytes = await upload.read(PERSONAL_UPLOAD_MAX_BYTES + 1)
|
|
if len(content_bytes) > PERSONAL_UPLOAD_MAX_BYTES:
|
|
logger.warning(f"Rejected oversized personal upload: {upload.filename!r}")
|
|
total_failed += 1
|
|
continue
|
|
|
|
def _index_upload():
|
|
with open(file_path, "wb") as f:
|
|
f.write(content_bytes)
|
|
|
|
ext = os.path.splitext(safe_name)[1].lower()
|
|
if ext == ".pdf":
|
|
from src.personal_docs import extract_pdf_text
|
|
text = extract_pdf_text(file_path)
|
|
else:
|
|
text = content_bytes.decode("utf-8", errors="replace")
|
|
|
|
if not text or not text.strip():
|
|
return 0, 1, None
|
|
|
|
indexed = 0
|
|
failed = 0
|
|
chunks = rag._split_into_chunks(text, chunk_size=500)
|
|
for i, chunk in enumerate(chunks):
|
|
metadata = {
|
|
"source": file_path,
|
|
"filename": safe_name,
|
|
"stored_filename": stored_name,
|
|
"directory": upload_dir,
|
|
"type": ext,
|
|
"chunk_id": i,
|
|
}
|
|
if user:
|
|
metadata["owner"] = user
|
|
if rag.add_document(chunk, metadata):
|
|
indexed += 1
|
|
else:
|
|
failed += 1
|
|
return indexed, failed, safe_name
|
|
|
|
indexed, failed, uploaded_name = await run_in_threadpool(_index_upload)
|
|
total_indexed += indexed
|
|
total_failed += failed
|
|
if uploaded_name:
|
|
uploaded_files.append(uploaded_name)
|
|
except Exception as e:
|
|
logger.error(f"Failed to upload/index {upload.filename}: {e}")
|
|
total_failed += 1
|
|
|
|
# Same transition, same lock: the tracking update must not land
|
|
# while another job is mid-write over the same state.
|
|
if uploaded_files and hasattr(personal_docs_manager, "add_directory"):
|
|
await run_in_threadpool(
|
|
personal_docs_manager.add_directory, upload_dir, index=False
|
|
)
|
|
|
|
return {
|
|
"success": True,
|
|
"uploaded": uploaded_files,
|
|
"indexed_count": total_indexed,
|
|
"failed_count": total_failed,
|
|
}
|
|
|
|
@router.delete("/file")
|
|
async def delete_file_from_rag(filepath: str = Query(...), owner: str = Depends(require_user), _admin: None = Depends(require_admin)):
|
|
"""Delete a specific file from RAG index and optionally from disk."""
|
|
try:
|
|
def _delete_file():
|
|
# Remove chunks from RAG vector store (best-effort)
|
|
removed = 0
|
|
rag = _rag()
|
|
if rag:
|
|
try:
|
|
removed = rag.delete_by_source(filepath)
|
|
except Exception as e:
|
|
logger.warning(f"RAG removal failed for {filepath}: {e}")
|
|
|
|
# Delete file from disk if it's in the caller's own uploads dir.
|
|
# Scope to the per-owner subdir, not the shared uploads root, so one
|
|
# admin can't delete another user's personal files by path.
|
|
deleted_from_disk = False
|
|
# allow_root=False: the per-owner upload directory itself is
|
|
# never a deletion target, only files under it.
|
|
try:
|
|
abs_target = confine(
|
|
_personal_upload_dir_for_owner(owner, create=False),
|
|
filepath,
|
|
allow_root=False,
|
|
)
|
|
except (ValueError, OSError):
|
|
abs_target = ""
|
|
if abs_target:
|
|
try:
|
|
os.remove(abs_target)
|
|
deleted_from_disk = True
|
|
except FileNotFoundError:
|
|
pass # already gone — race with another request or cleanup
|
|
|
|
# Exclude the file from the listing (persists across restarts)
|
|
personal_docs_manager.exclude_file(filepath)
|
|
return removed, deleted_from_disk
|
|
|
|
# Vector removal, the disk unlink and the exclusion write are one
|
|
# transition over the same state add_directory mutates (#5634), and
|
|
# all three block. Take the shared job lock BEFORE offloading, as
|
|
# add_directory does.
|
|
async with _index_job_lock:
|
|
removed, deleted_from_disk = await run_in_threadpool(_delete_file)
|
|
|
|
return {
|
|
"success": True,
|
|
"removed_chunks": removed,
|
|
"deleted_from_disk": deleted_from_disk,
|
|
}
|
|
except Exception as e:
|
|
logger.error(f"Failed to delete file {filepath}: {e}")
|
|
raise HTTPException(500, f"Failed to delete file: {str(e)}")
|
|
|
|
return router
|