|
# retoor <retoor@molodetz.nl>
|
|
|
|
import asyncio
|
|
import fcntl
|
|
import json
|
|
import logging
|
|
import os
|
|
import pty
|
|
import re
|
|
import signal
|
|
import struct
|
|
import termios
|
|
from typing import Annotated
|
|
|
|
from fastapi import Depends, APIRouter, Request, WebSocket, WebSocketDisconnect
|
|
from fastapi.responses import HTMLResponse, JSONResponse
|
|
|
|
from devplacepy.database import get_table, resolve_by_slug
|
|
from devplacepy.content import can_manage_instance, can_view_project_containers
|
|
from devplacepy.models import ContainerExecForm, ContainerInstanceForm
|
|
from devplacepy.responses import action_result, json_error, respond
|
|
from devplacepy.schemas import ContainersOut
|
|
from devplacepy.seo import base_seo_context
|
|
from devplacepy.utils import _user_from_session, is_admin, not_found, require_admin
|
|
from devplacepy.services.manager import service_manager
|
|
from devplacepy.services.containers import api, store
|
|
from devplacepy.services.containers.api import ContainerError
|
|
from devplacepy.services.containers.backend.base import WORKSPACE_MOUNT
|
|
from devplacepy.services.containers.runtime import get_backend
|
|
from devplacepy.services.audit import record as audit
|
|
|
|
from devplacepy.dependencies import json_or_form
|
|
from devplacepy.routers.projects.containers._shared import (
|
|
audit_instance,
|
|
fail,
|
|
instance_for,
|
|
manage_guard,
|
|
project_for,
|
|
slug_of,
|
|
)
|
|
|
|
logger = logging.getLogger(__name__)
|
|
router = APIRouter()
|
|
|
|
@router.get("/{project_slug}/containers", response_class=HTMLResponse)
|
|
async def containers_page(request: Request, project_slug: str):
|
|
user = require_admin(request)
|
|
project = project_for(project_slug, user)
|
|
slug = slug_of(project)
|
|
return respond(
|
|
request,
|
|
"containers.html",
|
|
{
|
|
**base_seo_context(
|
|
request,
|
|
title=f"Containers - {project.get('title', 'Project')}",
|
|
description="Container manager",
|
|
robots="noindex,nofollow",
|
|
breadcrumbs=[
|
|
{"name": "Home", "url": "/feed"},
|
|
{
|
|
"name": project.get("title", "Project"),
|
|
"url": f"/projects/{slug}",
|
|
},
|
|
{"name": "Containers", "url": f"/projects/{slug}/containers"},
|
|
],
|
|
),
|
|
"request": request,
|
|
"user": user,
|
|
"project": project,
|
|
"instances": store.list_instances(project["uid"]),
|
|
},
|
|
model=ContainersOut,
|
|
)
|
|
|
|
@router.get("/{project_slug}/containers/data")
|
|
async def containers_json(request: Request, project_slug: str):
|
|
user = require_admin(request)
|
|
project = project_for(project_slug, user)
|
|
return JSONResponse(
|
|
{
|
|
"instances": store.list_instances(project["uid"]),
|
|
}
|
|
)
|
|
|
|
# ---------------- instances ----------------
|
|
|
|
@router.post("/{project_slug}/containers/instances")
|
|
async def create_instance(
|
|
request: Request, project_slug: str, data: Annotated[ContainerInstanceForm, Depends(json_or_form(ContainerInstanceForm))]
|
|
):
|
|
user = require_admin(request)
|
|
project = project_for(project_slug, user)
|
|
try:
|
|
inst = await api.create_instance(
|
|
project,
|
|
name=data.name,
|
|
boot_command=data.boot_command,
|
|
boot_language=data.boot_language,
|
|
boot_script=data.boot_script,
|
|
run_as_uid=data.run_as_uid,
|
|
start_on_boot=data.start_on_boot,
|
|
env=data.env,
|
|
cpu_limit=data.cpu_limit,
|
|
mem_limit=data.mem_limit,
|
|
ports=data.ports,
|
|
volumes=data.volumes,
|
|
restart_policy=data.restart_policy,
|
|
autostart=data.autostart,
|
|
ingress_slug=data.ingress_slug,
|
|
ingress_port=data.ingress_port,
|
|
actor=("user", user["uid"]),
|
|
)
|
|
except ContainerError as exc:
|
|
return fail(exc)
|
|
audit_instance(
|
|
request, user, "container.instance.create", inst, project,
|
|
summary=f"admin {user['username']} created container instance {inst.get('name')} in project {project.get('title')}",
|
|
metadata={"ports": data.ports, "ingress_slug": data.ingress_slug, "restart_policy": data.restart_policy},
|
|
)
|
|
return action_result(
|
|
request, f"/projects/{slug_of(project)}/containers", data={"instance": inst}
|
|
)
|
|
|
|
@router.get("/{project_slug}/containers/instances/{uid}")
|
|
async def instance_detail(request: Request, project_slug: str, uid: str):
|
|
user = require_admin(request)
|
|
project = project_for(project_slug, user)
|
|
inst = instance_for(project, uid)
|
|
return JSONResponse(
|
|
{
|
|
"instance": inst,
|
|
"events": store.list_events(uid),
|
|
"schedules": store.list_schedules(uid),
|
|
"stats": api.instance_stats(uid),
|
|
"runtime": api.instance_runtime(inst),
|
|
}
|
|
)
|
|
|
|
_ACTIONS = {
|
|
"start": store.DESIRED_RUNNING,
|
|
"stop": store.DESIRED_STOPPED,
|
|
"pause": store.DESIRED_PAUSED,
|
|
"resume": store.DESIRED_RUNNING,
|
|
}
|
|
|
|
@router.post("/{project_slug}/containers/instances/{uid}/delete")
|
|
async def delete_instance(request: Request, project_slug: str, uid: str):
|
|
user = require_admin(request)
|
|
project = project_for(project_slug, user)
|
|
inst = instance_for(project, uid)
|
|
denied = manage_guard(request, user, project, inst, "container.instance.delete")
|
|
if denied:
|
|
return denied
|
|
api.mark_for_removal(inst, actor=("user", user["uid"]))
|
|
audit_instance(
|
|
request, user, "container.instance.delete", inst, project,
|
|
summary=f"admin {user['username']} removed container instance {inst.get('name')}",
|
|
)
|
|
return action_result(
|
|
request, f"/projects/{slug_of(project)}/containers", data={"uid": uid}
|
|
)
|
|
|
|
@router.post("/{project_slug}/containers/instances/{uid}/exec")
|
|
async def instance_exec(
|
|
request: Request,
|
|
project_slug: str,
|
|
uid: str,
|
|
data: Annotated[ContainerExecForm, Depends(json_or_form(ContainerExecForm))],
|
|
):
|
|
user = require_admin(request)
|
|
project = project_for(project_slug, user)
|
|
inst = instance_for(project, uid)
|
|
denied = manage_guard(request, user, project, inst, "container.instance.exec")
|
|
if denied:
|
|
return denied
|
|
if not inst.get("container_id"):
|
|
return json_error(400, "instance is not running")
|
|
result = await get_backend().exec(
|
|
inst["container_id"], ["/bin/sh", "-c", data.command]
|
|
)
|
|
store.record_event(
|
|
inst,
|
|
"exec",
|
|
"user",
|
|
user["uid"],
|
|
{"command": data.command, "exit_code": result.exit_code},
|
|
)
|
|
audit_instance(
|
|
request, user, "container.instance.exec", inst, project,
|
|
summary=f"admin {user['username']} ran a command in instance {inst.get('name')}",
|
|
metadata={"command": data.command, "exit_code": result.exit_code},
|
|
)
|
|
return JSONResponse({"exit_code": result.exit_code, "output": result.output})
|
|
|
|
@router.get("/{project_slug}/containers/instances/{uid}/logs")
|
|
async def instance_logs(request: Request, project_slug: str, uid: str, tail: int = 200):
|
|
user = require_admin(request)
|
|
project = project_for(project_slug, user)
|
|
inst = instance_for(project, uid)
|
|
if not inst.get("container_id"):
|
|
return JSONResponse({"logs": ""})
|
|
lines: list = []
|
|
|
|
async def collect(line: str) -> None:
|
|
lines.append(line)
|
|
|
|
await get_backend().logs(
|
|
inst["container_id"], follow=False, tail=max(1, min(tail, 2000)), on_log=collect
|
|
)
|
|
return JSONResponse({"logs": "\n".join(lines)})
|
|
|
|
@router.get("/{project_slug}/containers/instances/{uid}/metrics")
|
|
async def instance_metrics(request: Request, project_slug: str, uid: str):
|
|
user = require_admin(request)
|
|
project = project_for(project_slug, user)
|
|
instance_for(project, uid)
|
|
return JSONResponse(
|
|
{"metrics": store.recent_metrics(uid), "stats": api.instance_stats(uid)}
|
|
)
|
|
|
|
@router.post("/{project_slug}/containers/instances/{uid}/sync")
|
|
async def instance_sync(request: Request, project_slug: str, uid: str):
|
|
user = require_admin(request)
|
|
project = project_for(project_slug, user)
|
|
inst = instance_for(project, uid)
|
|
denied = manage_guard(request, user, project, inst, "container.instance.sync")
|
|
if denied:
|
|
return denied
|
|
try:
|
|
counts = await api.sync_workspace(inst, user)
|
|
except ContainerError as exc:
|
|
return fail(exc)
|
|
audit_instance(
|
|
request, user, "container.instance.sync", inst, project,
|
|
summary=f"admin {user['username']} synced the workspace of instance {inst.get('name')}",
|
|
metadata=counts,
|
|
)
|
|
return action_result(
|
|
request, f"/projects/{slug_of(project)}/containers", data=counts
|
|
)
|
|
|
|
@router.post("/{project_slug}/containers/instances/{uid}/start")
|
|
@router.post("/{project_slug}/containers/instances/{uid}/stop")
|
|
@router.post("/{project_slug}/containers/instances/{uid}/pause")
|
|
@router.post("/{project_slug}/containers/instances/{uid}/resume")
|
|
@router.post("/{project_slug}/containers/instances/{uid}/restart")
|
|
async def instance_action(request: Request, project_slug: str, uid: str):
|
|
user = require_admin(request)
|
|
project = project_for(project_slug, user)
|
|
inst = instance_for(project, uid)
|
|
action = request.url.path.rsplit("/", 1)[-1]
|
|
denied = manage_guard(request, user, project, inst, f"container.instance.{action}")
|
|
if denied:
|
|
return denied
|
|
actor = ("user", user["uid"])
|
|
if action == "restart":
|
|
api.request_restart(inst, actor=actor)
|
|
else:
|
|
api.set_desired_state(inst, _ACTIONS[action], actor=actor)
|
|
audit_instance(
|
|
request, user, f"container.instance.{action}", inst, project,
|
|
summary=f"admin {user['username']} {action} container instance {inst.get('name')}",
|
|
)
|
|
return action_result(
|
|
request,
|
|
f"/projects/{slug_of(project)}/containers",
|
|
data={"uid": uid, "action": action},
|
|
)
|
|
|
|
@router.websocket("/{project_slug}/containers/instances/{uid}/exec/ws")
|
|
async def instance_exec_ws(websocket: WebSocket, project_slug: str, uid: str):
|
|
await websocket.accept()
|
|
user = _user_from_session(websocket)
|
|
if not user or not is_admin(user):
|
|
await websocket.close(code=1008)
|
|
return
|
|
if not service_manager.owns_lock():
|
|
await websocket.close(code=1013)
|
|
return
|
|
project = resolve_by_slug(get_table("projects"), project_slug)
|
|
if not project or not can_view_project_containers(project, user):
|
|
await websocket.close(code=1011)
|
|
return
|
|
inst = store.get_instance(uid)
|
|
if (
|
|
not inst
|
|
or inst["project_uid"] != project["uid"]
|
|
or not inst.get("container_id")
|
|
):
|
|
await websocket.close(code=1011)
|
|
return
|
|
if not can_manage_instance(inst, project, user):
|
|
audit.record(
|
|
websocket,
|
|
"container.instance.shell.open",
|
|
user=user,
|
|
target_type="instance",
|
|
target_uid=inst["uid"],
|
|
target_label=inst.get("name"),
|
|
summary=f"admin {user['username']} denied shell access on instance {inst.get('name')}",
|
|
result="denied",
|
|
links=[audit.instance(inst["uid"], inst.get("name"))],
|
|
)
|
|
await websocket.close(code=1008)
|
|
return
|
|
if inst.get("status") != "running":
|
|
await websocket.send_text(
|
|
"\r\n[devplace] This container is "
|
|
f"{inst.get('status') or 'not running'}. Start it from the container page, "
|
|
"then open the terminal.\r\n"
|
|
)
|
|
await websocket.close(code=1011)
|
|
return
|
|
|
|
master, slave = pty.openpty()
|
|
session = (websocket.query_params.get("session") or "").strip()
|
|
if session and re.match(r"^[A-Za-z0-9_-]{1,64}$", session):
|
|
shell_cmd = (
|
|
f"if command -v tmux >/dev/null 2>&1; then exec tmux new-session -A -s {session}; "
|
|
"elif command -v bash >/dev/null 2>&1; then exec bash; else exec sh; fi"
|
|
)
|
|
else:
|
|
shell_cmd = (
|
|
"if command -v bash >/dev/null 2>&1; then exec bash; else exec sh; fi"
|
|
)
|
|
proc = await asyncio.create_subprocess_exec(
|
|
"docker",
|
|
"exec",
|
|
"-it",
|
|
"-w",
|
|
WORKSPACE_MOUNT,
|
|
inst["container_id"],
|
|
"/bin/sh",
|
|
"-c",
|
|
shell_cmd,
|
|
stdin=slave,
|
|
stdout=slave,
|
|
stderr=slave,
|
|
start_new_session=True,
|
|
)
|
|
os.close(slave)
|
|
loop = asyncio.get_event_loop()
|
|
audit.record(
|
|
websocket,
|
|
"container.instance.shell.open",
|
|
user=user,
|
|
target_type="instance",
|
|
target_uid=inst["uid"],
|
|
target_label=inst.get("name"),
|
|
metadata={"session": session or None},
|
|
summary=f"admin {user['username']} opened an interactive shell on instance {inst.get('name')}",
|
|
links=[audit.instance(inst["uid"], inst.get("name"))],
|
|
)
|
|
|
|
async def pump_out():
|
|
try:
|
|
while True:
|
|
data = await loop.run_in_executor(None, os.read, master, 4096)
|
|
if not data:
|
|
break
|
|
await websocket.send_text(data.decode("utf-8", "replace"))
|
|
except Exception:
|
|
pass
|
|
|
|
def set_winsize(rows: int, cols: int) -> None:
|
|
try:
|
|
fcntl.ioctl(
|
|
master, termios.TIOCSWINSZ, struct.pack("HHHH", rows, cols, 0, 0)
|
|
)
|
|
except OSError:
|
|
pass
|
|
try:
|
|
proc.send_signal(signal.SIGWINCH)
|
|
except (ProcessLookupError, OSError, ValueError):
|
|
pass
|
|
|
|
reader = asyncio.create_task(pump_out())
|
|
try:
|
|
while True:
|
|
text = await websocket.receive_text()
|
|
if text[:1] == "{":
|
|
try:
|
|
control = json.loads(text)
|
|
except ValueError:
|
|
control = None
|
|
if isinstance(control, dict) and control.get("type") == "resize":
|
|
rows = int(control.get("rows") or 0)
|
|
cols = int(control.get("cols") or 0)
|
|
if rows > 0 and cols > 0:
|
|
await loop.run_in_executor(None, set_winsize, rows, cols)
|
|
continue
|
|
await loop.run_in_executor(None, os.write, master, text.encode("utf-8"))
|
|
except WebSocketDisconnect:
|
|
pass
|
|
except Exception:
|
|
pass
|
|
finally:
|
|
reader.cancel()
|
|
try:
|
|
proc.kill()
|
|
except ProcessLookupError:
|
|
pass
|
|
os.close(master)
|
|
audit.record(
|
|
websocket,
|
|
"container.instance.shell.close",
|
|
user=user,
|
|
target_type="instance",
|
|
target_uid=inst["uid"],
|
|
target_label=inst.get("name"),
|
|
summary=f"admin {user['username']} closed the shell on instance {inst.get('name')}",
|
|
links=[audit.instance(inst["uid"], inst.get("name"))],
|
|
)
|