Files
devplacepy/devplacepy/routers/projects/containers/workspace.py
T
retoorandClaude Sonnet 5 4a13415b43
DevPlace CI / test (push) Failing after 26m53s
Fix blocking sync I/O on the event loop across containers, jobs, and xmlrpc
The container/workspace reachability probes (socket connect + HTTP check)
ran synchronously with real timeouts inside async request handlers and the
live-view relay's 3-4s ticks, freezing the whole event loop whenever a
container wasn't cleanly reachable - the likely cause of the periodic
app-wide stalls. Converted the probe chain (api._port_reachable/_http_probe,
editor_reachable, instance_runtime, provision.editor_ready/view) to real
async I/O and parallelized the admin container/workspace list decorators.

Also fixes: XmlrpcService.on_disable blocked up to 10s on a synchronous
subprocess.wait inside an async method (now matches TelegramService's
async-subprocess pattern); JobService._sweep_expired ran every job kind's
cleanup() - including shutil.rmtree on large directories - synchronously
on every tick, now offloaded via asyncio.to_thread for all job kinds at
once; and several smaller blocking reads/writes on request/service paths
(attachment-to-gitea mirroring, stealth chunked downloads, dbapi/isslop
file reads, job payload/report I/O) moved off the loop thread.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-09-03 12:59:10 +00:00

357 lines
13 KiB
Python

