From 5ac9f15de4cac26478e36f69206fba1decc9eae3 Mon Sep 17 00:00:00 2001 From: menzelj Date: Sun, 21 Jun 2026 20:28:49 +0000 Subject: [PATCH] 0.37.2: stream folder zip-downloads to fix 504 on large folders MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A big folder hit a 504 Gateway Timeout: the zip was built into a temp file *before* any response was sent, so for large folders the backend stayed silent past nginx's proxy_read_timeout. Now the zip is streamed as it's built, end to end: - file_service.open_archive() returns (filename, byte iterator); _iter_zip walks the dir and yields zip bytes incrementally via a small drain buffer, writing each file in 1 MiB chunks (bounded memory, valid CRCs). Same hardening as before — only real regular files; FIFOs/sockets/ devices/symlinks skipped without open(); per-file read errors skipped. - /api/files/download and /agent/files/download return a StreamingResponse (no temp file). The agent proxy streams the agent response straight through (agent_service.stream_download), pulling the first chunk eagerly so an offline/bad-token agent still yields a clean status before 200. - Files page: streamed downloads have no Content-Length, so the progress bar shows the running downloaded byte count ("Downloading … 12.3 MB") instead of a percentage, after the initial "Preparing …". Verified end to end via TestClient (200, application/zip, valid zip, 2 MiB file intact, FIFO skipped, no hang). Co-Authored-By: Claude Opus 4.8 --- backend/agent_app.py | 12 ++-- backend/routers/agents.py | 33 ++++++--- backend/routers/files.py | 19 ++++-- backend/services/agent_service.py | 29 ++++++++ backend/services/file_service.py | 108 +++++++++++++++++++----------- backend/version.py | 2 +- frontend/package.json | 2 +- frontend/src/pages/Files.tsx | 10 ++- 8 files changed, 149 insertions(+), 66 deletions(-) diff --git a/backend/agent_app.py b/backend/agent_app.py index 0abe167..cbf5a13 100644 --- a/backend/agent_app.py +++ b/backend/agent_app.py @@ -32,8 +32,7 @@ from fastapi import ( WebSocket, WebSocketDisconnect, ) -from fastapi.responses import FileResponse, JSONResponse -from starlette.background import BackgroundTask +from fastapi.responses import FileResponse, JSONResponse, StreamingResponse from pydantic import BaseModel from config import settings @@ -668,10 +667,11 @@ def files_read(path: str = Query(...)) -> dict: @app.get("/agent/files/download", dependencies=[Depends(verify_token)]) def files_download(path: str = Query(...)): if _file_guard(file_service.is_dir, path): - tmp, filename = _file_guard(file_service.archive_dir, path) - return FileResponse( - tmp, filename=filename, media_type="application/zip", - background=BackgroundTask(os.unlink, tmp), + filename, chunks = _file_guard(file_service.open_archive, path) + # Stream the zip as it's built (no temp file, starts immediately). + return StreamingResponse( + chunks, media_type="application/zip", + headers={"Content-Disposition": f'attachment; filename="{filename}"'}, ) real, filename = _file_guard(file_service.resolve_download, path) return FileResponse(real, filename=filename, media_type="application/octet-stream") diff --git a/backend/routers/agents.py b/backend/routers/agents.py index 4468942..440253d 100644 --- a/backend/routers/agents.py +++ b/backend/routers/agents.py @@ -6,7 +6,7 @@ import os import tempfile from fastapi import APIRouter, Depends, File, Form, HTTPException, Query, Request, UploadFile -from fastapi.responses import FileResponse +from fastapi.responses import FileResponse, StreamingResponse from pydantic import BaseModel from sqlmodel import Session, select from starlette.background import BackgroundTask @@ -791,19 +791,30 @@ async def agent_files_download( _user: User = Depends(get_current_user), ): agent = _get_or_404(session, agent_id) - tmp = tempfile.NamedTemporaryFile(delete=False) - tmp.close() + # Stream the agent's response straight through (works for single files and + # for on-the-fly folder zips), so nothing is staged to disk and the + # download starts immediately. Pull the first chunk eagerly so a failed + # agent (offline / bad token / 404) still surfaces a clean HTTP status + # before we commit to a 200 streaming response. + chunks = agent_service.stream_download( + session, agent, "/agent/files/download", params={"path": path} + ) try: - await agent_service.download_to_file( - session, agent, "/agent/files/download", tmp.name, params={"path": path} - ) + first = await chunks.__anext__() + except StopAsyncIteration: + first = b"" except AgentError as exc: - if os.path.exists(tmp.name): - os.unlink(tmp.name) _raise(exc) - return FileResponse( - tmp.name, media_type="application/octet-stream", filename=os.path.basename(path), - background=BackgroundTask(os.unlink, tmp.name), + + async def body(): + yield first + async for chunk in chunks: + yield chunk + + return StreamingResponse( + body(), + media_type="application/octet-stream", + headers={"Content-Disposition": f'attachment; filename="{os.path.basename(path)}"'}, ) diff --git a/backend/routers/files.py b/backend/routers/files.py index ad3319d..8fb6554 100644 --- a/backend/routers/files.py +++ b/backend/routers/files.py @@ -19,8 +19,7 @@ from fastapi import ( Request, UploadFile, ) -from fastapi.responses import FileResponse -from starlette.background import BackgroundTask +from fastapi.responses import FileResponse, StreamingResponse from pydantic import BaseModel from sqlmodel import Session @@ -32,6 +31,12 @@ from services import audit_service, device_service, file_service router = APIRouter(prefix="/api/files", tags=["files"]) +def _attachment(filename: str) -> str: + """A safe ``Content-Disposition`` value for an arbitrary filename.""" + safe = filename.replace("\\", "_").replace('"', "_") + return f'attachment; filename="{safe}"' + + def _ip(request: Request) -> str: return request.client.host if request.client else "unknown" @@ -71,10 +76,12 @@ def download( _user: User = Depends(get_current_user), ): if _guard(file_service.is_dir, path): - tmp, filename = _guard(file_service.archive_dir, path) - return FileResponse( - tmp, filename=filename, media_type="application/zip", - background=BackgroundTask(os.unlink, tmp), + filename, chunks = _guard(file_service.open_archive, path) + # Stream the zip as it's built so the response starts immediately + # (large folders no longer hit the proxy's read timeout). + return StreamingResponse( + chunks, media_type="application/zip", + headers={"Content-Disposition": _attachment(filename)}, ) real, filename = _guard(file_service.resolve_download, path) return FileResponse(real, filename=filename, media_type="application/octet-stream") diff --git a/backend/services/agent_service.py b/backend/services/agent_service.py index a51c73d..5e868d1 100644 --- a/backend/services/agent_service.py +++ b/backend/services/agent_service.py @@ -137,6 +137,35 @@ async def download_to_file( raise AgentError(502, "agent_unreachable", str(exc)) from exc +async def stream_download( + session: Session, + agent: Agent, + path: str, + *, + params: Optional[dict] = None, +): + """Stream a GET from the agent straight through, yielding chunks. + + Unlike :func:`download_to_file` this never buffers to disk, so a large + response (e.g. a folder zip the agent builds on the fly) starts flowing to + the browser immediately instead of being staged first. + """ + url = agent.url.rstrip("/") + path + headers = {"Authorization": f"Bearer {agent.token}"} + try: + async with httpx.AsyncClient(follow_redirects=True) as client: + async with client.stream("GET", url, headers=headers, params=params, timeout=None) as resp: + if resp.status_code >= 400: + text = (await resp.aread()).decode("utf-8", "replace") + _handle_status(session, agent, resp.status_code, text) + _handle_status(session, agent, resp.status_code) + async for chunk in resp.aiter_bytes(1024 * 256): + yield chunk + except httpx.HTTPError as exc: + _mark(session, agent, "offline") + raise AgentError(502, "agent_unreachable", str(exc)) from exc + + async def upload_file( session: Session, agent: Agent, diff --git a/backend/services/file_service.py b/backend/services/file_service.py index 2d9b36b..27becfb 100644 --- a/backend/services/file_service.py +++ b/backend/services/file_service.py @@ -10,8 +10,8 @@ from __future__ import annotations import os import shutil -import tempfile import zipfile +from collections.abc import Iterator from services.device_service import BrowseError, _is_allowed, _real_root @@ -182,54 +182,84 @@ def is_dir(path: str) -> bool: return os.path.isdir(_safe_real(path)) -def archive_dir(path: str) -> tuple[str, str]: - """Zip a directory (recursively) into a temp file. +class _ZipBuffer: + """A writable sink that hands out and clears whatever was written to it. - Returns ``(tmp_zip_path, download_filename)``. The caller is responsible - for deleting the temp file once it has been streamed to the client. + Lets us drive ``zipfile`` while draining its output incrementally so the + archive can be streamed to the client instead of buffered to disk. + """ + + def __init__(self) -> None: + self._buf = bytearray() + + def write(self, data: bytes) -> int: + self._buf += data + return len(data) + + def flush(self) -> None: # pragma: no cover - zipfile calls this + pass + + def take(self) -> bytes: + data = bytes(self._buf) + self._buf.clear() + return data + + +def open_archive(path: str) -> tuple[str, "Iterator[bytes]"]: + """Validate a directory and return ``(download_filename, byte_iterator)``. + + The iterator zips the directory recursively **on the fly**, yielding bytes + as they are produced so the response starts immediately (no waiting for the + whole archive to build → no gateway timeout) and memory stays bounded. Only regular files and real subdirectories are archived. Symlinks are skipped (no sandbox escape / loops); special files (FIFOs, sockets, - devices) are skipped too — opening a FIFO would block forever and a - socket can't be read at all. Files that can't be read (permissions, or - that vanish mid-walk) are skipped individually rather than aborting the - whole archive. + devices) are skipped too — opening a FIFO would block forever and a socket + can't be read at all. Files that can't be read (permissions, or that vanish + mid-walk) are skipped individually rather than aborting the whole archive. """ real = _safe_real(path) if not os.path.isdir(real): raise BrowseError(f"Not a directory: {path}") name = os.path.basename(path.rstrip("/")) or "root" + return f"{name}.zip", _iter_zip(real, name) - fd, tmp = tempfile.mkstemp(suffix=".zip") - os.close(fd) - try: - with zipfile.ZipFile(tmp, "w", zipfile.ZIP_DEFLATED) as zf: - for root, dirs, files in os.walk(real): - # Don't follow symlinked directories (avoids loops / escapes). - dirs[:] = [d for d in dirs if not os.path.islink(os.path.join(root, d))] - rel_root = os.path.relpath(root, real) - if not files and not dirs and rel_root != ".": - # Preserve otherwise-empty directories. - zf.writestr(os.path.join(name, rel_root) + "/", "") - for f in files: - full = os.path.join(root, f) - # os.path.isfile follows symlinks; combined with the islink - # check it admits only real regular files (skips FIFOs, - # sockets, devices and symlinks without ever open()-ing them). - if os.path.islink(full) or not os.path.isfile(full): - continue - arc = (os.path.join(name, rel_root, f) if rel_root != "." - else os.path.join(name, f)) - try: - zf.write(full, arc) - except OSError: - # Unreadable or vanished mid-walk — skip just this file. - continue - except OSError: - if os.path.exists(tmp): - os.unlink(tmp) - raise - return tmp, f"{name}.zip" + +def _iter_zip(real: str, name: str): + sink = _ZipBuffer() + with zipfile.ZipFile(sink, "w", zipfile.ZIP_DEFLATED) as zf: + for root, dirs, files in os.walk(real): + # Don't follow symlinked directories (avoids loops / escapes). + dirs[:] = [d for d in dirs if not os.path.islink(os.path.join(root, d))] + rel_root = os.path.relpath(root, real) + if not files and not dirs and rel_root != ".": + # Preserve otherwise-empty directories. + zf.writestr(os.path.join(name, rel_root) + "/", "") + if chunk := sink.take(): + yield chunk + for f in files: + full = os.path.join(root, f) + # os.path.isfile follows symlinks; combined with the islink + # check it admits only real regular files (skips FIFOs, sockets, + # devices and symlinks without ever open()-ing them). + if os.path.islink(full) or not os.path.isfile(full): + continue + arc = (os.path.join(name, rel_root, f) if rel_root != "." + else os.path.join(name, f)) + try: + info = zipfile.ZipInfo.from_file(full, arc) + info.compress_type = zipfile.ZIP_DEFLATED + with open(full, "rb") as src, zf.open(info, "w") as dest: + while buf := src.read(1024 * 1024): + dest.write(buf) + if chunk := sink.take(): + yield chunk + except OSError: + # Unreadable or vanished mid-walk — skip just this file. + continue + if chunk := sink.take(): + yield chunk + yield sink.take() def upload_target( diff --git a/backend/version.py b/backend/version.py index d0255c4..3e7c20b 100644 --- a/backend/version.py +++ b/backend/version.py @@ -1,3 +1,3 @@ """Single source of truth for the StackPilot release version.""" -APP_VERSION = "0.37.1" +APP_VERSION = "0.37.2" diff --git a/frontend/package.json b/frontend/package.json index 494b242..a92c642 100644 --- a/frontend/package.json +++ b/frontend/package.json @@ -1,7 +1,7 @@ { "name": "stackpilot-frontend", "private": true, - "version": "0.37.1", + "version": "0.37.2", "type": "module", "scripts": { "dev": "vite", diff --git a/frontend/src/pages/Files.tsx b/frontend/src/pages/Files.tsx index 2d709c1..9efc8d2 100644 --- a/frontend/src/pages/Files.tsx +++ b/frontend/src/pages/Files.tsx @@ -88,6 +88,7 @@ export function Files() { pct: number; kind: "upload" | "download"; indeterminate?: boolean; + loaded?: number; } | null>(null); const fileInput = useRef(null); const folderInput = useRef(null); @@ -222,6 +223,7 @@ export function Files() { pct: total ? (loaded / total) * 100 : 0, kind: "download", indeterminate: !total, + loaded, }), ) .catch((err) => toast.error(apiErrorMessage(err))) @@ -329,14 +331,18 @@ export function Files() { )} {progress.kind === "download" - ? progress.indeterminate + ? progress.indeterminate && !progress.loaded ? `Preparing ${progress.label}…` : `Downloading ${progress.label}` : `Uploading ${progress.label}`} - {progress.indeterminate ? "" : `${Math.round(progress.pct)}%`} + {progress.indeterminate + ? progress.loaded + ? formatBytes(progress.loaded) + : "" + : `${Math.round(progress.pct)}%`}