Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 14 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,20 @@ Two things are versioned separately from this file and worth knowing about:

## [Unreleased]

### Fixed

- **The Workspaces page no longer sits on its loading state before it draws.**
Its landing view, every file at once, read up to twenty-five container-backed
workspaces one after another and then fetched their image thumbnails one after
another, so the page waited for the sum of all those round trips. Workspaces
and thumbnails are now read side by side, on eight threads of their own shared
by the whole process, so several people opening the page against a host that
has stopped answering cannot hold up sign-in or anything else on the default
thread pool. A listing waits ten seconds for each call to a host instead of the
archive's sixty-second default, so a host that does not answer at all is
reported as unreadable after ten seconds; one that answers slowly can still
take longer over a deep walk. The "Count files" switch reads hosts the same way.

## [0.0.521] - 2026-10-06

### Fixed
Expand Down
205 changes: 161 additions & 44 deletions backend/app/services/sandbox_workspace.py
Original file line number Diff line number Diff line change
Expand Up @@ -25,19 +25,23 @@
import asyncio
import base64
import contextlib
import contextvars
import logging
import shlex
from binascii import Error as BinasciiError
from collections.abc import AsyncIterator
from collections.abc import AsyncIterator, Callable, Coroutine, Sequence
from concurrent.futures import ThreadPoolExecutor
from contextlib import asynccontextmanager
from dataclasses import dataclass, field, replace
from functools import partial
from io import BytesIO
from pathlib import PurePosixPath
from typing import Any
from uuid import UUID

from PIL import Image, ImageOps
from pydantic_ai_backends import FileData, FileInfo
from pydantic_ai_backends.remote.client import DEFAULT_TIMEOUT_SECONDS
from sqlalchemy.ext.asyncio import AsyncSession

from app.agents.capabilities.sandbox import SandboxConfig
Expand Down Expand Up @@ -75,6 +79,58 @@
"""Where a listing stops. A host holding a `node_modules` must not turn one
workspace into ten thousand rows on a page about twenty-five of them."""

_HOST_CALLS_AT_ONCE = 8
"""How many calls to sandbox hosts the whole process has in flight - one `ls` of
a walk or one `read_bytes` each count as one. The listings that read more than one
workspace used to read them one after another, so the page waited for the sum of
twenty-five walks; with this it waits for the slowest few."""

_HOST_CALLS = ThreadPoolExecutor(
max_workers=_HOST_CALLS_AT_ONCE, thread_name_prefix="workspace-host"
)
"""The threads host calls run on, and the bound on them.

Not the loop's default executor, which `to_thread` uses: that one also carries
`bcrypt` and pinned-host DNS, and it has `min(32, cpus + 4)` threads - six on a
two-core host. A host that stops answering holds a thread for the browse timeout,
or sixty seconds for a read, so a few people opening the Workspaces page at once
would leave sign-in queueing behind them. A per-request limit does not prevent
that, since every request brings its own; a pool of the process's own does, and
a dead host then slows only other host calls.

Unlike `app.core.blocking`, which carries file work, there is no admission gate:
what waits in this queue is a bound method and two strings rather than an upload
buffer, and a cancelled caller's call that has not started is dropped from it.
Threads start on first use and are joined at interpreter exit."""


async def _on_host_thread[T](call: Callable[..., T], *args: object) -> T:
"""Run a blocking host call on `_HOST_CALLS`, carrying the caller's context as
`to_thread` would, so its spans stay under the request's trace."""
context = contextvars.copy_context()
return await asyncio.get_running_loop().run_in_executor(
_HOST_CALLS, partial(context.run, call, *args)
)


_BROWSE_TIMEOUT_SECONDS = 10.0
"""How long a listing waits on a host. The archive's own default is sixty seconds,
and a host that has gone away held the whole Workspaces page for that long before
it could say "this host did not answer". A download keeps the default."""


async def _concurrently[T](calls: Sequence[Callable[[], Coroutine[object, object, T]]]) -> list[T]:
"""Every call's result in the order given, run side by side.

A `TaskGroup` rather than `gather`: a call that raises cancels the others and
the group waits for them, so none keeps running after the request has failed
and reaches its database session afterwards. How many reach a host at once is
`_HOST_CALLS`'s to bound, not this.
"""
async with asyncio.TaskGroup() as group:
tasks = [group.create_task(call()) for call in calls]
return [task.result() for task in tasks]


@dataclass(frozen=True)
class WorkspaceContents:
Expand Down Expand Up @@ -362,6 +418,10 @@ def __init__(self, db: AsyncSession) -> None:
self.db = db
self.connections = SandboxConnectionService(db)
self._connections_by_id: dict[UUID, ResolvedConnection] = {}
# The listings read hosts concurrently, and resolving a connection is a
# query on the one session this service holds - which must not be used by
# two coroutines at once.
self._resolving = asyncio.Lock()