# retoor <retoor@molodetz.nl>
from typing import Annotated
from fastapi import APIRouter, Depends, Form, Request, WebSocket
from starlette.responses import Response
from devplacepy.content import (
can_manage_workspace,
can_open_workspace,
)
from devplacepy.database import get_table, resolve_by_slug
from devplacepy.dependencies import json_or_form
from devplacepy.models import EditorPrefsForm, TunnelForm
from devplacepy.responses import action_result, json_error, respond
from devplacepy.schemas import WorkspaceOut
from devplacepy.services.audit import record as audit
from devplacepy.services.containers import activity, api, forward, store
from devplacepy.services.containers.api import ContainerError
from devplacepy.services.containers.workspace import editor, provision, quota, tunnels
from devplacepy.services.containers.workspace.provision import WorkspaceError
from devplacepy.utils import not_found, require_user
from ._shared import audit_instance, fail
router = APIRouter()
def _restart_required(instance: dict, profile: editor.EditorProfile) -> bool:
if instance.get("status") != store.ST_RUNNING:
return False
return editor.restart_required(instance, profile)
def _project_or_404(slug: str) -> dict:
project = resolve_by_slug(get_table("projects"), slug)
if not project:
raise not_found("Project not found")
return project
def _workspace_or_404(project: dict, user: dict) -> dict:
instance = provision.find_for_project(project["uid"], user["uid"])
if not instance:
raise not_found("No workspace for this project")
return instance
def _guard(request: Request, project: dict, user: dict, event_key: str) -> None:
if can_open_workspace(project, user):
return
audit.record(
request,
event_key,
user=user,
target_type="project",
target_uid=project["uid"],
target_label=project.get("title"),
summary=f"{user['username']} denied workspace access",
result="denied",
)
raise not_found("Workspaces are not available for this project")
@router.get("/{slug}/workspace")
async def workspace_page(request: Request, slug: str):
user = require_user(request)
if isinstance(user, Response):
return user
project = _project_or_404(slug)
_guard(request, project, user, "container.workspace.open")
instance = provision.find_for_project(project["uid"], user["uid"])
limits = quota.resolve(user["uid"])
profile = editor.resolve(user["uid"], instance)
context = {
"project": project,
"workspace": await provision.view(instance) if instance else None,
"has_workspace": bool(instance),
"viewer_can_workspace": True,
"workspace_count": provision.count_for_owner(user["uid"]),
"max_workspaces": limits.max_workspaces,
"unlimited_workspaces": limits.unlimited,
"editor_url": (
f"/projects/{slug}/containers/instances/{instance['uid']}/code/"
if instance
else ""
),
"editor_password": (
api.ensure_editor_password(instance) if instance else ""
),
"editor": editor.view(user["uid"], instance),
"restart_required": (
_restart_required(instance, profile) if instance else False
),
"user": user,
}
return respond(request, "workspace.html", context, model=WorkspaceOut)
@router.post("/{slug}/workspace")
async def workspace_open(request: Request, slug: str):
user = require_user(request)
if isinstance(user, Response):
return user
project = _project_or_404(slug)
_guard(request, project, user, "container.workspace.quota.block")
try:
instance = await provision.ensure(project, user)
instance = provision.resume(instance)
except WorkspaceError as error:
audit.record(
request,
"container.workspace.quota.block",
user=user,
target_type="project",
target_uid=project["uid"],
summary=str(error),
result="denied",
)
return json_error(400, str(error))
except ContainerError as error:
return fail(error)
audit_instance(
request,
user,
"container.workspace.create",
instance,
project,
summary=f"{user['username']} opened workspace for {project.get('title')}",
)
provision.write_manifest(instance)
return action_result(
request, f"/projects/{slug}/workspace", data=await provision.view(instance)
)
@router.post("/{slug}/workspace/stop")
async def workspace_stop(request: Request, slug: str):
user = require_user(request)
if isinstance(user, Response):
return user
project = _project_or_404(slug)
instance = _workspace_or_404(project, user)
if not can_manage_workspace(instance, project, user):
return json_error(403, "Not allowed to manage this workspace")
provision.stop(instance)
audit_instance(request, user, "container.workspace.stop", instance, project)
return action_result(request, f"/projects/{slug}/workspace")
@router.post("/{slug}/workspace/delete")
async def workspace_delete(request: Request, slug: str):
user = require_user(request)
if isinstance(user, Response):
return user
project = _project_or_404(slug)
instance = _workspace_or_404(project, user)
if not can_manage_workspace(instance, project, user):
return json_error(403, "Not allowed to manage this workspace")
for row in tunnels.list_for_instance(instance["uid"]):
tunnels.soft_delete(row["uid"], user["uid"])
store.delete_instance(instance["uid"], user["uid"])
audit_instance(request, user, "container.workspace.delete", instance, project)
return action_result(request, f"/projects/{slug}/workspace")
@router.get("/{slug}/workspace/editor")
async def editor_prefs_read(request: Request, slug: str):
user = require_user(request)
if isinstance(user, Response):
return user
project = _project_or_404(slug)
instance = _workspace_or_404(project, user)
if not can_manage_workspace(instance, project, user):
return json_error(403, "Not allowed to manage this workspace")
profile = editor.resolve(user["uid"], instance)
return {
"editor": editor.view(user["uid"], instance),
"restart_required": _restart_required(instance, profile),
}
@router.post("/{slug}/workspace/editor")
async def editor_prefs_write(
request: Request,
slug: str,
data: Annotated[EditorPrefsForm, Depends(json_or_form(EditorPrefsForm))],
):
user = require_user(request)
if isinstance(user, Response):
return user
project = _project_or_404(slug)
instance = _workspace_or_404(project, user)
if not can_manage_workspace(instance, project, user):
return json_error(403, "Not allowed to manage this workspace")
owner_uid = instance.get("workspace_owner_uid") or user["uid"]
if data.reset:
editor.reset_prefs(owner_uid, user["uid"])
summary = f"{user['username']} reset their workspace editor preferences"
else:
editor.save_prefs(
owner_uid, data.model_dump(exclude={"reset"}, exclude_unset=True)
)
summary = f"{user['username']} updated their workspace editor preferences"
audit_instance(
request,
user,
"container.workspace.editor.update",
instance,
project,
summary=summary,
)
profile = editor.resolve(owner_uid, instance)
return action_result(
request,
f"/projects/{slug}/workspace",
data={
"editor": editor.view(owner_uid, instance),
"restart_required": _restart_required(instance, profile),
},
)
@router.get("/{slug}/workspace/tunnels")
async def tunnel_list(request: Request, slug: str):
user = require_user(request)
if isinstance(user, Response):
return user
project = _project_or_404(slug)
instance = _workspace_or_404(project, user)
if not can_manage_workspace(instance, project, user):
return json_error(403, "Not allowed to manage this workspace")
return {"tunnels": tunnels.list_for_instance(instance["uid"])}
@router.post("/{slug}/workspace/tunnels")
async def tunnel_create(
request: Request, slug: str, data: Annotated[TunnelForm, Form()]
):
user = require_user(request)
if isinstance(user, Response):
return user
project = _project_or_404(slug)
instance = _workspace_or_404(project, user)
if not can_manage_workspace(instance, project, user):
return json_error(403, "Not allowed to manage this workspace")
try:
row = provision.publish_tunnel(
instance, data.label, data.container_port, user["uid"]
)
except provision.WorkspaceError as error:
return json_error(400, str(error))
audit_instance(
request,
user,
"container.tunnel.create",
instance,
project,
metadata={"hostname": row["hostname"], "port": data.container_port},
)
return action_result(request, f"/projects/{slug}/workspace", data=row)
@router.delete("/{slug}/workspace/tunnels/{uid}")
@router.post("/{slug}/workspace/tunnels/{uid}/delete")
async def tunnel_delete(request: Request, slug: str, uid: str):
user = require_user(request)
if isinstance(user, Response):
return user
project = _project_or_404(slug)
instance = _workspace_or_404(project, user)
if not can_manage_workspace(instance, project, user):
return json_error(403, "Not allowed to manage this workspace")
row = tunnels.get(uid)
if not row or row.get("instance_uid") != instance["uid"]:
raise not_found("Tunnel not found")
tunnels.soft_delete(uid, user["uid"])
audit_instance(
request,
user,
"container.tunnel.delete",
instance,
project,
metadata={"hostname": row.get("hostname", "")},
)
provision.write_manifest(instance)
return action_result(request, f"/projects/{slug}/workspace")
def _editor_guard(request: Request, slug: str, uid: str):
user = require_user(request)
if isinstance(user, Response):
return None, None, user
project = _project_or_404(slug)
instance = store.get_instance(uid)
if not instance or not instance.get("is_workspace"):
raise not_found("Workspace not found")
if not can_manage_workspace(instance, project, user):
return None, None, json_error(403, "Not allowed to open this workspace")
return project, instance, None
@router.api_route(
"/{slug}/containers/instances/{uid}/code", methods=forward.METHODS
)
@router.api_route(
"/{slug}/containers/instances/{uid}/code/{path:path}", methods=forward.METHODS
)
async def editor_proxy(request: Request, slug: str, uid: str, path: str = ""):
project, instance, denial = _editor_guard(request, slug, uid)
if denial is not None:
return denial
if instance.get("suspended_at"):
return Response(
"this workspace is suspended", status_code=403, media_type="text/plain"
)
if instance.get("status") != store.ST_RUNNING:
return Response(
"this workspace is not running", status_code=409, media_type="text/plain"
)
host, port = provision.editor_target(instance)
if not host or not port:
return Response(
"the editor has no reachable port", status_code=502, media_type="text/plain"
)
activity.touch(instance["uid"])
prefix = f"/projects/{slug}/containers/instances/{uid}/code"
return await forward.proxy_http(request, host, port, path, prefix=prefix)
@router.websocket("/{slug}/containers/instances/{uid}/code")
@router.websocket("/{slug}/containers/instances/{uid}/code/{path:path}")
async def editor_proxy_ws(
websocket: WebSocket, slug: str, uid: str, path: str = ""
):
from devplacepy.utils import get_current_user
user = get_current_user(websocket)
project = resolve_by_slug(get_table("projects"), slug)
instance = store.get_instance(uid)
if not user or not project or not instance or not instance.get("is_workspace"):
await websocket.close(code=1008)
return
if not can_manage_workspace(instance, project, user):
await websocket.close(code=1008)
return
if instance.get("suspended_at") or instance.get("status") != store.ST_RUNNING:
await websocket.close(code=1011)
return
host, port = provision.editor_target(instance)
if not host or not port:
await websocket.close(code=1011)
return
activity.touch(instance["uid"])
prefix = f"/projects/{slug}/containers/instances/{uid}/code"
await forward.proxy_ws(websocket, host, port, path, prefix=prefix)