iUUUUpdatexz
DevPlace CI / test (push) Failing after 1h2m3s

This commit is contained in:
2026-08-07 10:53:43 +02:00
parent b777a5b9d0
commit 21f6ae0615
57 changed files with 4050 additions and 165 deletions
+2
View File
@@ -8,6 +8,7 @@ from devplacepy.routers.admin import (
backups,
bots,
containers,
workspaces,
devii_tasks,
game,
gateway_configs,
@@ -42,3 +43,4 @@ router.include_router(devii_tasks.router)
router.include_router(game.router)
router.include_router(services.router, prefix="/services")
router.include_router(containers.router, prefix="/containers")
router.include_router(workspaces.router)
+2 -1
View File
@@ -10,6 +10,7 @@ from devplacepy.database import (
build_pagination,
get_post_counts_by_user_uids,
invalidate_admins_cache,
is_account_active,
)
from devplacepy.utils import (
require_admin,
@@ -196,7 +197,7 @@ async def admin_user_toggle(request: Request, uid: str):
if _is_senior_admin(admin, user):
return _deny_senior(request, admin, uid, user, "admin.user.active.disable")
if user:
new_state = not user.get("is_active", True)
new_state = not is_account_active(user)
users.update({"uid": uid, "is_active": new_state}, ["uid"])
clear_user_cache(uid)
logger.info(
+282
View File
@@ -0,0 +1,282 @@
# retoor <retoor@molodetz.nl>
from typing import Annotated
from fastapi import APIRouter, Depends, Request
from fastapi.responses import JSONResponse
from devplacepy.database import db, get_table, get_users_by_uids
from devplacepy.dependencies import json_or_form
from devplacepy.models import (
WorkspaceFlagForm,
WorkspaceQuotaForm,
WorkspaceSuspendForm,
)
from devplacepy.responses import action_result, json_error, respond
from devplacepy.schemas import AdminWorkspacesOut
from devplacepy.seo import base_seo_context
from devplacepy.services.audit import record as audit
from devplacepy.services.containers import store
from devplacepy.services.containers.workspace import flags, provision, quota, tunnels
from devplacepy.utils import create_notification, generate_uid, not_found, require_admin
router = APIRouter()
def _decorate(rows: list[dict]) -> list[dict]:
owner_uids = {row.get("workspace_owner_uid") for row in rows if row.get("workspace_owner_uid")}
owners = get_users_by_uids(list(owner_uids)) if owner_uids else {}
projects = {}
if "projects" in db.tables:
project_uids = {row.get("project_uid") for row in rows if row.get("project_uid")}
for uid in project_uids:
found = get_table("projects").find_one(uid=uid)
if found:
projects[uid] = found
decorated = []
for row in rows:
view = provision.view(row)
owner = owners.get(row.get("workspace_owner_uid", "")) or {}
project = projects.get(row.get("project_uid", "")) or {}
view["owner_username"] = owner.get("username", "")
view["project_title"] = project.get("title", "")
view["project_slug"] = project.get("slug", "") or project.get("uid", "")
decorated.append(view)
return decorated
def _all_workspaces() -> list[dict]:
return list(get_table("instances").find(is_workspace=1, deleted_at=None))
def _instance_or_404(uid: str) -> dict:
instance = store.get_instance(uid)
if not instance or not instance.get("is_workspace"):
raise not_found("Workspace not found")
return instance
def _audit(request: Request, admin: dict, event_key: str, instance: dict, **extra):
audit.record(
request,
event_key,
user=admin,
target_type="instance",
target_uid=instance["uid"],
target_label=instance.get("name"),
summary=f"{admin['username']} {event_key} workspace {instance.get('name')}",
links=[audit.instance(instance["uid"], instance.get("name"))],
**extra,
)
@router.get("/workspaces")
async def admin_workspaces(request: Request):
admin = require_admin(request)
if not isinstance(admin, dict):
return admin
context = {
"workspaces": _decorate(_all_workspaces()),
"flags": flags.list_flags(),
"admin_section": "workspaces",
"user": admin,
**base_seo_context(
request,
title="Workspaces - Admin",
description="Administer dev workspaces.",
robots="noindex,nofollow",
breadcrumbs=[
{"name": "Home", "url": "/feed"},
{"name": "Admin", "url": "/admin"},
{"name": "Workspaces", "url": "/admin/workspaces"},
],
),
}
return respond(request, "admin_workspaces.html", context, model=AdminWorkspacesOut)
@router.get("/workspaces/data")
async def admin_workspaces_data(request: Request):
admin = require_admin(request)
if not isinstance(admin, dict):
return admin
return JSONResponse(
{"workspaces": _decorate(_all_workspaces()), "flags": flags.list_flags()}
)
@router.post("/workspaces/{uid}/suspend")
async def admin_workspace_suspend(
request: Request,
uid: str,
data: Annotated[WorkspaceSuspendForm, Depends(json_or_form(WorkspaceSuspendForm))],
):
admin = require_admin(request)
if not isinstance(admin, dict):
return admin
instance = _instance_or_404(uid)
reason = (data.reason or "").strip()
if not reason:
return json_error("a reason is required and is shown to the owner", 400)
provision.suspend(instance, admin["uid"], reason)
_audit(request, admin, "container.workspace.suspend", instance, metadata={"reason": reason})
owner = instance.get("workspace_owner_uid", "")
if owner:
create_notification(
owner,
"workspace",
f"Workspace {instance.get('name', '')} was suspended: {reason}",
instance["uid"],
"/projects",
)
return action_result(request, "/admin/workspaces")
@router.post("/workspaces/{uid}/unsuspend")
async def admin_workspace_unsuspend(request: Request, uid: str):
admin = require_admin(request)
if not isinstance(admin, dict):
return admin
instance = _instance_or_404(uid)
provision.unsuspend(instance)
_audit(request, admin, "container.workspace.unsuspend", instance)
owner = instance.get("workspace_owner_uid", "")
if owner:
create_notification(
owner,
"workspace",
f"Workspace {instance.get('name', '')} is available again.",
instance["uid"],
"/projects",
)
return action_result(request, "/admin/workspaces")
@router.post("/workspaces/{uid}/stop")
async def admin_workspace_stop(request: Request, uid: str):
admin = require_admin(request)
if not isinstance(admin, dict):
return admin
instance = _instance_or_404(uid)
provision.stop(instance)
_audit(request, admin, "container.workspace.stop", instance)
return action_result(request, "/admin/workspaces")
@router.post("/workspaces/{uid}/start")
async def admin_workspace_start(request: Request, uid: str):
admin = require_admin(request)
if not isinstance(admin, dict):
return admin
instance = _instance_or_404(uid)
store.update_instance(instance["uid"], {"desired_state": "running"})
_audit(request, admin, "container.workspace.resume", instance)
return action_result(request, "/admin/workspaces")
@router.post("/workspaces/{uid}/delete")
async def admin_workspace_delete(request: Request, uid: str):
admin = require_admin(request)
if not isinstance(admin, dict):
return admin
instance = _instance_or_404(uid)
for row in tunnels.list_for_instance(instance["uid"]):
tunnels.soft_delete(row["uid"], admin["uid"])
store.delete_instance(instance["uid"], admin["uid"])
_audit(request, admin, "container.workspace.delete", instance)
return action_result(request, "/admin/workspaces")
@router.post("/workspaces/{uid}/flag")
async def admin_workspace_flag(
request: Request,
uid: str,
data: Annotated[WorkspaceFlagForm, Depends(json_or_form(WorkspaceFlagForm))],
):
admin = require_admin(request)
if not isinstance(admin, dict):
return admin
instance = _instance_or_404(uid)
row = flags.raise_flag(
instance, data.kind or flags.KIND_MANUAL, data.severity, data.detail
)
_audit(request, admin, "container.workspace.flag.raise", instance,
metadata={"kind": data.kind, "severity": data.severity})
owner = instance.get("workspace_owner_uid", "")
if owner:
create_notification(
owner,
"workspace",
f"Workspace {instance.get('name', '')} was flagged: {data.detail or data.kind}",
instance["uid"],
"/projects",
)
return action_result(request, "/admin/workspaces", data=row)
@router.post("/workspaces/flags/{flag_uid}/resolve")
async def admin_flag_resolve(request: Request, flag_uid: str, status: str = "resolved"):
admin = require_admin(request)
if not isinstance(admin, dict):
return admin
if not flags.set_status(flag_uid, status, admin["uid"]):
raise not_found("Flag not found")
event = (
"container.workspace.flag.dismiss"
if status == "dismissed"
else "container.workspace.flag.resolve"
)
audit.record(
request,
event,
user=admin,
target_type="workspace_flag",
target_uid=flag_uid,
summary=f"{admin['username']} set flag {flag_uid} to {status}",
)
return action_result(request, "/admin/workspaces")
@router.post("/workspaces/quota")
async def admin_workspace_quota(
request: Request,
data: Annotated[WorkspaceQuotaForm, Depends(json_or_form(WorkspaceQuotaForm))],
):
admin = require_admin(request)
if not isinstance(admin, dict):
return admin
if not data.owner_id:
return json_error("owner_id is required", 400)
table = get_table(quota.RULES_TABLE)
existing = table.find_one(
owner_kind="user", owner_id=data.owner_id, deleted_at=None
)
payload = {key: getattr(data, key) for key in quota.RULE_COLUMNS}
if existing:
table.update({"uid": existing["uid"], "label": data.label, **payload}, ["uid"])
uid = existing["uid"]
else:
uid = generate_uid()
table.insert(
{
"uid": uid,
"owner_kind": "user",
"owner_id": data.owner_id,
"label": data.label,
"created_at": "",
"updated_at": "",
"deleted_at": None,
"deleted_by": None,
**payload,
}
)
audit.record(
request,
"container.workspace.settings.update",
user=admin,
target_type="user",
target_uid=data.owner_id,
summary=f"{admin['username']} updated workspace quota",
metadata=payload,
)
return action_result(request, "/admin/workspaces", data={"uid": uid, **payload})
+2 -2
View File
@@ -4,7 +4,7 @@ import logging
from typing import Annotated
from fastapi import Depends, APIRouter, Request
from fastapi.responses import HTMLResponse
from devplacepy.database import get_table, get_int_setting
from devplacepy.database import get_table, get_int_setting, is_account_active
from devplacepy.templating import templates
from devplacepy.utils import (
verify_password_async,
@@ -59,7 +59,7 @@ async def login(request: Request, data: Annotated[LoginForm, Depends(json_or_for
if not user or not await verify_password_async(password, user["password_hash"]):
errors.append("Invalid email or password")
elif not user.get("is_active", True):
elif not is_account_active(user):
errors.append("Account is deactivated")
if errors:
+2 -2
View File
@@ -6,7 +6,7 @@ from typing import Annotated
from fastapi import Depends, APIRouter, Request
from fastapi.responses import JSONResponse
from devplacepy.database import get_table
from devplacepy.database import get_table, is_account_active
from devplacepy.utils import verify_password_async, get_current_user
from devplacepy.models import LoginForm
from devplacepy.dependencies import json_or_form
@@ -59,7 +59,7 @@ async def token(
status_code=401,
)
if not user.get("is_active", True):
if not is_account_active(user):
audit.record(
request,
"auth.token.failure",
+2 -2
View File
@@ -9,7 +9,7 @@ from fastapi.responses import Response
from devplacepy.cache import TTLCache
from devplacepy.config import SECONDS_PER_DAY
from devplacepy.database import get_table, get_setting
from devplacepy.database import get_table, get_setting, is_account_active
from devplacepy.utils import verify_password_async, register_account_async
from devplacepy.services.audit import record as audit
from devplacepy.services.devrant.params import merge_params
@@ -42,7 +42,7 @@ async def auth_token(request: Request):
)
if (
not user
or not user.get("is_active", True)
or not is_account_active(user)
or not await verify_password_async(password, user["password_hash"])
):
audit.record(
@@ -2,8 +2,9 @@
from fastapi import APIRouter
from devplacepy.routers.projects.containers import instances, schedules
from devplacepy.routers.projects.containers import instances, schedules, workspace
router = APIRouter()
router.include_router(instances.router)
router.include_router(schedules.router)
router.include_router(workspace.router)
@@ -0,0 +1,278 @@
# retoor <retoor@molodetz.nl>
from typing import Annotated
from fastapi import APIRouter, 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.models import 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, forward, store
from devplacepy.services.containers.workspace import provision, quota, tunnels
from devplacepy.services.containers.workspace.provision import WorkspaceError
from devplacepy.utils import not_found, require_user
from ._shared import audit_instance
router = APIRouter()
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"])
context = {
"project": project,
"workspace": 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,
"editor_url": (
f"/projects/{slug}/containers/instances/{instance['uid']}/code/"
if instance
else ""
),
"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(str(error), 400)
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=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("Not allowed to manage this workspace", 403)
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("Not allowed to manage this workspace", 403)
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/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("Not allowed to manage this workspace", 403)
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("Not allowed to manage this workspace", 403)
if data.container_port <= 0:
return json_error("container_port must be between 1 and 65535", 400)
limits = quota.resolve(instance.get("workspace_owner_uid", ""), instance)
if limits.max_tunnels and tunnels.count_for_instance(
instance["uid"]
) >= limits.max_tunnels:
return json_error(f"tunnel limit reached ({limits.max_tunnels})", 400)
row = tunnels.create(instance, data.label, data.container_port, user["uid"])
if not row:
return json_error("could not create tunnel", 400)
audit_instance(
request,
user,
"container.tunnel.create",
instance,
project,
metadata={"hostname": row["hostname"], "port": data.container_port},
)
provision.write_manifest(instance)
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("Not allowed to manage this workspace", 403)
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("Not allowed to open this workspace", 403)
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)
if instance.get("status") != store.ST_RUNNING:
return Response("this workspace is not running", status_code=409)
host, port = provision.editor_target(instance)
if not host or not port:
return Response("the editor has no reachable port", status_code=502)
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"])
await forward.proxy_ws(websocket, host, port, path)
+6 -137
View File
@@ -1,70 +1,18 @@
# retoor <retoor@molodetz.nl>
import asyncio
import logging
import httpx
import websockets
from fastapi import APIRouter, Request, WebSocket
from starlette.responses import Response
from devplacepy.services.containers import api, store
from devplacepy.utils import not_found
from devplacepy.services.audit import record as audit
from devplacepy.services.containers import api, forward, store
from devplacepy.utils import not_found
logger = logging.getLogger(__name__)
router = APIRouter()
HOP_HEADERS = {
"connection",
"keep-alive",
"proxy-authenticate",
"proxy-authorization",
"te",
"trailers",
"transfer-encoding",
"upgrade",
"host",
"content-length",
"content-encoding",
}
METHODS = ["GET", "POST", "PUT", "DELETE", "PATCH", "OPTIONS", "HEAD"]
def _forward_headers(request: Request, prefix: str) -> dict:
headers = {k: v for k, v in request.headers.items() if k.lower() not in HOP_HEADERS}
headers["X-Forwarded-Prefix"] = prefix
headers["X-Script-Name"] = prefix
headers["X-Forwarded-Host"] = request.headers.get(
"host", request.url.hostname or ""
)
headers["X-Forwarded-Proto"] = request.headers.get(
"x-forwarded-proto", request.url.scheme
)
headers["Accept-Encoding"] = "identity"
return headers
def _inject_base(body: bytes, prefix: str) -> bytes:
lowered = body.lower()
if b"<base" in lowered:
return body
tag = f'<base href="{prefix}/">'.encode()
head = lowered.find(b"<head")
anchor = (
lowered.find(b">", head)
if head != -1
else lowered.find(b">", lowered.find(b"<html"))
)
if anchor == -1:
return tag + body
return body[: anchor + 1] + tag + body[anchor + 1 :]
def _rewrite_location(value: str, prefix: str) -> str:
if value.startswith("/") and not value.startswith("//"):
return prefix + value
return value
METHODS = forward.METHODS
def _resolve(slug: str):
@@ -95,41 +43,9 @@ async def proxy_http(request: Request, slug: str, path: str = ""):
summary=f"request proxied to instance {instance.get('name')} via ingress {slug}",
links=[audit.instance(instance["uid"], instance.get("name"))],
)
prefix = f"/p/{slug}"
url = f"http://{host}:{port}/{path}"
headers = _forward_headers(request, prefix)
body = await request.body()
try:
async with httpx.AsyncClient(timeout=60.0, follow_redirects=False) as client:
upstream = await client.request(
request.method,
url,
params=request.query_params,
headers=headers,
content=body,
)
except httpx.HTTPError as exc:
return Response(f"upstream error: {exc}", status_code=502)
out_headers = {
k: v
for k, v in upstream.headers.items()
if k.lower() not in HOP_HEADERS and k.lower() != "set-cookie"
}
if "location" in out_headers:
out_headers["location"] = _rewrite_location(out_headers["location"], prefix)
content_type = upstream.headers.get("content-type", "")
content = upstream.content
if "text/html" in content_type.lower():
content = _inject_base(content, prefix)
response = Response(
content=content,
status_code=upstream.status_code,
headers=out_headers,
media_type=content_type or None,
return await forward.proxy_http(
request, host, port, path, prefix=f"/p/{slug}", timeout=60.0
)
for cookie in upstream.headers.get_list("set-cookie"):
response.headers.append("set-cookie", cookie)
return response
@router.websocket("/{slug}")
@@ -139,9 +55,6 @@ async def proxy_ws(websocket: WebSocket, slug: str, path: str = ""):
if instance is None or not host or not port:
await websocket.close(code=1011)
return
upstream_url = f"ws://{host}:{port}/{path}"
if websocket.url.query:
upstream_url += f"?{websocket.url.query}"
await websocket.accept()
audit.record(
websocket,
@@ -154,48 +67,4 @@ async def proxy_ws(websocket: WebSocket, slug: str, path: str = ""):
summary=f"websocket proxied to instance {instance.get('name')} via ingress {slug}",
links=[audit.instance(instance["uid"], instance.get("name"))],
)
try:
async with websockets.connect(
upstream_url, open_timeout=10, max_size=None
) as upstream:
await _pump(websocket, upstream)
except Exception as exc:
logger.debug("ws proxy %s failed: %s", slug, exc)
try:
await websocket.close(code=1011)
except Exception:
pass
async def _pump(client_ws: WebSocket, upstream) -> None:
async def client_to_upstream():
try:
while True:
message = await client_ws.receive()
if message["type"] == "websocket.disconnect":
break
if message.get("text") is not None:
await upstream.send(message["text"])
elif message.get("bytes") is not None:
await upstream.send(message["bytes"])
except Exception:
pass
finally:
await upstream.close()
async def upstream_to_client():
try:
async for message in upstream:
if isinstance(message, (bytes, bytearray)):
await client_ws.send_bytes(bytes(message))
else:
await client_ws.send_text(message)
except Exception:
pass
finally:
try:
await client_ws.close()
except Exception:
pass
await asyncio.gather(client_to_upstream(), upstream_to_client())
await forward.proxy_ws(websocket, host, port, path, accepted=True)
+67
View File
@@ -0,0 +1,67 @@
# retoor <retoor@molodetz.nl>
import logging
from fastapi import APIRouter, Request, WebSocket
from starlette.responses import Response
from devplacepy.services.containers import activity, api, forward, store
from devplacepy.services.containers.workspace import naming, tunnels
logger = logging.getLogger(__name__)
router = APIRouter()
METHODS = forward.METHODS
def resolve(host: str):
if not naming.is_tunnel_host(host):
return None, None, None, None
row = tunnels.by_hostname(host)
if not row or row.get("status") not in tunnels.SERVING_STATUSES:
return None, None, None, None
instance = store.get_instance(row.get("instance_uid", ""))
if not instance or instance.get("deleted_at"):
return None, None, None, None
if instance.get("suspended_at"):
return row, instance, None, None
if instance.get("status") != store.ST_RUNNING:
return row, instance, None, None
gateway, _ = api.proxy_target(instance)
host_port = _published_host_port(instance, int(row.get("container_port") or 0))
return row, instance, gateway, host_port
def _published_host_port(instance: dict, container_port: int) -> int:
import json
for mapping in json.loads(instance.get("ports_json") or "[]"):
if int(mapping.get("container") or 0) == container_port:
return int(mapping.get("host") or 0)
return 0
async def handle_http(request: Request, path: str) -> Response:
host = request.headers.get("host", "")
row, instance, gateway, port = resolve(host)
if row is None:
return Response("no tunnel is published at this address", status_code=404)
if instance is not None and instance.get("suspended_at"):
return Response("this workspace is suspended", status_code=403)
if not gateway or not port:
return Response("the tunnel has no reachable port", status_code=502)
response = await forward.proxy_http(request, gateway, port, path)
size = len(response.body) if hasattr(response, "body") and response.body else 0
activity.touch(instance["uid"], egress_bytes=size)
tunnels.record_hit(row["uid"], size)
return response
async def handle_ws(websocket: WebSocket, path: str) -> None:
host = websocket.headers.get("host", "")
row, instance, gateway, port = resolve(host)
if row is None or instance is None or not gateway or not port:
await websocket.close(code=1011)
return
activity.touch(instance["uid"])
await forward.proxy_ws(websocket, gateway, port, path)