async def open(
self,
Expand Down Expand Up @@ -950,6 +1010,7 @@ async def measured(
unreadable = 0
read = 0
remote = 0
rows: list[AgentWorkspace] = []
for overview in overviews:
row = overview.row
if row.backend != "state":
Expand All @@ -958,7 +1019,8 @@ async def measured(
if remote >= limit:
continue
remote += 1
contents = await self._entries(ctx, row)
rows.append(row)
for row, contents in zip(rows, await self._entries_of(ctx, rows), strict=True):
if contents.unreadable_reason is not None:
unreadable += 1
continue
Expand Down Expand Up @@ -996,20 +1058,38 @@ async def flat_files(self, ctx: AuthContext, *, limit: int = 25) -> FlatFileList
"""
overviews = await self.visible_to(ctx)
files: list[FlatEntry] = []
unreadable = 0
# A budget across the whole request, not per workspace: a host's thumbnail
# is a `read_bytes` for that file, and twenty-five workspaces of photographs
# would be a page that fetches two hundred images to draw them 64 pixels
# wide. The tiles past it fall back to the grey glyph, which is what every
# container-backed image looked like before.
budget = HOST_THUMBNAIL_BUDGET
for overview in overviews[:limit]:
contents = await self._entries(ctx, overview.row)
if contents.unreadable_reason is not None:
unreadable += 1
opened = overviews[:limit]
listed = await self._entries_of(ctx, [overview.row for overview in opened])
readable: list[tuple[WorkspaceOverview, list[FileInfo]]] = [
(
overview,
browsable([entry for entry in contents.entries if not entry.get("is_dir")]),
)
for overview, contents in zip(opened, listed, strict=True)
if contents.unreadable_reason is None
]
# Shared out in listing order before anything is fetched, so which tiles are
# drawn does not depend on which host happened to answer first.
shares: list[int] = []
for overview, entries in readable:
if overview.row.backend == "state":
shares.append(0)
continue
entries = browsable([entry for entry in contents.entries if not entry.get("is_dir")])
previews, thumbnails, budget = await self._tiles(ctx, overview.row, entries, budget)
picked, budget = thumbnail_picks(entries, budget)
shares.append(len(picked))
tiled = await _concurrently(
[
partial(self._tiles, ctx, overview.row, entries, share)
for (overview, entries), share in zip(readable, shares, strict=True)
]
)
for (overview, entries), (previews, thumbnails, _left) in zip(readable, tiled, strict=True):
files.extend(
FlatEntry(
overview=overview,
Expand All @@ -1021,11 +1101,24 @@ async def flat_files(self, ctx: AuthContext, *, limit: int = 25) -> FlatFileList
)
return FlatFileListing(
files=files,
workspaces_read=min(len(overviews), limit) - unreadable,
unreadable=unreadable,
workspaces_read=len(readable),
unreadable=len(opened) - len(readable),
truncated=len(overviews) > limit,
)

async def _entries_of(
self, ctx: AuthContext, rows: list[AgentWorkspace]
) -> list[WorkspaceContents]:
"""What each of several workspaces holds, in the order they were given.

Concurrently, because a container-backed workspace is a walk of round trips
to its host, and these listings used to make them one workspace after
another - twenty-five walks end to end before the Workspaces page could draw
anything. A host that will not answer still comes back as
`unreadable_reason`, never as a raise, so one of them cannot fail the rest.
"""
return await _concurrently([partial(self._entries, ctx, row) for row in rows])

async def _host_thumbnails(
self,
ctx: AuthContext,
Expand All @@ -1045,29 +1138,24 @@ async def _host_thumbnails(
one, or a host that stops answering mid-grid, leaves that tile with the
glyph every container-backed image had before this existed.
"""
drawn: dict[str, str] = {}
for entry in entries:
if budget <= 0:
break
path = str(entry.get("path"))
if PurePosixPath(path).suffix.lower() not in THUMBNAIL_SUFFIXES:
continue
size = entry.get("size")
if size is None or int(size) > THUMBNAIL_SOURCE_LIMIT:
continue
budget -= 1
try:
raw = await self._read_bytes_from(ctx, row, path)
except Exception:
logger.warning("workspace_thumbnail_unreadable", extra={"path": path})
continue
if raw is None:
continue
scaled = thumbnail_of(path, raw)
if scaled is not None:
drawn[path] = scaled
paths, budget = thumbnail_picks(entries, budget)
tiles = await _concurrently(
[partial(self._host_thumbnail, ctx, row, path) for path in paths]
)
drawn = {path: tile for path, tile in zip(paths, tiles, strict=True) if tile is not None}
return drawn, budget

async def _host_thumbnail(self, ctx: AuthContext, row: AgentWorkspace, path: str) -> str | None:
"""One host image drawn as a tile, or `None` for the glyph."""
try:
raw = await self._read_bytes_from(ctx, row, path)
except Exception:
logger.warning("workspace_thumbnail_unreadable", extra={"path": path})
return None
if raw is None:
return None
return thumbnail_of(path, raw)

async def files_of(
self, ctx: AuthContext, workspace_id: UUID
) -> tuple[AgentWorkspace, WorkspaceContents]:
Expand Down Expand Up @@ -1274,7 +1362,7 @@ async def _read_bytes_from(
if archive is None:
return None
try:
return await asyncio.to_thread(
return await _on_host_thread(
archive.read_bytes, row.session_id or row.scope_key, path
)
except Exception as exc:
Expand Down Expand Up @@ -1334,9 +1422,11 @@ async def _walk(archive: Any, session: str) -> tuple[list[FileInfo], bool]:
not available here: it needs a running session, and browsing deliberately
never starts one.

`to_thread`, because `WorkspaceArchive` is a synchronous `httpx.Client`.
Off the loop, because `WorkspaceArchive` is a synchronous `httpx.Client`:
`flat_files` runs this for up to 25 workspaces in one request, so a blocking
call would hold the loop for all 25 - and not only for that request.
call would hold the loop for all 25 - and not only for that request. On
`_HOST_CALLS` rather than `to_thread`, so a host that stops answering cannot
take the default executor's threads; see there.

Returns:
What was found, and whether either bound stopped it - because a
Expand All @@ -1349,7 +1439,7 @@ async def _walk(archive: Any, session: str) -> tuple[list[FileInfo], bool]:
while queue:
path, depth = queue.pop(0)
try:
entries = await asyncio.to_thread(archive.ls, session, path)
entries = await _on_host_thread(archive.ls, session, path)
except Exception:
# The root's failure is the host's and belongs to the caller, which
# turns it into `unreadable_reason`. One directory below it refusing
Expand Down Expand Up @@ -1379,7 +1469,7 @@ async def _walk(archive: Any, session: str) -> tuple[list[FileInfo], bool]:

async def _remote_entries(self, ctx: AuthContext, row: AgentWorkspace) -> WorkspaceContents:
try:
async with self._archive(ctx, row) as archive:
async with self._archive(ctx, row, timeout=_BROWSE_TIMEOUT_SECONDS) as archive:
if archive is None:
return WorkspaceContents(entries=[])
session = row.session_id or row.scope_key
Expand Down Expand Up @@ -1464,7 +1554,9 @@ def _state_bytes(row: AgentWorkspace, path: str) -> bytes | None:
return backend.read_bytes(path)

@asynccontextmanager
async def _archive(self, ctx: AuthContext, row: AgentWorkspace) -> AsyncIterator[Any | None]:
async def _archive(
self, ctx: AuthContext, row: AgentWorkspace, *, timeout: float = DEFAULT_TIMEOUT_SECONDS
) -> AsyncIterator[Any | None]:
"""A reader for the host volume behind a container-backed workspace.

`None` when the workspace has no connection left to ask - the host was
Expand All @@ -1488,7 +1580,7 @@ async def _archive(self, ctx: AuthContext, row: AgentWorkspace) -> AsyncIterator
if resolved.kind != "docker" or not resolved.row.base_url:
yield None
return
archive = WorkspaceArchive(resolved.row.base_url, token=resolved.token)
archive = WorkspaceArchive(resolved.row.base_url, token=resolved.token, timeout=timeout)
try:
yield archive
finally:
Expand All @@ -1510,11 +1602,12 @@ async def _connection(self, ctx: AuthContext, connection_id: UUID) -> ResolvedCo
process. Keyed on the id rather than the row, so two workspaces on one
connection share the answer and two connections do not.
"""
cached = self._connections_by_id.get(connection_id)
if cached is None:
cached = await self.connections.resolve(ctx, connection_id)
self._connections_by_id[connection_id] = cached
return cached
async with self._resolving:
cached = self._connections_by_id.get(connection_id)
if cached is None:
cached = await self.connections.resolve(ctx, connection_id)
self._connections_by_id[connection_id] = cached
return cached


def _absent(row: AgentWorkspace, contents: WorkspaceContents, path: str) -> bool:
Expand Down Expand Up @@ -1658,6 +1751,30 @@ def stored_thumbnail(path: str, data: FileData | None) -> str | None:
return thumbnail_of(path, raw)


def thumbnail_picks(entries: list[FileInfo], budget: int) -> tuple[list[str], int]:
"""Which of a host's images a listing will fetch, and the budget left.

A stored image arrives base64 in the document the listing already read, so
drawing it costs nothing; a host's is a `read_bytes` for that one file. So this
is bounded and the size is checked *before* fetching, off the listing entry - a
20 MB photograph would be read in full to be thrown away by the scaler's own
ceiling.
"""
picked: list[str] = []
for entry in entries:
if budget <= 0:
break
path = str(entry.get("path"))
if PurePosixPath(path).suffix.lower() not in THUMBNAIL_SUFFIXES:
continue
size = entry.get("size")
if size is None or int(size) > THUMBNAIL_SOURCE_LIMIT:
continue
budget -= 1
picked.append(path)
return picked, budget


def thumbnail_of(path: str, raw: bytes) -> str | None:
"""An image's bytes scaled to a data URI, wherever they came from.

Expand Down
Loading
Loading