mirror of
https://github.com/bytedance/deer-flow.git
synced 2026-07-27 08:28:00 +00:00
* fix(gateway): offload gateway upload file IO Move Gateway upload router filesystem work off the asyncio event loop by using a dedicated ContextVar-preserving file IO executor. Use async sandbox acquisition for non-mounted sandbox uploads and offload remote sandbox sync together with host file reads. Add blocking-IO regression coverage for upload, list, delete, and remote sandbox sync paths. * fix(gateway): align file IO worker env var prefix Rename the file IO executor worker-count environment variable from DEERFLOW_FILE_IO_WORKERS to DEER_FLOW_FILE_IO_WORKERS to match the repo's existing runtime configuration prefix convention.
52 lines
1.6 KiB
Python
52 lines
1.6 KiB
Python
"""Dedicated async offload helper for filesystem work."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import atexit
|
|
import contextvars
|
|
import functools
|
|
import logging
|
|
import os
|
|
from collections.abc import Callable
|
|
from concurrent.futures import ThreadPoolExecutor
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
|
|
def _default_file_io_workers() -> int:
|
|
raw = os.getenv("DEER_FLOW_FILE_IO_WORKERS")
|
|
if raw:
|
|
try:
|
|
workers = int(raw)
|
|
if workers > 0:
|
|
return workers
|
|
except ValueError:
|
|
pass
|
|
logger.warning("Invalid DEER_FLOW_FILE_IO_WORKERS value; using default file IO worker count")
|
|
return min(32, (os.cpu_count() or 1) + 4)
|
|
|
|
|
|
_FILE_IO_EXECUTOR = ThreadPoolExecutor(max_workers=_default_file_io_workers(), thread_name_prefix="file-io")
|
|
|
|
|
|
def _shutdown_file_io_executor() -> None:
|
|
_FILE_IO_EXECUTOR.shutdown(wait=False, cancel_futures=True)
|
|
|
|
|
|
atexit.register(_shutdown_file_io_executor)
|
|
|
|
|
|
async def run_file_io[**P, T](func: Callable[P, T], /, *args: P.args, **kwargs: P.kwargs) -> T:
|
|
"""Run blocking filesystem-oriented work on the dedicated file IO pool.
|
|
|
|
``asyncio.to_thread`` copies ``ContextVar`` values automatically; raw
|
|
``loop.run_in_executor`` does not. Copy the current context explicitly so
|
|
user-scoped helpers such as ``get_effective_user_id()`` keep working inside
|
|
the worker thread.
|
|
"""
|
|
loop = asyncio.get_running_loop()
|
|
ctx = contextvars.copy_context()
|
|
call = functools.partial(func, *args, **kwargs)
|
|
return await loop.run_in_executor(_FILE_IO_EXECUTOR, ctx.run, call)
|