feat: add container manager CLI commands and docker-compose overlay for admin container lifecycle

This commit is contained in:
2026-06-09 04:41:27 +00:00
parent 317126efc2
commit c565e3b55c
83 changed files with 6078 additions and 84 deletions
+7 -7
View File
@@ -33,10 +33,14 @@ class BotBrowser:
self._page = await self._context.new_page()
async def close(self) -> None:
context, self._context = self._context, None
browser, self._browser = self._browser, None
playwright, self._playwright = self._playwright, None
self._page = None
for label, resource, shutdown in (
("context", self._context, lambda r: r.close()),
("browser", self._browser, lambda r: r.close()),
("playwright", self._playwright, lambda r: r.stop()),
("context", context, lambda r: r.close()),
("browser", browser, lambda r: r.close()),
("playwright", playwright, lambda r: r.stop()),
):
if resource is None:
continue
@@ -44,10 +48,6 @@ class BotBrowser:
await shutdown(resource)
except Exception as e:
logger.debug("Browser %s close error: %s", label, e)
self._page = None
self._context = None
self._browser = None
self._playwright = None
@property
def page(self) -> Page:
+15 -2
View File
@@ -135,8 +135,9 @@ class BotsService(BaseService):
await self._stop_fleet()
async def _stop_fleet(self) -> None:
for slot in list(self._fleet):
await self._stop_slot(slot)
slots = list(self._fleet)
if slots:
await asyncio.gather(*(self._stop_slot(slot) for slot in slots), return_exceptions=True)
async def _stop_slot(self, slot: int) -> None:
entry = self._fleet.pop(slot, None)
@@ -150,8 +151,20 @@ class BotsService(BaseService):
pass
except Exception as e:
self.log(f"bot{slot} stop error: {e}")
await self._close_bot(slot, entry.get("bot"))
self.log(f"Stopped bot{slot}")
async def _close_bot(self, slot: int, bot) -> None:
browser = getattr(bot, "b", None)
if browser is None:
return
try:
await asyncio.wait_for(browser.close(), timeout=15)
except (asyncio.CancelledError, asyncio.TimeoutError):
pass
except Exception as e:
self.log(f"bot{slot} browser close error: {e}")
def collect_metrics(self) -> dict:
rows = []
total_cost = 0.0
@@ -0,0 +1 @@
# retoor <retoor@molodetz.nl>
+366
View File
@@ -0,0 +1,366 @@
# retoor <retoor@molodetz.nl>
import asyncio
import json
import re
import socket
from pathlib import Path
import httpx
from devplacepy import config, project_files
from devplacepy.services.containers import store
from devplacepy.services.containers.locks import NUMBERING_LOCK
from devplacepy.services.containers.templates import DEFAULT_DOCKERFILE
from devplacepy.services.containers.backend.base import Mount, PortMapping, RunSpec
from devplacepy.services.devii.tasks.schedule import Schedule, next_run, now_utc, to_iso, from_iso
from devplacepy.services.jobs import queue
IMAGE_NAME_RE = re.compile(r"^[a-z0-9][a-z0-9._-]{0,62}$")
INGRESS_SLUG_RE = re.compile(r"^[a-z0-9][a-z0-9-]{0,62}$")
MEM_RE = re.compile(r"^\d+(\.\d+)?[bkmgBKMG]?$")
CPU_RE = re.compile(r"^\d+(\.\d+)?$")
ENV_KEY_RE = re.compile(r"^[A-Za-z_][A-Za-z0-9_]*$")
INSTANCE_LABEL = "devplace.instance"
PROJECT_LABEL = "devplace.project"
HOST_PORT_MIN = 20001
HOST_PORT_MAX = 65535
class ContainerError(ValueError):
pass
def validate_image_name(name: str) -> str:
name = (name or "").strip().lower()
if not IMAGE_NAME_RE.match(name):
raise ContainerError("name must be lowercase letters, digits, '.', '_' or '-' (max 63 chars)")
return name
def _validate_limits(cpu_limit: str, mem_limit: str) -> None:
if cpu_limit and not CPU_RE.match(str(cpu_limit)):
raise ContainerError("cpu limit must be a number, e.g. 1 or 1.5")
if mem_limit and not MEM_RE.match(str(mem_limit)):
raise ContainerError("memory limit must look like 512m, 1g, or a byte count")
def parse_ports(value) -> list:
ports = []
if not value:
return ports
items = value if isinstance(value, list) else str(value).replace(",", "\n").splitlines()
for item in items:
item = str(item).strip()
if not item:
continue
proto = "tcp"
if "/" in item:
item, proto = item.split("/", 1)
if ":" in item:
host, container = item.split(":", 1)
else:
host, container = "0", item
if not host.isdigit() or not container.isdigit():
raise ContainerError(f"port '{item}' must be numeric host:container or a bare container port")
ports.append(PortMapping(int(host), int(container), proto.strip() or "tcp"))
return ports
def used_host_ports() -> set:
ports = set()
for instance in store.all_instances():
for mapping in json.loads(instance.get("ports_json") or "[]"):
host = int(mapping.get("host") or 0)
if host:
ports.add(host)
return ports
def _host_port_free(port: int) -> bool:
with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as sock:
sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
try:
sock.bind(("0.0.0.0", port))
return True
except OSError:
return False
def allocate_host_port(reserved: set) -> int:
for port in range(HOST_PORT_MIN, HOST_PORT_MAX + 1):
if port in reserved:
continue
if _host_port_free(port):
return port
raise ContainerError(f"no free host port available in range {HOST_PORT_MIN}-{HOST_PORT_MAX}")
def assign_host_ports(port_list: list) -> list:
reserved = used_host_ports()
assigned = []
for mapping in port_list:
host = mapping.host
if not host:
host = allocate_host_port(reserved)
elif host in reserved:
raise ContainerError(f"host port {host} is already published by another instance")
reserved.add(host)
assigned.append(PortMapping(host, mapping.container, mapping.proto))
return assigned
def parse_env(value) -> dict:
env = {}
if not value:
return env
if isinstance(value, dict):
items = value.items()
else:
items = (line.split("=", 1) for line in str(value).splitlines() if "=" in line)
for key, val in items:
key = str(key).strip()
if not ENV_KEY_RE.match(key):
raise ContainerError(f"invalid environment variable name: {key}")
env[key] = str(val)
return env
# ---------------- dockerfiles + builds ----------------
async def enqueue_build(dockerfile: dict, version: dict, *, also_latest: bool = True, owner=("system", "system")) -> dict:
async with NUMBERING_LOCK:
head = store.get_dockerfile(dockerfile["uid"])
build_number = int(head.get("latest_build_number") or 0) + 1
store.update_dockerfile(head["uid"], {"latest_build_number": build_number})
image_tag = f"{head['name']}:{build_number}"
build = store.create_build(head, version, build_number, image_tag, also_latest)
job_uid = queue.enqueue("container_build", {"build_uid": build["uid"]}, owner[0], owner[1], f"build {image_tag}")
store.update_build(build["uid"], {"job_uid": job_uid})
return store.get_build(build["uid"])
async def create_dockerfile(project: dict, user: dict, *, name: str, description: str = "",
tags: str = "", content: str = None) -> dict:
name = validate_image_name(name)
for existing in store.list_dockerfiles(project["uid"]):
if existing["name"] == name:
raise ContainerError(f"a Dockerfile named '{name}' already exists in this project")
content = content if content is not None and content.strip() else DEFAULT_DOCKERFILE
dockerfile = store.create_dockerfile(project["uid"], user, name, description, tags)
version = store.create_version(dockerfile, content, user)
build = await enqueue_build(dockerfile, version, owner=("user", user["uid"]))
return {"dockerfile": store.get_dockerfile(dockerfile["uid"]), "version": version, "build": build}
async def save_version(dockerfile: dict, user: dict, content: str, *, force: bool = False) -> dict:
if not force and store.content_hash(content) == dockerfile.get("content_hash"):
return {"version": None, "build": None, "changed": False}
version = store.create_version(dockerfile, content, user)
build = await enqueue_build(store.get_dockerfile(dockerfile["uid"]), version, owner=("user", user["uid"]))
return {"version": version, "build": build, "changed": True}
# ---------------- instances ----------------
def validate_ingress(slug: str, port, port_list) -> tuple:
slug = (slug or "").strip().lower()
if not slug:
return "", 0
if not INGRESS_SLUG_RE.match(slug):
raise ContainerError("ingress slug must be lowercase letters, digits, or '-' (max 63 chars)")
for other in store.all_instances():
if other.get("ingress_slug") == slug:
raise ContainerError(f"ingress slug '{slug}' is already in use")
ingress_port = int(port) if port else 0
container_ports = {p.container for p in port_list}
if ingress_port and ingress_port not in container_ports:
raise ContainerError(f"ingress_port {ingress_port} must be one of the container ports you mapped")
if not ingress_port and len(container_ports) != 1:
raise ContainerError("set ingress_port to choose which mapped container port to publish")
return slug, ingress_port
async def create_instance(project: dict, dockerfile: dict, build: dict, *, name: str,
boot_command: str = "", env="", cpu_limit: str = "", mem_limit: str = "",
ports="", volumes="", restart_policy: str = "never",
autostart: bool = True, ingress_slug: str = "", ingress_port=None,
actor=("system", "system")) -> dict:
if build.get("status") != store.BUILD_SUCCESS:
raise ContainerError("the selected build has not completed successfully")
if restart_policy not in store.RESTART_POLICIES:
raise ContainerError(f"restart policy must be one of {', '.join(store.RESTART_POLICIES)}")
_validate_limits(cpu_limit, mem_limit)
port_list = assign_host_ports(parse_ports(ports))
env_map = parse_env(env)
name = (name or "").strip()
if not name:
raise ContainerError("instance name is required")
ingress_slug, ingress_port = validate_ingress(ingress_slug, ingress_port, port_list)
workspace = Path(config.CONTAINER_WORKSPACES_DIR) / project["uid"]
await asyncio.to_thread(project_files.export_to_dir, project["uid"], "", str(workspace))
row = {
"project_uid": project["uid"], "dockerfile_uid": dockerfile["uid"], "build_uid": build["uid"],
"name": name, "boot_command": boot_command or "",
"env_json": json.dumps(env_map),
"cpu_limit": str(cpu_limit or ""), "mem_limit": str(mem_limit or ""),
"ports_json": json.dumps([{"host": p.host, "container": p.container, "proto": p.proto} for p in port_list]),
"volumes_json": volumes if isinstance(volumes, str) else json.dumps(volumes or []),
"restart_policy": restart_policy,
"ingress_slug": ingress_slug, "ingress_port": ingress_port,
"desired_state": store.DESIRED_RUNNING if autostart else store.DESIRED_STOPPED,
"status": store.ST_CREATED,
"workspace_dir": str(workspace),
}
instance = store.create_instance(row)
store.record_event(instance, "created", actor[0], actor[1], {"build": build["uid"]})
return instance
def set_desired_state(instance: dict, desired: str, *, actor=("system", "system")) -> dict:
if desired not in (store.DESIRED_RUNNING, store.DESIRED_STOPPED, store.DESIRED_PAUSED):
raise ContainerError("desired state must be running, stopped, or paused")
store.update_instance(instance["uid"], {"desired_state": desired})
store.record_event(instance, f"desire_{desired}", actor[0], actor[1])
return store.get_instance(instance["uid"])
def request_restart(instance: dict, *, actor=("system", "system")) -> dict:
store.update_instance(instance["uid"], {"desired_state": store.DESIRED_RUNNING, "status": store.ST_RESTARTING})
store.record_event(instance, "restart", actor[0], actor[1])
return store.get_instance(instance["uid"])
def mark_for_removal(instance: dict, *, actor=("system", "system")) -> None:
store.update_instance(instance["uid"], {"desired_state": store.DESIRED_STOPPED, "status": store.ST_REMOVING})
store.record_event(instance, "remove", actor[0], actor[1])
def run_spec_for(instance: dict, image_tag: str) -> RunSpec:
env = json.loads(instance.get("env_json") or "{}")
ports = [PortMapping(p["host"], p["container"], p.get("proto", "tcp"))
for p in json.loads(instance.get("ports_json") or "[]")]
mounts = [Mount(instance["workspace_dir"], "/app", "rw")]
for extra in json.loads(instance.get("volumes_json") or "[]"):
if isinstance(extra, dict) and extra.get("host") and extra.get("container"):
mounts.append(Mount(extra["host"], extra["container"], extra.get("mode", "rw")))
command = []
boot = (instance.get("boot_command") or "").strip()
if boot:
command = ["/bin/sh", "-c", boot]
return RunSpec(
image=image_tag, name=instance["slug"],
labels={INSTANCE_LABEL: instance["uid"], PROJECT_LABEL: instance["project_uid"]},
env=env, cpu_limit=instance.get("cpu_limit", ""), mem_limit=instance.get("mem_limit", ""),
ports=ports, mounts=mounts, restart_policy=instance.get("restart_policy", "never"), command=command,
)
async def sync_workspace(instance: dict, user: dict) -> int:
workspace = instance.get("workspace_dir")
if not workspace:
raise ContainerError("instance has no workspace")
count = await asyncio.to_thread(project_files.import_from_dir, instance["project_uid"], workspace, user)
store.record_event(instance, "sync", "user", user["uid"], {"imported": count})
return count
def add_schedule(instance: dict, action: str, schedule: Schedule) -> dict:
if action not in ("start", "stop"):
raise ContainerError("schedule action must be start or stop")
first = schedule.first_run(now_utc())
return store.create_schedule(instance, action, schedule.columns(), to_iso(first))
# ---------------- aggregation ----------------
def _percentile(values: list, pct: float) -> float:
if not values:
return 0.0
ordered = sorted(values)
index = min(len(ordered) - 1, int(round((pct / 100.0) * (len(ordered) - 1))))
return float(ordered[index])
def dockerfile_stats(dockerfile_uid: str) -> dict:
builds = store.list_builds(dockerfile_uid, limit=1000)
done = [b for b in builds if b["status"] in (store.BUILD_SUCCESS, store.BUILD_FAILED)]
success = [b for b in done if b["status"] == store.BUILD_SUCCESS]
durations = [int(b.get("duration_ms") or 0) for b in success if b.get("duration_ms")]
return {
"total_builds": len(builds),
"success": len(success),
"failed": len(done) - len(success),
"success_rate": round(len(success) / len(done) * 100, 1) if done else 0.0,
"avg_build_ms": int(sum(durations) / len(durations)) if durations else 0,
"p95_build_ms": int(_percentile(durations, 95)),
}
def instance_stats(instance_uid: str) -> dict:
metrics = store.recent_metrics(instance_uid, limit=720)
cpu = [m.get("cpu_pct", 0) for m in metrics]
mem = [m.get("mem_bytes", 0) for m in metrics]
return {
"samples": len(metrics),
"cpu_avg": round(sum(cpu) / len(cpu), 2) if cpu else 0.0,
"cpu_p95": round(_percentile(cpu, 95), 2),
"mem_max": max(mem) if mem else 0,
"mem_avg": int(sum(mem) / len(mem)) if mem else 0,
}
def _port_reachable(host: str, port: int, timeout: float = 0.3) -> bool:
try:
with socket.create_connection((host, port), timeout=timeout):
return True
except OSError:
return False
def _http_probe(host: str, port: int, timeout: float = 1.0) -> str:
try:
with httpx.Client(timeout=timeout) as client:
response = client.get(f"http://{host}:{port}/")
return f"HTTP {response.status_code}"
except Exception as exc: # noqa: BLE001 - diagnostic, any failure is informative
return f"unreachable: {type(exc).__name__}"
def _ingress_host_port(instance: dict, port_maps: list) -> int:
ingress_port = int(instance.get("ingress_port") or 0)
if ingress_port:
for mapping in port_maps:
if int(mapping.get("container") or 0) == ingress_port:
return int(mapping.get("host") or 0)
return 0
return int(port_maps[0].get("host") or 0) if port_maps else 0
def instance_runtime(instance: dict) -> dict:
boot = (instance.get("boot_command") or "").strip()
port_maps = json.loads(instance.get("ports_json") or "[]")
host = config.CONTAINER_PROXY_HOST
ports = []
for mapping in port_maps:
host_port = int(mapping.get("host") or 0)
ports.append({
"container": int(mapping.get("container") or 0),
"host": host_port,
"proto": mapping.get("proto", "tcp"),
"reachable": _port_reachable(host, host_port) if host_port else False,
})
ingress_host_port = _ingress_host_port(instance, port_maps)
ingress_serving = _http_probe(host, ingress_host_port) if ingress_host_port else "no ingress port mapped"
return {
"command": boot or "image CMD (no boot_command set)",
"ports": ports,
"ingress_port": int(instance.get("ingress_port") or 0),
"ingress_serving": ingress_serving,
"status_ok_for_ingress": instance.get("status") == store.ST_RUNNING,
"restart_count": int(instance.get("restart_count") or 0),
"exit_code": instance.get("exit_code"),
"container_id": (instance.get("container_id") or "")[:12],
}
@@ -0,0 +1,23 @@
# retoor <retoor@molodetz.nl>
from devplacepy.services.containers.backend.base import (
Backend,
BuildResult,
ExecResult,
Mount,
PortMapping,
PsRow,
RunSpec,
StatsSample,
)
__all__ = [
"Backend",
"BuildResult",
"ExecResult",
"Mount",
"PortMapping",
"PsRow",
"RunSpec",
"StatsSample",
]
@@ -0,0 +1,118 @@
# retoor <retoor@molodetz.nl>
from __future__ import annotations
from abc import ABC, abstractmethod
from dataclasses import dataclass, field
from typing import Awaitable, Callable, Optional
LogCallback = Callable[[str], Awaitable[None]]
@dataclass
class PortMapping:
host: int
container: int
proto: str = "tcp"
@dataclass
class Mount:
host: str
container: str
mode: str = "rw"
@dataclass
class RunSpec:
image: str
name: str
labels: dict = field(default_factory=dict)
env: dict = field(default_factory=dict)
cpu_limit: str = ""
mem_limit: str = ""
ports: list = field(default_factory=list)
mounts: list = field(default_factory=list)
restart_policy: str = "no"
command: list = field(default_factory=list)
@dataclass
class BuildResult:
success: bool
image_id: str = ""
error: str = ""
@dataclass
class ExecResult:
exit_code: int
output: str = ""
@dataclass
class StatsSample:
cpu_pct: float = 0.0
mem_bytes: int = 0
mem_pct: float = 0.0
net_rx: int = 0
net_tx: int = 0
blk_read: int = 0
blk_write: int = 0
@dataclass
class PsRow:
container_id: str
name: str
state: str
status: str
exit_code: Optional[int]
labels: dict = field(default_factory=dict)
class Backend(ABC):
@abstractmethod
async def build(self, *, context_dir: str, dockerfile_text: str, tags: list,
on_log: Optional[LogCallback] = None, network: str = "") -> BuildResult: ...
@abstractmethod
async def run(self, spec: RunSpec) -> str: ...
@abstractmethod
async def start(self, cid: str) -> None: ...
@abstractmethod
async def stop(self, cid: str, timeout: int = 10) -> None: ...
@abstractmethod
async def restart(self, cid: str) -> None: ...
@abstractmethod
async def pause(self, cid: str) -> None: ...
@abstractmethod
async def unpause(self, cid: str) -> None: ...
@abstractmethod
async def rm(self, cid: str, force: bool = False) -> None: ...
@abstractmethod
async def exec(self, cid: str, cmd: list, *, tty: bool = False,
on_log: Optional[LogCallback] = None) -> ExecResult: ...
@abstractmethod
async def logs(self, cid: str, *, follow: bool = False, tail: int = 200,
on_log: Optional[LogCallback] = None) -> None: ...
@abstractmethod
async def stats_once(self, cids: list) -> dict: ...
@abstractmethod
async def ps(self, *, label_filter: str = "devplace.instance") -> list: ...
@abstractmethod
async def inspect(self, cid: str) -> dict: ...
@abstractmethod
async def image_prune(self) -> None: ...
@@ -0,0 +1,267 @@
# retoor <retoor@molodetz.nl>
from __future__ import annotations
import asyncio
import json
import logging
import re
import tempfile
from pathlib import Path
from typing import Optional
from devplacepy.services.containers.backend.base import (
Backend,
BuildResult,
ExecResult,
LogCallback,
PsRow,
RunSpec,
StatsSample,
)
logger = logging.getLogger("containers.docker")
DOCKER = "docker"
DEFAULT_TIMEOUT = 600
BUILD_TIMEOUT = 3600
_EXIT_RE = re.compile(r"Exited \((\d+)\)")
_SIZE_RE = re.compile(r"([0-9.]+)\s*([kKmMgGtT]?i?)B")
def build_run_argv(spec: RunSpec) -> list:
argv = [DOCKER, "run", "-d", "--name", spec.name]
for key, value in spec.labels.items():
argv += ["--label", f"{key}={value}"]
for key, value in spec.env.items():
argv += ["-e", f"{key}={value}"]
if spec.cpu_limit:
argv += ["--cpus", str(spec.cpu_limit)]
if spec.mem_limit:
argv += ["--memory", str(spec.mem_limit)]
for port in spec.ports:
argv += ["-p", f"{port.host}:{port.container}/{port.proto}"]
for mount in spec.mounts:
argv += ["-v", f"{mount.host}:{mount.container}:{mount.mode}"]
if spec.restart_policy and spec.restart_policy != "never":
policy = "no" if spec.restart_policy in ("no", "never") else spec.restart_policy
argv += ["--restart", policy]
argv.append(spec.image)
argv += list(spec.command)
return argv
def build_image_argv(context_dir: str, dockerfile_path: str, tags: list, iidfile: str, network: str = "") -> list:
argv = [DOCKER, "build", "--label", "devplace.build=1", "--iidfile", iidfile, "-f", dockerfile_path]
if network:
argv += ["--network", network]
for tag in tags:
argv += ["-t", tag]
argv.append(context_dir)
return argv
def parse_size(text: str) -> int:
match = _SIZE_RE.search(text or "")
if not match:
return 0
value = float(match.group(1))
unit = match.group(2).lower().rstrip("i")
factor = {"": 1, "k": 1024, "m": 1024 ** 2, "g": 1024 ** 3, "t": 1024 ** 4}.get(unit, 1)
return int(value * factor)
def parse_stats_line(row: dict) -> StatsSample:
def _pair(text: str) -> tuple:
parts = (text or "").split("/")
left = parse_size(parts[0]) if parts else 0
right = parse_size(parts[1]) if len(parts) > 1 else 0
return left, right
rx, tx = _pair(row.get("NetIO", ""))
rd, wr = _pair(row.get("BlockIO", ""))
return StatsSample(
cpu_pct=float((row.get("CPUPerc", "0%") or "0%").rstrip("%") or 0),
mem_bytes=_pair(row.get("MemUsage", ""))[0],
mem_pct=float((row.get("MemPerc", "0%") or "0%").rstrip("%") or 0),
net_rx=rx, net_tx=tx, blk_read=rd, blk_write=wr,
)
class DockerCliBackend(Backend):
def __init__(self, binary: str = DOCKER, timeout: int = DEFAULT_TIMEOUT,
build_timeout: int = BUILD_TIMEOUT) -> None:
self._binary = binary
self._timeout = timeout
self._build_timeout = build_timeout
async def _run(self, argv: list, *, timeout: Optional[int] = None) -> tuple:
proc = await asyncio.create_subprocess_exec(
*argv, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE)
try:
out, err = await asyncio.wait_for(proc.communicate(), timeout=timeout or self._timeout)
except asyncio.TimeoutError:
proc.kill()
raise RuntimeError(f"docker command timed out: {' '.join(argv[:3])}")
return proc.returncode, out.decode("utf-8", "replace"), err.decode("utf-8", "replace")
async def _stream(self, argv: list, on_log: Optional[LogCallback], *, timeout: Optional[int] = None) -> int:
proc = await asyncio.create_subprocess_exec(
*argv, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.STDOUT)
async def pump():
assert proc.stdout is not None
async for raw in proc.stdout:
line = raw.decode("utf-8", "replace").rstrip("\n")
if on_log is not None:
await on_log(line)
return await proc.wait()
try:
return await asyncio.wait_for(pump(), timeout=timeout or self._timeout)
except asyncio.TimeoutError:
proc.kill()
raise RuntimeError(f"docker stream timed out: {' '.join(argv[:3])}")
async def _check(self, argv: list) -> str:
code, out, err = await self._run(argv)
if code != 0:
raise RuntimeError(f"{' '.join(argv[:2])} failed ({code}): {(err or out).strip()[:500]}")
return out.strip()
async def build(self, *, context_dir, dockerfile_text, tags, on_log=None, network="") -> BuildResult:
ctx = Path(context_dir)
dockerfile_path = ctx / "Dockerfile.devplace"
iid_path = await asyncio.to_thread(self._prepare_context, dockerfile_path, dockerfile_text)
argv = build_image_argv(str(ctx), str(dockerfile_path), tags, iid_path, network)
try:
code = await self._stream(argv, on_log, timeout=self._build_timeout)
if code != 0:
return BuildResult(success=False, error=f"docker build exited {code}")
image_id = (await asyncio.to_thread(Path(iid_path).read_text)).strip()
return BuildResult(success=True, image_id=image_id)
finally:
await asyncio.to_thread(self._cleanup_context, iid_path, dockerfile_path)
@staticmethod
def _prepare_context(dockerfile_path: Path, dockerfile_text: str) -> str:
dockerfile_path.write_text(dockerfile_text, encoding="utf-8")
iid = tempfile.NamedTemporaryFile(suffix=".iid", delete=False)
iid.close()
return iid.name
@staticmethod
def _cleanup_context(iid_path: str, dockerfile_path: Path) -> None:
Path(iid_path).unlink(missing_ok=True)
dockerfile_path.unlink(missing_ok=True)
async def run(self, spec: RunSpec) -> str:
return await self._check(build_run_argv(spec))
async def start(self, cid: str) -> None:
await self._check([self._binary, "start", cid])
async def stop(self, cid: str, timeout: int = 10) -> None:
await self._check([self._binary, "stop", "-t", str(timeout), cid])
async def restart(self, cid: str) -> None:
await self._check([self._binary, "restart", cid])
async def pause(self, cid: str) -> None:
await self._check([self._binary, "pause", cid])
async def unpause(self, cid: str) -> None:
await self._check([self._binary, "unpause", cid])
async def rm(self, cid: str, force: bool = False) -> None:
argv = [self._binary, "rm"]
if force:
argv.append("-f")
argv.append(cid)
await self._check(argv)
async def exec(self, cid, cmd, *, tty=False, on_log=None) -> ExecResult:
argv = [self._binary, "exec"]
if tty:
argv.append("-t")
argv.append(cid)
argv += list(cmd)
if on_log is not None:
code = await self._stream(argv, on_log)
return ExecResult(exit_code=code)
code, out, err = await self._run(argv)
return ExecResult(exit_code=code or 0, output=(out + err))
async def logs(self, cid, *, follow=False, tail=200, on_log=None) -> None:
argv = [self._binary, "logs", "--tail", str(tail)]
if follow:
argv.append("--follow")
argv.append(cid)
await self._stream(argv, on_log)
async def stats_once(self, cids: list) -> dict:
if not cids:
return {}
argv = [self._binary, "stats", "--no-stream", "--format", "{{json .}}", *cids]
code, out, _ = await self._run(argv, timeout=30)
samples = {}
if code != 0:
return samples
for line in out.splitlines():
line = line.strip()
if not line:
continue
try:
row = json.loads(line)
except ValueError:
continue
name = row.get("Name") or row.get("ID") or ""
samples[name] = parse_stats_line(row)
return samples
async def ps(self, *, label_filter: str = "devplace.instance") -> list:
argv = [self._binary, "ps", "-a", "--filter", f"label={label_filter}", "--format", "{{json .}}"]
code, out, _ = await self._run(argv, timeout=30)
rows = []
if code != 0:
return rows
for line in out.splitlines():
line = line.strip()
if not line:
continue
try:
row = json.loads(line)
except ValueError:
continue
labels = {}
for pair in (row.get("Labels", "") or "").split(","):
if "=" in pair:
key, value = pair.split("=", 1)
labels[key] = value
status = row.get("Status", "") or ""
match = _EXIT_RE.search(status)
rows.append(PsRow(
container_id=row.get("ID", ""),
name=row.get("Names", ""),
state=(row.get("State", "") or "").lower(),
status=status,
exit_code=int(match.group(1)) if match else None,
labels=labels,
))
return rows
async def inspect(self, cid: str) -> dict:
code, out, _ = await self._run([self._binary, "inspect", cid], timeout=30)
if code != 0:
return {}
try:
data = json.loads(out)
return data[0] if isinstance(data, list) and data else {}
except ValueError:
return {}
async def image_prune(self) -> None:
try:
await self._run([self._binary, "image", "prune", "-f", "--filter", "label=devplace.build"], timeout=120)
except RuntimeError as exc:
logger.warning("image prune failed: %s", exc)
@@ -0,0 +1,99 @@
# retoor <retoor@molodetz.nl>
from __future__ import annotations
from devplacepy.services.containers.backend.base import (
Backend,
BuildResult,
ExecResult,
PsRow,
RunSpec,
StatsSample,
)
class FakeBackend(Backend):
def __init__(self) -> None:
self.containers: dict = {}
self.built_images: list = []
self.removed: list = []
self.pruned = 0
self._counter = 0
self.fail_build = False
async def build(self, *, context_dir, dockerfile_text, tags, on_log=None, network="") -> BuildResult:
if on_log is not None:
await on_log(f"FAKE build {tags} net={network}")
if self.fail_build:
return BuildResult(success=False, error="fake build failure")
self.built_images.append(list(tags))
return BuildResult(success=True, image_id=f"sha256:fake{len(self.built_images)}")
async def run(self, spec: RunSpec) -> str:
self._counter += 1
cid = f"fake{self._counter:012d}"
self.containers[cid] = {
"name": spec.name, "state": "running", "labels": dict(spec.labels),
"exit_code": None, "image": spec.image, "spec": spec,
}
return cid
def _set_state(self, cid: str, state: str, exit_code=None) -> None:
if cid in self.containers:
self.containers[cid]["state"] = state
if exit_code is not None:
self.containers[cid]["exit_code"] = exit_code
async def start(self, cid: str) -> None:
self._set_state(cid, "running")
async def stop(self, cid: str, timeout: int = 10) -> None:
self._set_state(cid, "exited", exit_code=0)
async def restart(self, cid: str) -> None:
self._set_state(cid, "running")
async def pause(self, cid: str) -> None:
self._set_state(cid, "paused")
async def unpause(self, cid: str) -> None:
self._set_state(cid, "running")
async def rm(self, cid: str, force: bool = False) -> None:
self.removed.append(cid)
self.containers.pop(cid, None)
async def exec(self, cid, cmd, *, tty=False, on_log=None) -> ExecResult:
line = f"FAKE exec {' '.join(cmd)}"
if on_log is not None:
await on_log(line)
return ExecResult(exit_code=0, output=line)
async def logs(self, cid, *, follow=False, tail=200, on_log=None) -> None:
if on_log is not None:
await on_log(f"FAKE logs for {cid}")
async def stats_once(self, cids: list) -> dict:
return {self.containers[cid]["name"]: StatsSample(cpu_pct=1.0, mem_bytes=1024)
for cid in cids if cid in self.containers}
async def ps(self, *, label_filter: str = "devplace.instance") -> list:
key, _, value = label_filter.partition("=")
rows = []
for cid, data in self.containers.items():
labels = data["labels"]
if key not in labels:
continue
if value and labels.get(key) != value:
continue
rows.append(PsRow(
container_id=cid, name=data["name"], state=data["state"],
status=data["state"], exit_code=data["exit_code"], labels=labels,
))
return rows
async def inspect(self, cid: str) -> dict:
return self.containers.get(cid, {})
async def image_prune(self) -> None:
self.pruned += 1
@@ -0,0 +1,101 @@
# retoor <retoor@molodetz.nl>
import asyncio
import shutil
import time
from pathlib import Path
from devplacepy import config, project_files
from devplacepy.services.base import ConfigField
from devplacepy.services.containers import store
from devplacepy.services.containers.runtime import get_backend
from devplacepy.services.jobs.base import JobService
LOG_CAP = 256 * 1024
def _cap(text: str) -> str:
return text[-LOG_CAP:]
class ContainerBuildService(JobService):
kind = "container_build"
default_enabled = False
title = "Container builds"
description = (
"Builds container images for project Dockerfiles asynchronously via the docker CLI, "
"streaming logs and recording every build in the permanent build history."
)
def __init__(self):
super().__init__(name="container_build", interval_seconds=2)
self.config_fields = self.config_fields + [
ConfigField("container_build_network", "Build network", type="str", default="host",
help="docker build --network value so the builder can reach the internet "
"(e.g. 'host' for PyPI access). Empty uses the docker default.",
group="Builds"),
]
async def process(self, job: dict) -> dict:
build_uid = job["payload"].get("build_uid")
build = store.get_build(build_uid)
if build is None:
raise RuntimeError("build row missing")
if build["status"] == store.BUILD_CANCELLED:
return {}
version = store.get_version(build["dockerfile_version_uid"])
dockerfile = store.get_dockerfile(build["dockerfile_uid"])
if version is None or dockerfile is None:
store.update_build(build_uid, {"status": store.BUILD_FAILED, "error": "version or dockerfile missing"})
return {}
store.update_build(build_uid, {"status": store.BUILD_BUILDING, "started_at": store.now()})
ctx = Path(config.CONTAINER_BUILD_CONTEXTS_DIR) / build_uid
await asyncio.to_thread(project_files.export_to_dir, build["project_uid"], "", str(ctx))
lines: list = []
async def on_log(line: str) -> None:
lines.append(line)
if len(lines) % 20 == 0:
store.update_build(build_uid, {"logs": _cap("\n".join(lines))})
tags = [build["image_tag"]]
if build.get("also_latest"):
tags.append(f"{dockerfile['name']}:latest")
network = str(self.get_config().get("container_build_network", "host") or "")
started = time.monotonic()
try:
result = await get_backend().build(
context_dir=str(ctx), dockerfile_text=version["content"], tags=tags, on_log=on_log, network=network)
except Exception as exc:
result = None
error = str(exc)
else:
error = result.error
finally:
await asyncio.to_thread(shutil.rmtree, str(ctx), ignore_errors=True)
duration_ms = int((time.monotonic() - started) * 1000)
logs = _cap("\n".join(lines))
latest = store.get_build(build_uid)
if latest and latest["status"] == store.BUILD_CANCELLED:
store.update_build(build_uid, {"logs": logs, "completed_at": store.now(), "duration_ms": duration_ms})
self.log(f"build {build['image_tag']} cancelled")
return {"item_count": len(lines)}
if result is not None and result.success:
store.update_build(build_uid, {
"status": store.BUILD_SUCCESS, "image_id": result.image_id, "logs": logs,
"completed_at": store.now(), "duration_ms": duration_ms,
})
self.log(f"build {build['image_tag']} succeeded ({duration_ms}ms)")
else:
store.update_build(build_uid, {
"status": store.BUILD_FAILED, "error": (error or "build failed")[:2000], "logs": logs,
"completed_at": store.now(), "duration_ms": duration_ms,
})
self.log(f"build {build['image_tag']} failed")
return {"item_count": len(lines)}
+5
View File
@@ -0,0 +1,5 @@
# retoor <retoor@molodetz.nl>
import asyncio
NUMBERING_LOCK = asyncio.Lock()
+16
View File
@@ -0,0 +1,16 @@
# retoor <retoor@molodetz.nl>
_backend = None
def get_backend():
global _backend
if _backend is None:
from devplacepy.services.containers.backend.docker_cli import DockerCliBackend
_backend = DockerCliBackend()
return _backend
def set_backend(backend) -> None:
global _backend
_backend = backend
+218
View File
@@ -0,0 +1,218 @@
# retoor <retoor@molodetz.nl>
import json
from devplacepy.services.base import BaseService, ConfigField
from devplacepy.services.containers import api, store
from devplacepy.services.containers.runtime import get_backend
from devplacepy.services.devii.tasks.schedule import next_run, now_utc, to_iso
AUTO_RESTART = ("always", "on-failure", "unless-stopped")
class ContainerService(BaseService):
default_enabled = False
min_interval = 1
title = "Containers"
description = (
"Reconciles desired container state against the docker daemon: launches, stops, restarts, "
"reaps orphans, fires schedules, and samples per-instance metrics. Requires access to the "
"docker socket."
)
config_fields = [
ConfigField("container_metrics_every", "Metrics sample every (ticks)", type="int",
default=1, minimum=1, maximum=60,
help="Sample docker stats once every N reconcile ticks.", group="Containers"),
]
def __init__(self):
super().__init__(name="containers", interval_seconds=5)
self._metric_tick = 0
async def run_once(self) -> None:
backend = get_backend()
try:
ps_rows = await backend.ps(label_filter=api.INSTANCE_LABEL)
except Exception as exc:
self.log(f"docker ps failed: {exc}")
return
by_uid = {}
for row in ps_rows:
uid = row.labels.get(api.INSTANCE_LABEL)
if uid:
by_uid[uid] = row
instances = store.all_instances()
known = {inst["uid"] for inst in instances}
for uid, row in by_uid.items():
if uid not in known:
try:
await backend.rm(row.container_id, force=True)
self.log(f"reaped orphan container {row.name}")
except Exception as exc:
self.log(f"orphan rm failed for {row.name}: {exc}")
for inst in instances:
try:
await self._reconcile(backend, inst, by_uid.get(inst["uid"]))
except Exception as exc:
self.log(f"reconcile {inst.get('name')} failed: {exc}")
await self._fire_schedules()
self._metric_tick += 1
if self._metric_tick % max(1, self.get_config().get("container_metrics_every", 1)) == 0:
await self._sample_metrics(backend, by_uid)
async def _reconcile(self, backend, inst, ps) -> None:
uid = inst["uid"]
desired = inst["desired_state"]
status = inst["status"]
if status == store.ST_REMOVING:
if ps is not None:
await backend.rm(ps.container_id, force=True)
store.delete_instance(uid)
self.log(f"removed instance {inst['name']}")
return
if desired == store.DESIRED_RUNNING:
if ps is None:
await self._launch(backend, inst)
elif ps.state == "running":
if status != store.ST_RUNNING:
store.update_instance(uid, {"status": store.ST_RUNNING, "container_id": ps.container_id,
"started_at": inst.get("started_at") or store.now()})
elif ps.state == "paused":
await backend.unpause(ps.container_id)
store.update_instance(uid, {"status": store.ST_RUNNING})
elif ps.state == "created":
await backend.start(ps.container_id)
store.update_instance(uid, {"status": store.ST_RUNNING})
else:
await self._handle_exit(backend, inst, ps)
elif desired == store.DESIRED_STOPPED:
if ps is not None and ps.state in ("running", "restarting", "paused"):
await backend.stop(ps.container_id)
if status != store.ST_STOPPED:
store.update_instance(uid, {"status": store.ST_STOPPED, "stopped_at": store.now()})
elif desired == store.DESIRED_PAUSED:
if ps is not None and ps.state == "running":
await backend.pause(ps.container_id)
store.update_instance(uid, {"status": store.ST_PAUSED})
async def _launch(self, backend, inst) -> None:
build = store.get_build(inst["build_uid"])
if build is None or build.get("status") != store.BUILD_SUCCESS:
store.update_instance(inst["uid"], {"status": store.ST_CRASHED})
store.record_event(inst, "launch_failed", "reconciler", "", {"reason": "no successful build"})
self.log(f"instance {inst['name']} has no successful build")
return
spec = api.run_spec_for(inst, build["image_tag"])
try:
cid = await backend.run(spec)
except Exception as exc:
store.update_instance(inst["uid"], {"status": store.ST_CRASHED})
store.record_event(inst, "launch_failed", "reconciler", "", {"reason": str(exc)})
self.log(f"launch {inst['name']} failed: {exc}")
return
store.update_instance(inst["uid"], {"container_id": cid, "status": store.ST_RUNNING, "started_at": store.now()})
store.record_event(inst, "start", "reconciler", "")
self.log(f"launched instance {inst['name']}")
async def _exit_logs(self, backend, container_id: str) -> str:
lines: list = []
async def collect(line: str) -> None:
lines.append(line)
try:
await backend.logs(container_id, follow=False, tail=20, on_log=collect)
except Exception:
return ""
return "\n".join(lines)[-2000:]
async def _handle_exit(self, backend, inst, ps) -> None:
uid = inst["uid"]
exit_code = ps.exit_code or 0
logs = await self._exit_logs(backend, ps.container_id)
policy = inst.get("restart_policy", "never")
if policy in AUTO_RESTART and not (policy == "on-failure" and exit_code == 0):
try:
await backend.start(ps.container_id)
store.update_instance(uid, {"status": store.ST_RUNNING,
"restart_count": int(inst.get("restart_count") or 0) + 1})
store.record_event(inst, "policy_restart", "reconciler", "",
{"exit_code": exit_code, "logs": logs})
return
except Exception as exc:
self.log(f"policy restart {inst['name']} failed: {exc}")
terminal = store.ST_CRASHED if exit_code != 0 else store.ST_STOPPED
store.update_instance(uid, {"status": terminal, "desired_state": store.DESIRED_STOPPED,
"exit_code": exit_code, "stopped_at": store.now()})
if terminal == store.ST_CRASHED:
store.record_event(inst, "crash", "reconciler", "", {"exit_code": exit_code, "logs": logs})
async def _fire_schedules(self) -> None:
now = now_utc()
now_iso = to_iso(now)
for sched in store.due_schedules(now_iso):
inst = store.get_instance(sched["instance_uid"])
if inst is None:
store.delete_schedule(sched["uid"])
continue
action = sched["action"]
desired = store.DESIRED_RUNNING if action == "start" else store.DESIRED_STOPPED
store.update_instance(inst["uid"], {"desired_state": desired})
store.record_event(inst, f"schedule_{action}", "scheduler", "")
cols = json.loads(sched.get("schedule_json") or "{}")
run_count = int(sched.get("run_count") or 0) + 1
upcoming = next_run(cols.get("kind"), cols.get("every_seconds"), cols.get("cron"), now)
changes = {"last_run_at": now_iso, "run_count": run_count}
max_runs = cols.get("max_runs")
if upcoming is None or (max_runs and run_count >= max_runs):
changes["enabled"] = 0
changes["next_run_at"] = ""
else:
changes["next_run_at"] = to_iso(upcoming)
store.update_schedule(sched["uid"], changes)
self.log(f"schedule fired: {action} {inst['name']}")
async def _sample_metrics(self, backend, by_uid) -> None:
running = {uid: row for uid, row in by_uid.items() if row.state == "running"}
if not running:
return
try:
samples = await backend.stats_once([row.container_id for row in running.values()])
except Exception as exc:
self.log(f"docker stats failed: {exc}")
return
name_to_uid = {row.name: uid for uid, row in running.items()}
for name, sample in samples.items():
uid = name_to_uid.get(name)
if uid:
store.insert_metric(uid, sample)
def collect_metrics(self) -> dict:
instances = store.all_instances()
counts = {}
for inst in instances:
counts[inst["status"]] = counts.get(inst["status"], 0) + 1
running = counts.get(store.ST_RUNNING, 0)
stats = [
{"label": "Instances", "value": len(instances)},
{"label": "Running", "value": running},
{"label": "Stopped", "value": counts.get(store.ST_STOPPED, 0)},
{"label": "Crashed", "value": counts.get(store.ST_CRASHED, 0)},
{"label": "Paused", "value": counts.get(store.ST_PAUSED, 0)},
]
rows = [[
inst.get("name", "")[:32], inst.get("status", ""), inst.get("desired_state", ""),
inst.get("restart_policy", ""), int(inst.get("restart_count") or 0),
] for inst in instances[:15]]
table = {"columns": ["Instance", "Status", "Desired", "Policy", "Restarts"], "rows": rows}
return {"stats": stats, "table": table}
+277
View File
@@ -0,0 +1,277 @@
# retoor <retoor@molodetz.nl>
import hashlib
import json
from datetime import datetime, timezone
from devplacepy.database import db, get_table
from devplacepy.utils import generate_uid, make_combined_slug
DF_ACTIVE = "active"
DF_ARCHIVED = "archived"
BUILD_QUEUED = "queued"
BUILD_BUILDING = "building"
BUILD_SUCCESS = "success"
BUILD_FAILED = "failed"
BUILD_CANCELLED = "cancelled"
ST_CREATED = "created"
ST_STARTING = "starting"
ST_RUNNING = "running"
ST_STOPPED = "stopped"
ST_PAUSED = "paused"
ST_CRASHED = "crashed"
ST_RESTARTING = "restarting"
ST_REMOVING = "removing"
ST_REMOVED = "removed"
DESIRED_RUNNING = "running"
DESIRED_STOPPED = "stopped"
DESIRED_PAUSED = "paused"
RESTART_POLICIES = ("never", "always", "on-failure", "unless-stopped")
METRICS_RING = 720
def now() -> str:
return datetime.now(timezone.utc).isoformat()
def content_hash(text: str) -> str:
return hashlib.sha256((text or "").encode("utf-8")).hexdigest()
def _exists(table: str) -> bool:
return table in db.tables
# ---------------- dockerfiles ----------------
def create_dockerfile(project_uid, user, name, description, tags) -> dict:
uid = generate_uid()
row = {
"uid": uid, "project_uid": project_uid, "user_uid": user["uid"],
"name": name, "slug": make_combined_slug(name, uid),
"description": description or "", "tags": tags or "", "author": user.get("username", ""),
"current_version": 0, "latest_build_number": 0,
"content_hash": "", "status": DF_ACTIVE,
"created_at": now(), "updated_at": now(),
}
get_table("dockerfiles").insert(row)
return get_dockerfile(uid)
def get_dockerfile(uid_or_slug: str):
table = get_table("dockerfiles")
if not _exists("dockerfiles"):
return None
return table.find_one(uid=uid_or_slug) or table.find_one(slug=uid_or_slug)
def list_dockerfiles(project_uid: str) -> list:
if not _exists("dockerfiles"):
return []
return sorted(get_table("dockerfiles").find(project_uid=project_uid),
key=lambda r: r.get("created_at", ""), reverse=True)
def update_dockerfile(uid: str, changes: dict) -> None:
get_table("dockerfiles").update({"uid": uid, "updated_at": now(), **changes}, ["uid"])
# ---------------- versions (immutable) ----------------
def create_version(dockerfile: dict, content: str, user: dict) -> dict:
version = int(dockerfile.get("current_version") or 0) + 1
uid = generate_uid()
get_table("dockerfile_versions").insert({
"uid": uid, "dockerfile_uid": dockerfile["uid"], "project_uid": dockerfile["project_uid"],
"version": version, "content": content, "content_hash": content_hash(content),
"created_at": now(), "created_by": user["uid"],
})
update_dockerfile(dockerfile["uid"], {"current_version": version, "content_hash": content_hash(content)})
return get_table("dockerfile_versions").find_one(uid=uid)
def get_version(uid: str):
if not _exists("dockerfile_versions"):
return None
return get_table("dockerfile_versions").find_one(uid=uid)
def current_version(dockerfile: dict):
if not _exists("dockerfile_versions"):
return None
return get_table("dockerfile_versions").find_one(
dockerfile_uid=dockerfile["uid"], version=int(dockerfile.get("current_version") or 0))
def list_versions(dockerfile_uid: str) -> list:
if not _exists("dockerfile_versions"):
return []
return sorted(get_table("dockerfile_versions").find(dockerfile_uid=dockerfile_uid),
key=lambda r: r.get("version", 0), reverse=True)
# ---------------- builds ----------------
def create_build(dockerfile: dict, version: dict, build_number: int, image_tag: str, also_latest: bool) -> dict:
uid = generate_uid()
get_table("builds").insert({
"uid": uid, "dockerfile_uid": dockerfile["uid"], "dockerfile_version_uid": version["uid"],
"project_uid": dockerfile["project_uid"], "build_number": build_number,
"image_tag": image_tag, "also_latest": 1 if also_latest else 0,
"status": BUILD_QUEUED, "job_uid": "", "image_id": "", "logs": "", "error": "",
"duration_ms": 0, "queued_at": now(), "started_at": "", "completed_at": "", "created_at": now(),
})
return get_build(uid)
def get_build(uid: str):
if not _exists("builds"):
return None
return get_table("builds").find_one(uid=uid)
def update_build(uid: str, changes: dict) -> None:
get_table("builds").update({"uid": uid, **changes}, ["uid"])
def list_builds(dockerfile_uid: str, limit: int = 50) -> list:
if not _exists("builds"):
return []
rows = sorted(get_table("builds").find(dockerfile_uid=dockerfile_uid),
key=lambda r: r.get("build_number", 0), reverse=True)
return rows[:limit]
# ---------------- instances ----------------
def create_instance(row: dict) -> dict:
uid = generate_uid()
base = {
"uid": uid, "slug": make_combined_slug(row.get("name", "instance"), uid),
"container_id": "", "status": ST_CREATED, "exit_code": 0, "restart_count": 0,
"ingress_slug": "", "ingress_port": 0,
"started_at": "", "stopped_at": "", "created_at": now(), "updated_at": now(),
}
base.update(row)
base["uid"] = uid
get_table("instances").insert(base)
return get_instance(uid)
def get_instance(uid: str):
table = get_table("instances")
if not _exists("instances"):
return None
return table.find_one(uid=uid) or table.find_one(slug=uid)
def list_instances(project_uid: str = None, dockerfile_uid: str = None) -> list:
if not _exists("instances"):
return []
filters = {}
if project_uid:
filters["project_uid"] = project_uid
if dockerfile_uid:
filters["dockerfile_uid"] = dockerfile_uid
return sorted(get_table("instances").find(**filters),
key=lambda r: r.get("created_at", ""), reverse=True)
def all_instances() -> list:
if not _exists("instances"):
return []
return list(get_table("instances").find())
def find_instance_by_ingress(slug: str):
if not slug or not _exists("instances"):
return None
matches = list(get_table("instances").find(ingress_slug=slug))
for instance in matches:
if instance.get("status") == ST_RUNNING:
return instance
return matches[0] if matches else None
def update_instance(uid: str, changes: dict) -> None:
get_table("instances").update({"uid": uid, "updated_at": now(), **changes}, ["uid"])
def delete_instance(uid: str) -> None:
get_table("instances").delete(uid=uid)
# ---------------- events / metrics / schedules ----------------
def record_event(instance: dict, event: str, actor_kind: str, actor_id: str, detail: dict = None) -> None:
get_table("instance_events").insert({
"uid": generate_uid(), "instance_uid": instance["uid"], "project_uid": instance.get("project_uid", ""),
"event": event, "actor_kind": actor_kind, "actor_id": actor_id or "",
"detail": json.dumps(detail or {}), "created_at": now(),
})
def list_events(instance_uid: str, limit: int = 100) -> list:
if not _exists("instance_events"):
return []
rows = sorted(get_table("instance_events").find(instance_uid=instance_uid),
key=lambda r: r.get("created_at", ""), reverse=True)
return rows[:limit]
def insert_metric(instance_uid: str, sample) -> None:
table = get_table("instance_metrics")
table.insert({
"uid": generate_uid(), "instance_uid": instance_uid, "ts": now(),
"cpu_pct": sample.cpu_pct, "mem_bytes": sample.mem_bytes, "mem_pct": sample.mem_pct,
"net_rx": sample.net_rx, "net_tx": sample.net_tx,
"blk_read": sample.blk_read, "blk_write": sample.blk_write,
})
rows = sorted(table.find(instance_uid=instance_uid), key=lambda r: r.get("ts", ""))
if len(rows) > METRICS_RING:
for old in rows[:len(rows) - METRICS_RING]:
table.delete(uid=old["uid"])
def recent_metrics(instance_uid: str, limit: int = 120) -> list:
if not _exists("instance_metrics"):
return []
rows = sorted(get_table("instance_metrics").find(instance_uid=instance_uid),
key=lambda r: r.get("ts", ""))
return rows[-limit:]
def create_schedule(instance: dict, action: str, schedule_columns: dict, next_run_at: str) -> dict:
uid = generate_uid()
get_table("instance_schedules").insert({
"uid": uid, "instance_uid": instance["uid"], "project_uid": instance.get("project_uid", ""),
"action": action, "schedule_json": json.dumps(schedule_columns), "enabled": 1,
"next_run_at": next_run_at, "last_run_at": "", "run_count": 0,
"created_at": now(), "updated_at": now(),
})
return get_table("instance_schedules").find_one(uid=uid)
def list_schedules(instance_uid: str) -> list:
if not _exists("instance_schedules"):
return []
return list(get_table("instance_schedules").find(instance_uid=instance_uid))
def due_schedules(now_iso: str) -> list:
if not _exists("instance_schedules"):
return []
return [r for r in get_table("instance_schedules").find(enabled=1)
if r.get("next_run_at") and r["next_run_at"] <= now_iso]
def update_schedule(uid: str, changes: dict) -> None:
get_table("instance_schedules").update({"uid": uid, "updated_at": now(), **changes}, ["uid"])
def delete_schedule(uid: str) -> None:
get_table("instance_schedules").delete(uid=uid)
@@ -0,0 +1,50 @@
# retoor <retoor@molodetz.nl>
DEFAULT_DOCKERFILE = """# syntax=docker/dockerfile:1
FROM python:3.13-slim-bookworm
ENV PYTHONUNBUFFERED=1 \\
PYTHONDONTWRITEBYTECODE=1
RUN useradd --create-home --shell /bin/bash app
WORKDIR /app
COPY . /app
RUN chown -R app:app /app
USER app
EXPOSE 8000
CMD ["python3", "-m", "http.server", "8000", "--directory", "/app"]
"""
PLAYWRIGHT_DOCKERFILE = """# syntax=docker/dockerfile:1
FROM python:3.13-slim-bookworm AS base
ENV PYTHONUNBUFFERED=1 \\
PYTHONDONTWRITEBYTECODE=1 \\
PIP_NO_CACHE_DIR=1 \\
PIP_DISABLE_PIP_VERSION_CHECK=1 \\
DEBIAN_FRONTEND=noninteractive \\
PLAYWRIGHT_BROWSERS_PATH=/opt/playwright
RUN apt-get update && apt-get install -y --no-install-recommends \\
git curl ca-certificates build-essential libpq-dev \\
&& rm -rf /var/lib/apt/lists/*
FROM base AS deps
RUN pip install \\
numpy pandas requests \\
fastapi "uvicorn[standard]" pydantic \\
sqlalchemy psycopg2-binary redis \\
httpx aiofiles python-dotenv \\
playwright \\
&& playwright install --with-deps chromium
FROM deps AS runtime
RUN useradd --create-home --shell /bin/bash app
WORKDIR /app
COPY . /app
RUN chown -R app:app /app /opt/playwright
USER app
CMD ["sleep", "infinity"]
"""
+117 -2
View File
@@ -205,6 +205,38 @@ ACTIONS: tuple[Action, ...] = (
summary="Delete a project",
params=(path("project_slug", "Exact project slug copied from a /projects/... link in a listing response; do not build it from the title."),),
),
Action(
name="project_set_private",
method="POST",
path="/projects/{project_slug}/private",
summary="Mark a project private (owner-only) or public",
description=(
"Set value=true to hide the project from everyone except its owner (and administrators), "
"or value=false to make it public again. Only the project owner may change this."
),
params=(
path("project_slug", "Project slug or uid."),
body("value", "true to make the project private, false to make it public.", required=True),
),
),
Action(
name="project_set_readonly",
method="POST",
path="/projects/{project_slug}/readonly",
summary="Mark a project read-only (immutable files) or writable again",
description=(
"Set value=true to make every file in the project immutable - no writes, edits, line "
"edits, moves, deletes, or uploads succeed afterwards, from anyone including you - or "
"value=false to allow changes again. This is a significant change: you MUST ask the user "
"for explicit confirmation BEFORE calling it, and only pass confirm=true once they have "
"agreed. Only the project owner may change this."
),
params=(
path("project_slug", "Project slug or uid."),
body("value", "true to make the project read-only, false to make it writable.", required=True),
body("confirm", "Must be true, set only after the user has explicitly confirmed.", required=True),
),
),
Action(
name="project_list_files",
method="GET",
@@ -230,14 +262,97 @@ ACTIONS: tuple[Action, ...] = (
name="project_write_file",
method="POST",
path="/projects/{project_slug}/files/write",
summary="Create or overwrite a text file in a project (parent directories are created automatically)",
description="The primary tool for building a project: write any text file by path. Missing parent directories are created recursively.",
summary="Create or overwrite a whole text file in a project (parent directories are created automatically)",
description=(
"Replaces the ENTIRE file with the content you send. Creating a new file needs no prior "
"read, but overwriting an existing one requires reading it first (project_read_file). "
"For an existing file prefer the surgical line tools (project_replace_lines, "
"project_insert_lines, project_delete_lines, project_append_file) instead of rewriting it, "
"and never batch several writes in one turn. Use this only to create a new file or fully "
"rewrite a small one. Missing parent directories are created recursively."
),
params=(
path("project_slug", "Project slug or uid."),
body("path", "Relative file path, e.g. src/app/main.py.", required=True),
body("content", "Full file content.", required=True),
),
),
Action(
name="project_read_lines",
method="GET",
path="/projects/{project_slug}/files/lines",
summary="Read a 1-indexed line range of a text file in a project",
description=(
"Returns {path, start, end, total_lines, lines, content} for the requested range. Use it "
"to inspect part of a large file before editing, and to learn total_lines so you can target "
"the right range with the line-edit tools."
),
params=(
path("project_slug", "Project slug or uid."),
query("path", "Relative file path inside the project.", required=True),
query("start", "First line to read (1-indexed, default 1)."),
query("end", "Last line to read (inclusive); omit for end of file."),
),
requires_auth=False,
),
Action(
name="project_replace_lines",
method="POST",
path="/projects/{project_slug}/files/replace-lines",
summary="Replace an inclusive 1-indexed line range of a text file with new content",
description=(
"Surgically rewrites lines start..end (inclusive) with content (which may be any number of "
"lines, or empty to delete the range). The preferred way to edit an existing file: it leaves "
"the rest of the file untouched and never overflows the model output limit."
),
params=(
path("project_slug", "Project slug or uid."),
body("path", "Relative file path.", required=True),
body("start", "First line to replace (1-indexed).", required=True),
body("end", "Last line to replace (inclusive).", required=True),
body("content", "Replacement text for those lines (empty deletes them)."),
),
),
Action(
name="project_insert_lines",
method="POST",
path="/projects/{project_slug}/files/insert-lines",
summary="Insert content before a 1-indexed line of a text file",
description=(
"Inserts content before line 'at' without touching existing lines. Use at=1 to prepend and "
"at=total_lines+1 to insert at the end."
),
params=(
path("project_slug", "Project slug or uid."),
body("path", "Relative file path.", required=True),
body("at", "Insert before this 1-indexed line (1 prepends, total+1 appends).", required=True),
body("content", "Text to insert.", required=True),
),
),
Action(
name="project_delete_lines",
method="POST",
path="/projects/{project_slug}/files/delete-lines",
summary="Delete an inclusive 1-indexed line range from a text file",
params=(
path("project_slug", "Project slug or uid."),
body("path", "Relative file path.", required=True),
body("start", "First line to delete (1-indexed).", required=True),
body("end", "Last line to delete (inclusive).", required=True),
),
),
Action(
name="project_append_file",
method="POST",
path="/projects/{project_slug}/files/append",
summary="Append content to the end of a text file in a project",
description="Adds content as new lines at the end of the file. Use it to grow a large file across turns without resending the whole thing.",
params=(
path("project_slug", "Project slug or uid."),
body("path", "Relative file path.", required=True),
body("content", "Text to append.", required=True),
),
),
Action(
name="project_upload_file",
method="POST",
@@ -17,6 +17,7 @@ CHUNK_ACTIONS: tuple[Action, ...] = (
),
handler="chunks",
requires_auth=False,
read_only=True,
params=(
Param(
name="chunk_id",
@@ -23,6 +23,7 @@ CLIENT_ACTIONS: tuple[Action, ...] = (
description=CLIENT + " Use this to understand where the user is and what they are looking at before acting or guiding them.",
handler="client",
requires_auth=False,
read_only=True,
),
Action(
name="run_js",
@@ -0,0 +1,130 @@
# retoor <retoor@molodetz.nl>
from .spec import Action, Param
def arg(name: str, description: str, required: bool = False, kind: str = "string") -> Param:
return Param(name=name, location="body", description=description, required=required, type=kind)
SLUG = arg("project_slug", "Project slug or uid that owns the container resources.", required=True)
CONTAINER_ACTIONS: tuple[Action, ...] = (
Action(
name="container_list_dockerfiles",
method="LOCAL", path="", handler="container", requires_admin=True, read_only=True,
summary="List the Dockerfiles, builds, and instances of a project",
params=(SLUG,),
),
Action(
name="container_create_dockerfile",
method="LOCAL", path="", handler="container", requires_admin=True,
summary="Create a Dockerfile in a project and queue its first build",
description="With no content a lean python-slim default builds in seconds and serves /app on port 8000 (long-lived, ingress-ready). Pass full content for anything else; for browser automation start from python-slim and add 'pip install playwright && playwright install --with-deps chromium'.",
params=(
SLUG,
arg("name", "Image name, lowercase letters/digits/.-_ (max 63 chars).", required=True),
arg("description", "Optional description."),
arg("tags", "Optional comma separated tags."),
arg("content", "Optional full Dockerfile content."),
),
),
Action(
name="container_update_dockerfile",
method="LOCAL", path="", handler="container", requires_admin=True,
summary="Replace a Dockerfile's content, creating a new immutable version and queueing a build",
params=(
SLUG,
arg("dockerfile", "Dockerfile name, slug, or uid.", required=True),
arg("content", "New Dockerfile content.", required=True),
),
),
Action(
name="container_build",
method="LOCAL", path="", handler="container", requires_admin=True,
summary="Rebuild the current version of a Dockerfile",
params=(SLUG, arg("dockerfile", "Dockerfile name, slug, or uid.", required=True)),
),
Action(
name="container_build_status",
method="LOCAL", path="", handler="container", requires_admin=True, read_only=True,
summary="Get a build's status and recent log tail",
params=(SLUG, arg("build_uid", "Build uid.", required=True)),
),
Action(
name="container_list_instances",
method="LOCAL", path="", handler="container", requires_admin=True, read_only=True,
summary="List a project's container instances and their status",
params=(SLUG,),
),
Action(
name="container_create_instance",
method="LOCAL", path="", handler="container", requires_admin=True,
summary="Create and start a container instance from a Dockerfile's latest successful build",
params=(
SLUG,
arg("dockerfile", "Dockerfile name, slug, or uid.", required=True),
arg("name", "Instance name.", required=True),
arg("boot_command", "Optional command to run on boot, e.g. 'python app.py'."),
arg("restart_policy", "never, always, on-failure, or unless-stopped."),
arg("env", "Optional env vars as KEY=VALUE lines."),
arg("ports", "Port maps per line or comma separated. Use a bare container port (e.g. '8899') to auto-assign a unique host port above 20000, or 'host:container' to pin one."),
arg("cpu_limit", "Optional CPU limit, e.g. 1 or 1.5."),
arg("mem_limit", "Optional memory limit, e.g. 512m or 1g."),
arg("autostart", "Start immediately ('true' or 'false', default true)."),
arg("ingress_slug", "Optional public ingress slug; the service is then reachable at /p/<slug>."),
arg("ingress_port", "Container port to publish at /p/<slug> (must be one of the mapped ports).", kind="integer"),
),
),
Action(
name="container_instance_action",
method="LOCAL", path="", handler="container", requires_admin=True,
summary="Control an instance: start, stop, restart, pause, resume, delete, or sync",
description="sync imports the container /app workspace back into the project files.",
params=(
SLUG,
arg("instance", "Instance name, slug, or uid.", required=True),
arg("action", "start, stop, restart, pause, resume, delete, or sync.", required=True),
),
),
Action(
name="container_logs",
method="LOCAL", path="", handler="container", requires_admin=True, read_only=True,
summary="Read the recent logs of a running instance",
params=(
SLUG,
arg("instance", "Instance name, slug, or uid.", required=True),
arg("tail", "Number of log lines (default 200).", kind="integer"),
),
),
Action(
name="container_exec",
method="LOCAL", path="", handler="container", requires_admin=True,
summary="Run a one-shot command inside a running instance and return its output",
params=(
SLUG,
arg("instance", "Instance name, slug, or uid.", required=True),
arg("command", "Command to run, e.g. 'pip list'.", required=True),
),
),
Action(
name="container_stats",
method="LOCAL", path="", handler="container", requires_admin=True, read_only=True,
summary="Get aggregated resource and runtime statistics for an instance",
params=(SLUG, arg("instance", "Instance name, slug, or uid.", required=True)),
),
Action(
name="container_schedule",
method="LOCAL", path="", handler="container", requires_admin=True,
summary="Schedule a start or stop of an instance (cron, interval, or one-time)",
params=(
SLUG,
arg("instance", "Instance name, slug, or uid.", required=True),
arg("action", "start or stop.", required=True),
arg("kind", "once, interval, or cron.", required=True),
arg("cron", "Cron expression for kind=cron, e.g. '0 2 * * *'."),
arg("run_at", "ISO time for kind=once, e.g. 2026-06-15T02:00:00."),
arg("every_seconds", "Interval seconds for kind=interval.", kind="integer"),
),
),
)
@@ -18,6 +18,7 @@ COST_ACTIONS: tuple[Action, ...] = (
),
handler="cost",
requires_auth=False,
read_only=True,
),
Action(
name="cost_stats",
@@ -33,5 +34,6 @@ COST_ACTIONS: tuple[Action, ...] = (
handler="cost",
requires_auth=True,
requires_admin=True,
read_only=True,
),
)
@@ -30,9 +30,25 @@ from .spec import Action, Catalog
MUTATING_METHODS = ("POST", "DELETE", "PUT", "PATCH")
CONFIRM_REQUIRED = {"project_set_readonly"}
logger = logging.getLogger("devii.dispatch")
def _is_confirmed(arguments: dict[str, Any]) -> bool:
return str(arguments.get("confirm", "")).strip().lower() in ("true", "1", "yes", "on")
def confirmation_error(name: str, arguments: dict[str, Any]) -> ToolInputError | None:
if name in CONFIRM_REQUIRED and not _is_confirmed(arguments):
return ToolInputError(
"Setting a project read-only makes every file immutable and blocks all further "
"edits. Ask the user to confirm this explicitly first, then call again with "
"confirm=true."
)
return None
class Dispatcher:
def __init__(
self,
@@ -59,6 +75,25 @@ class Dispatcher:
self._cost = CostController(quota_provider=quota_provider)
self._chunks = ChunkController(settings)
self._rsearch = RsearchController(settings)
from ..container import ContainerController
self._container = ContainerController(client)
self._read_files: set[tuple[str, str]] = set()
@staticmethod
def _file_key(arguments: dict[str, Any]) -> tuple[str, str] | None:
from devplacepy.project_files import normalize_path, ProjectFileError as _PFError
raw = arguments.get("path")
if not raw:
return None
try:
path = normalize_path(raw)
except _PFError:
return None
return str(arguments.get("project_slug", "")), path
def is_read_only(self, name: str) -> bool:
action = self._actions.get(name)
return bool(action and action.is_read_only)
async def dispatch(self, name: str, arguments: dict[str, Any]) -> str:
action = self._actions.get(name)
@@ -77,6 +112,9 @@ class Dispatcher:
"This information is restricted to administrators.",
tool=name,
)
guard = confirmation_error(name, arguments)
if guard is not None:
raise guard
resource_key = self._resource_key(action, arguments)
if resource_key:
cached = serve_resource(resource_key, self._settings.max_response_chars)
@@ -136,6 +174,9 @@ class Dispatcher:
if action.handler == "rsearch":
return await self._rsearch.dispatch(action.name, arguments)
if action.handler == "container":
return await self._container.dispatch(action.name, arguments)
if action.handler == "avatar":
if self._avatar is None:
return error_result(
@@ -204,7 +245,30 @@ class Dispatcher:
return None
return None
async def _file_exists(self, arguments: dict[str, Any]) -> bool:
slug = str(arguments.get("project_slug", "")).strip()
path = str(arguments.get("path", "")).strip()
if not slug or not path:
return False
response = await self._client.call(
method="GET",
path=f"/projects/{quote(slug, safe='')}/files/raw",
params={"path": path},
headers={"X-Requested-With": "fetch"},
)
return response.status_code == 200
async def _run_http(self, action: Action, arguments: dict[str, Any]) -> str:
if action.name == "project_write_file":
key = self._file_key(arguments)
if key is not None and key not in self._read_files and await self._file_exists(arguments):
raise ToolInputError(
f"Read '{key[1]}' before overwriting it. It already exists; call "
"project_read_file first. For an existing file prefer the line tools "
"(project_replace_lines, project_insert_lines, project_delete_lines, "
"project_append_file); project_write_file replaces the entire file."
)
url_path, params, data, file_field = self._build_request(action, arguments)
headers = {"X-Requested-With": "fetch"} if action.ajax else None
@@ -216,6 +280,10 @@ class Dispatcher:
file_field=file_field,
headers=headers,
)
if action.name in ("project_read_file", "project_read_lines", "project_write_file"):
key = self._file_key(arguments)
if key is not None:
self._read_files.add(key)
if action.method in MUTATING_METHODS:
record_mutation(action.name)
store = get_store()
@@ -17,6 +17,7 @@ DOCS_ACTIONS: tuple[Action, ...] = (
),
handler="docs",
requires_auth=False,
read_only=True,
params=(
Param(
name="query",
@@ -18,6 +18,7 @@ FETCH_ACTIONS: tuple[Action, ...] = (
),
handler="fetch",
requires_auth=False,
read_only=True,
params=(
Param(
name="url",
@@ -23,6 +23,7 @@ RSEARCH_ACTIONS: tuple[Action, ...] = (
),
handler="rsearch",
requires_auth=False,
read_only=True,
params=(
Param(name="query", location="body", description="The web search query.", required=True),
Param(name="count", location="body", description="Number of results (1-100, default 10).", type="integer"),
@@ -43,6 +44,7 @@ RSEARCH_ACTIONS: tuple[Action, ...] = (
),
handler="rsearch",
requires_auth=False,
read_only=True,
params=(
Param(name="query", location="body", description="The question or prompt to answer.", required=True),
Param(name="content", location="body", description="Let the AI read full page content while answering.", type="boolean"),
@@ -61,6 +63,7 @@ RSEARCH_ACTIONS: tuple[Action, ...] = (
),
handler="rsearch",
requires_auth=False,
read_only=True,
params=(
Param(name="prompt", location="body", description="The prompt to send.", required=True),
Param(name="json", location="body", description="Force a valid-JSON-only response.", type="boolean"),
@@ -78,6 +81,7 @@ RSEARCH_ACTIONS: tuple[Action, ...] = (
),
handler="rsearch",
requires_auth=False,
read_only=True,
params=(
Param(name="url", location="body", description="Public URL of the image to describe.", required=True),
),
+6 -1
View File
@@ -29,10 +29,15 @@ class Action:
requires_admin: bool = False
handler: Literal[
"http", "login", "logout", "status", "task", "agentic", "avatar", "client", "fetch",
"docs", "cost", "chunks", "rsearch"
"docs", "cost", "chunks", "rsearch", "container"
] = "http"
freeform_body: bool = False
ajax: bool = False
read_only: bool = False
@property
def is_read_only(self) -> bool:
return self.read_only or self.method.upper() == "GET"
def tool_schema(self) -> dict[str, Any]:
properties: dict[str, Any] = {}
+38
View File
@@ -70,6 +70,44 @@ SYSTEM_PROMPT = (
"trial-and-error, and never tell the user that a page or capability does not exist without "
"confirming against the docs. If a guess returns a 404, that means you guessed - search the docs "
"instead of concluding the feature is missing.\n\n"
"WRITING AND EDITING FILES\n"
"Creating a new file needs no prior read, but overwriting an existing one does: call "
"project_read_file first. project_write_file replaces the ENTIRE file and is rejected when the "
"path already exists and you have not read it this session, so use it only to create a new file "
"or fully rewrite a small one. To change an existing "
"file, prefer the surgical line tools - project_read_lines to inspect a range and learn total_lines, "
"then project_replace_lines, project_insert_lines, project_delete_lines, or project_append_file to "
"edit just the affected lines. These leave the rest of the file untouched and let you build or modify "
"very large files across turns without resending the whole thing. A single model completion has a "
"maximum output length: emit only ONE file write per turn and never batch several writes into one "
"response, or their combined content overflows the completion and the last file is truncated. If a "
"tool result has error 'tool_input_truncated', your output was cut off - resend that single write or "
"line edit on its own.\n\n"
"PROJECT PRIVACY AND READ-ONLY\n"
"A project owner can mark a project private (project_set_private) so only the owner and "
"administrators can see it, and read-only (project_set_readonly) so every file becomes immutable. "
"When a project is read-only, all file writes fail with 'project is read-only'; to edit such a "
"project you must first ask the owner for permission and set it writable again. Marking a project "
"read-only is significant and largely locks it down, so you MUST obtain explicit user confirmation "
"before calling project_set_readonly with value=true, and only then pass confirm=true; never set a "
"project read-only on your own initiative.\n\n"
"CONTAINERS (admin only)\n"
"When you are an administrator you can manage a project's Docker containers with the container_* "
"tools: author Dockerfiles (a high-quality default is provided), build versioned images "
"asynchronously, create and control instances "
"(container_instance_action: start/stop/restart/pause/resume/delete/sync), read logs, run one-shot "
"commands (container_exec), inspect stats, and schedule starts/stops. Builds and instances run on "
"the host docker daemon, so confirm destructive actions and never expose backend or infrastructure "
"detail. These tools are unavailable to non-administrators. "
"Builds are asynchronous and can take a while for heavy base images. "
"container_build_status reports a 'status' (queued, building, success, failed, cancelled) "
"and a 'terminal' flag: a status of queued or building is NORMAL in-progress, not an error - keep "
"polling patiently (wait between polls) and do NOT change the Dockerfile or start a new build until "
"'terminal' is true. Only treat status 'failed' (or a 'build_error' field) as a failure to fix. "
"A running instance can be published with an ingress slug; the container tools then return an "
"'ingress_url'. To reach a published service, use that 'ingress_url' (the configured public URL, "
"e.g. https://host/p/<slug>) - never localhost or 127.0.0.1, which web tools refuse. If ingress_url "
"is a relative /p/<slug>, the admin has not set the public Site URL in admin settings yet.\n\n"
"REMOTE WEB TOOLS (rsearch)\n"
"The rsearch_* tools (rsearch, rsearch_answer, rsearch_chat, rsearch_describe_image) reach an "
"EXTERNAL public web/AI service, not this platform. They are not platform-specific, so platform "
+25 -5
View File
@@ -37,6 +37,7 @@ ITERATION_LIMIT_MESSAGE = "[stopped] Maximum iterations reached without a final
LABEL_KEYS = (
"path",
"post_slug",
"project_slug",
"gist_slug",
@@ -75,6 +76,12 @@ def call_label(call: dict[str, Any]) -> str:
target_type = arguments.get("target_type")
target = clean(arguments["target_uid"])
return f"{target_type}/{target}" if target_type else target
if arguments.get("command"):
return clean(arguments["command"])
if arguments.get("instance"):
instance = clean(arguments["instance"])
action = arguments.get("action")
return f"{clean(action)} {instance}" if action else instance
for key in LABEL_KEYS:
value = arguments.get(key)
if value not in (None, ""):
@@ -122,10 +129,18 @@ async def _run_tool_call(dispatcher: Any, call: dict[str, Any]) -> str:
raw_arguments = function.get("arguments") or "{}"
try:
arguments = json.loads(raw_arguments) if isinstance(raw_arguments, str) else raw_arguments
if not isinstance(arguments, dict):
raise ValueError("arguments must be an object")
except ValueError as exc:
return json.dumps({"error": "tool_input_error", "message": f"Invalid arguments: {exc}"})
except json.JSONDecodeError as exc:
return json.dumps({
"error": "tool_input_truncated",
"message": (
f"The arguments for {name or 'this tool'} were cut off and could not be parsed "
f"({exc.msg} at position {exc.pos}); the model output hit its length limit. "
"Emit only one write tool call per turn (do not batch several file writes into a "
"single response) and resend this one call on its own."
),
})
if not isinstance(arguments, dict):
return json.dumps({"error": "tool_input_error", "message": "Invalid arguments: arguments must be an object"})
return await dispatcher.dispatch(name, arguments)
@@ -169,7 +184,12 @@ async def react_loop(
tool_calls = message.get("tool_calls") or []
if tool_calls:
if plan_required and state.plan is None and tool_calls[0]["function"]["name"] != "plan":
needs_plan = plan_required and state.plan is None and any(
call["function"]["name"] != "plan"
and not dispatcher.is_read_only(call["function"]["name"])
for call in tool_calls
)
if needs_plan:
trace("plan-gate")
for call in tool_calls:
messages.append(
@@ -0,0 +1,5 @@
# retoor <retoor@molodetz.nl>
from devplacepy.services.devii.container.controller import ContainerController
__all__ = ["ContainerController"]
@@ -0,0 +1,217 @@
# retoor <retoor@molodetz.nl>
import json
import logging
from typing import Any
from devplacepy.database import get_table, resolve_by_slug
from devplacepy.services.containers import api, store
from devplacepy.services.containers.api import ContainerError
from devplacepy.services.containers.runtime import get_backend
from devplacepy.services.devii.errors import ToolInputError
from devplacepy.services.devii.tasks.schedule import Schedule, from_iso
logger = logging.getLogger("devii.container")
class ContainerController:
def __init__(self, client: Any = None) -> None:
self._client = client
def _ingress_url(self, instance: dict):
slug = instance.get("ingress_slug")
if not slug:
return None
from devplacepy.seo import public_base_url
base = public_base_url()
return f"{base}/p/{slug}" if base else f"/p/{slug}"
def _actor_user(self) -> dict:
username = getattr(self._client, "username", None)
if username:
user = get_table("users").find_one(username=username)
if user:
return user
return {"uid": "admin", "username": username or "admin"}
def _project(self, arguments: dict) -> dict:
slug = str(arguments.get("project_slug", "")).strip()
project = resolve_by_slug(get_table("projects"), slug) if slug else None
if not project:
raise ToolInputError(f"project not found: {slug}")
return project
def _dockerfile(self, project: dict, ref: str) -> dict:
df = store.get_dockerfile(ref)
if df is None or df["project_uid"] != project["uid"]:
for candidate in store.list_dockerfiles(project["uid"]):
if candidate["name"] == ref:
return candidate
raise ToolInputError(f"dockerfile not found: {ref}")
return df
def _instance(self, project: dict, ref: str) -> dict:
inst = store.get_instance(ref)
if inst is None or inst["project_uid"] != project["uid"]:
for candidate in store.list_instances(project["uid"]):
if candidate["name"] == ref:
return candidate
raise ToolInputError(f"instance not found: {ref}")
return inst
def _build_view(self, build: dict) -> dict:
if not build:
return {}
terminal = build["status"] in (store.BUILD_SUCCESS, store.BUILD_FAILED, store.BUILD_CANCELLED)
view = {
"uid": build["uid"], "status": build["status"], "terminal": terminal,
"build_number": build.get("build_number"), "image_tag": build.get("image_tag"),
"image_id": build.get("image_id") or None, "duration_ms": build.get("duration_ms"),
}
if build.get("error"):
view["build_error"] = build["error"]
return view
async def dispatch(self, name: str, arguments: dict[str, Any]) -> str:
handler = getattr(self, "_" + name[len("container_"):], None) if name.startswith("container_") else None
if handler is None:
raise ToolInputError(f"unknown container tool: {name}")
try:
return await handler(arguments)
except ContainerError as exc:
return json.dumps({"error": "container_error", "message": str(exc)})
async def _list_dockerfiles(self, arguments) -> str:
project = self._project(arguments)
return json.dumps({
"dockerfiles": store.list_dockerfiles(project["uid"]),
"instances": store.list_instances(project["uid"]),
}, ensure_ascii=False, default=str)
async def _create_dockerfile(self, arguments) -> str:
project = self._project(arguments)
result = await api.create_dockerfile(
project, self._actor_user(), name=str(arguments.get("name", "")),
description=str(arguments.get("description", "")), tags=str(arguments.get("tags", "")),
content=arguments.get("content") or None)
return json.dumps({"dockerfile": result["dockerfile"], "build": self._build_view(result["build"])}, default=str)
async def _update_dockerfile(self, arguments) -> str:
project = self._project(arguments)
df = self._dockerfile(project, str(arguments.get("dockerfile", "")))
result = await api.save_version(df, self._actor_user(), str(arguments.get("content", "")))
return json.dumps({"changed": result["changed"], "build": self._build_view(result["build"])}, default=str)
async def _build(self, arguments) -> str:
project = self._project(arguments)
df = self._dockerfile(project, str(arguments.get("dockerfile", "")))
version = store.current_version(df)
if version is None:
raise ToolInputError("dockerfile has no version to build")
build = await api.enqueue_build(df, version, owner=("user", self._actor_user()["uid"]))
return json.dumps({"build": self._build_view(build)}, default=str)
async def _build_status(self, arguments) -> str:
build = store.get_build(str(arguments.get("build_uid", "")))
if build is None:
raise ToolInputError("build not found")
view = self._build_view(build)
view["logs_tail"] = (build.get("logs") or "")[-4000:]
return json.dumps(view, default=str)
async def _list_instances(self, arguments) -> str:
project = self._project(arguments)
instances = store.list_instances(project["uid"])
for inst in instances:
inst["ingress_url"] = self._ingress_url(inst)
return json.dumps({"instances": instances}, default=str)
async def _create_instance(self, arguments) -> str:
project = self._project(arguments)
df = self._dockerfile(project, str(arguments.get("dockerfile", "")))
successes = [b for b in store.list_builds(df["uid"]) if b["status"] == store.BUILD_SUCCESS]
if not successes:
raise ToolInputError("dockerfile has no successful build yet")
autostart = str(arguments.get("autostart", "true")).lower() not in ("false", "0", "no", "off")
inst = await api.create_instance(
project, df, successes[0], name=str(arguments.get("name", "")),
boot_command=str(arguments.get("boot_command", "")), env=arguments.get("env", ""),
cpu_limit=str(arguments.get("cpu_limit", "")), mem_limit=str(arguments.get("mem_limit", "")),
ports=arguments.get("ports", ""), restart_policy=str(arguments.get("restart_policy", "never")),
autostart=autostart, ingress_slug=str(arguments.get("ingress_slug", "")),
ingress_port=arguments.get("ingress_port"), actor=("user", self._actor_user()["uid"]))
return json.dumps({"instance": inst, "ingress_url": self._ingress_url(inst)}, default=str)
async def _instance_action(self, arguments) -> str:
project = self._project(arguments)
inst = self._instance(project, str(arguments.get("instance", "")))
action = str(arguments.get("action", "")).lower()
actor = ("user", self._actor_user()["uid"])
if action == "delete":
api.mark_for_removal(inst, actor=actor)
elif action == "sync":
count = await api.sync_workspace(inst, self._actor_user())
return json.dumps({"status": "synced", "imported": count})
elif action == "restart":
api.request_restart(inst, actor=actor)
elif action in ("start", "resume"):
api.set_desired_state(inst, store.DESIRED_RUNNING, actor=actor)
elif action == "stop":
api.set_desired_state(inst, store.DESIRED_STOPPED, actor=actor)
elif action == "pause":
api.set_desired_state(inst, store.DESIRED_PAUSED, actor=actor)
else:
raise ToolInputError(f"unknown action: {action}")
return json.dumps({"status": "ok", "action": action, "instance": inst["uid"]})
async def _logs(self, arguments) -> str:
project = self._project(arguments)
inst = self._instance(project, str(arguments.get("instance", "")))
if not inst.get("container_id"):
return json.dumps({"logs": "", "status": inst["status"]})
lines: list = []
async def collect(line: str) -> None:
lines.append(line)
tail = int(arguments.get("tail", 200) or 200)
await get_backend().logs(inst["container_id"], follow=False, tail=max(1, min(tail, 2000)), on_log=collect)
return json.dumps({"logs": "\n".join(lines)})
async def _exec(self, arguments) -> str:
project = self._project(arguments)
inst = self._instance(project, str(arguments.get("instance", "")))
if not inst.get("container_id"):
raise ToolInputError("instance is not running")
command = str(arguments.get("command", "")).strip()
if not command:
raise ToolInputError("command is required")
result = await get_backend().exec(inst["container_id"], ["/bin/sh", "-c", command])
actor = self._actor_user()
store.record_event(inst, "exec", "user", actor["uid"], {"command": command, "exit_code": result.exit_code})
return json.dumps({"exit_code": result.exit_code, "output": result.output})
async def _stats(self, arguments) -> str:
project = self._project(arguments)
inst = self._instance(project, str(arguments.get("instance", "")))
return json.dumps({
"instance": inst["name"], "status": inst["status"],
"ingress_url": self._ingress_url(inst),
"runtime": api.instance_runtime(inst),
"stats": api.instance_stats(inst["uid"]),
"dockerfile_stats": api.dockerfile_stats(inst["dockerfile_uid"]),
}, default=str)
async def _schedule(self, arguments) -> str:
project = self._project(arguments)
inst = self._instance(project, str(arguments.get("instance", "")))
run_at = arguments.get("run_at")
try:
schedule = Schedule(
kind=str(arguments.get("kind", "")), cron=arguments.get("cron") or None,
run_at=from_iso(run_at) if run_at else None,
every_seconds=arguments.get("every_seconds"))
except ValueError as exc:
raise ToolInputError(str(exc))
sched = api.add_schedule(inst, str(arguments.get("action", "")), schedule)
return json.dumps({"schedule": sched}, default=str)
+2
View File
@@ -6,6 +6,7 @@ from .actions.avatar_actions import AVATAR_ACTIONS
from .actions.catalog import ACTIONS
from .actions.chunk_actions import CHUNK_ACTIONS
from .actions.client_actions import CLIENT_ACTIONS
from .actions.container_actions import CONTAINER_ACTIONS
from .actions.cost_actions import COST_ACTIONS
from .actions.docs_actions import DOCS_ACTIONS
from .actions.fetch_actions import FETCH_ACTIONS
@@ -25,4 +26,5 @@ CATALOG = Catalog(
+ COST_ACTIONS
+ CHUNK_ACTIONS
+ RSEARCH_ACTIONS
+ CONTAINER_ACTIONS
)
+16 -1
View File
@@ -108,6 +108,7 @@ class DeviiSession:
self._buffer: list[dict[str, Any]] = []
self._turn_tool_calls = 0
self._pending_session: str | None = None
self._turns: set[asyncio.Task] = set()
def restore_history(self, messages: list[dict[str, Any]]) -> None:
if messages:
@@ -222,7 +223,21 @@ class DeviiSession:
await self._emit({"type": "clear"}, buffer=False)
def spawn_turn(self, text: str) -> None:
asyncio.create_task(self._run_turn(text))
task = asyncio.create_task(self._run_turn(text))
self._turns.add(task)
task.add_done_callback(self._turns.discard)
async def cancel_turns(self) -> None:
turns = [task for task in self._turns if not task.done()]
if not turns:
return
for task in turns:
task.cancel()
await asyncio.gather(*turns, return_exceptions=True)
async def reset(self) -> None:
await self.cancel_turns()
await self.reset_conversation()
async def _run_turn(self, text: str) -> None:
turn_id = uuid_utils.uuid7().hex
+3 -3
View File
@@ -8,15 +8,15 @@ import shutil
import sys
from pathlib import Path
from devplacepy.config import BASE_DIR, STATIC_DIR
from devplacepy.config import BASE_DIR, DATA_DIR
from devplacepy import project_files
from devplacepy.services.jobs.base import JobService
from devplacepy.utils import generate_uid, slugify
logger = logging.getLogger(__name__)
ZIPS_DIR = STATIC_DIR / "uploads" / "zips"
STAGING_DIR = STATIC_DIR / "uploads" / "zip_staging"
ZIPS_DIR = DATA_DIR / "zips"
STAGING_DIR = DATA_DIR / "zip_staging"
WORKER_MODULE = "devplacepy.services.jobs.zip_worker"