forked from retoor/devplacepy
feat: add project fork async job service with cli prune/clear commands and shared container image build target
This commit is contained in:
+23
-98
@@ -6,19 +6,16 @@ import requests
|
||||
|
||||
from tests.conftest import BASE_URL, login_user, run_async
|
||||
from devplacepy.database import db, get_table, init_db, refresh_snapshot
|
||||
from devplacepy import project_files
|
||||
from devplacepy import config, project_files
|
||||
from devplacepy.services.containers import api, store, runtime
|
||||
from devplacepy.services.containers.backend.base import Mount, PortMapping, RunSpec
|
||||
from devplacepy.services.containers.backend.docker_cli import build_run_argv, parse_size
|
||||
from devplacepy.services.containers.backend.fake import FakeBackend
|
||||
from devplacepy.services.containers.build_service import ContainerBuildService
|
||||
from devplacepy.services.containers.service import ContainerService
|
||||
from devplacepy.services.containers.api import INSTANCE_LABEL
|
||||
from devplacepy.services.devii.tasks.schedule import Schedule, now_utc, to_iso
|
||||
from devplacepy.services.jobs import queue
|
||||
|
||||
_CONTAINER_TABLES = ("dockerfiles", "dockerfile_versions", "builds", "instances",
|
||||
"instance_events", "instance_metrics", "instance_schedules")
|
||||
_CONTAINER_TABLES = ("instances", "instance_events", "instance_metrics", "instance_schedules")
|
||||
|
||||
|
||||
@pytest.fixture(autouse=True)
|
||||
@@ -32,41 +29,26 @@ def env(tmp_path, monkeypatch):
|
||||
fake = FakeBackend()
|
||||
runtime.set_backend(fake)
|
||||
monkeypatch.setattr("devplacepy.config.CONTAINER_WORKSPACES_DIR", tmp_path / "ws")
|
||||
monkeypatch.setattr("devplacepy.config.CONTAINER_BUILD_CONTEXTS_DIR", tmp_path / "ctx")
|
||||
pid = "ctest-p1"
|
||||
project = {"uid": pid, "slug": "ctest", "title": "C", "user_uid": "ctest-u1"}
|
||||
user = {"uid": "ctest-u1", "username": "ctestadmin"}
|
||||
project_files.write_text_file(pid, user, "app.py", "print(1)\n")
|
||||
yield {"fake": fake, "project": project, "user": user}
|
||||
runtime.set_backend(None)
|
||||
for table in _CONTAINER_TABLES + ("jobs",):
|
||||
for table in _CONTAINER_TABLES:
|
||||
if table in db.tables:
|
||||
rows = [r for r in get_table(table).find()]
|
||||
for row in rows:
|
||||
if str(row.get("project_uid", row.get("kind", ""))).startswith("ctest") or row.get("kind") == "container_build":
|
||||
for row in [r for r in get_table(table).find()]:
|
||||
if str(row.get("project_uid", "")).startswith("ctest"):
|
||||
get_table(table).delete(uid=row["uid"])
|
||||
for row in list(get_table("project_files").find()):
|
||||
if str(row.get("project_uid", "")).startswith("ctest"):
|
||||
get_table("project_files").delete(uid=row["uid"])
|
||||
|
||||
|
||||
def _drain_builds():
|
||||
async def drive():
|
||||
service = ContainerBuildService()
|
||||
for _ in range(80):
|
||||
await service.run_once()
|
||||
refresh_snapshot()
|
||||
pending = [j for j in queue.list_jobs(kind="container_build") if j["status"] in ("pending", "running")]
|
||||
if not pending and not service._inflight:
|
||||
return
|
||||
await asyncio.sleep(0.02)
|
||||
run_async(drive())
|
||||
|
||||
|
||||
# ---------------- backend argv ----------------
|
||||
|
||||
def test_build_run_argv_exact():
|
||||
spec = RunSpec(image="myapp:3", name="inst", labels={INSTANCE_LABEL: "u1"}, env={"A": "B"},
|
||||
spec = RunSpec(image="ppy:latest", name="inst", labels={INSTANCE_LABEL: "u1"}, env={"A": "B"},
|
||||
cpu_limit="1.5", mem_limit="512m", ports=[PortMapping(8080, 80)],
|
||||
mounts=[Mount("/ws", "/app")], restart_policy="on-failure", command=["python", "app.py"])
|
||||
argv = build_run_argv(spec)
|
||||
@@ -76,11 +58,11 @@ def test_build_run_argv_exact():
|
||||
assert "-p" in argv and "8080:80/tcp" in argv
|
||||
assert "-v" in argv and "/ws:/app:rw" in argv
|
||||
assert "--restart" in argv and "on-failure" in argv
|
||||
assert argv[-3:] == ["myapp:3", "python", "app.py"]
|
||||
assert argv[-3:] == ["ppy:latest", "python", "app.py"]
|
||||
|
||||
|
||||
def test_never_policy_not_passed_to_docker():
|
||||
spec = RunSpec(image="i:1", name="n", restart_policy="never")
|
||||
spec = RunSpec(image="ppy:latest", name="n", restart_policy="never")
|
||||
assert "--restart" not in build_run_argv(spec)
|
||||
|
||||
|
||||
@@ -89,51 +71,29 @@ def test_parse_size():
|
||||
assert parse_size("512MB") == 512 * 1024 ** 2
|
||||
|
||||
|
||||
# ---------------- versioning + builds ----------------
|
||||
# ---------------- instance creation ----------------
|
||||
|
||||
def test_create_dockerfile_and_autobuild(env):
|
||||
result = run_async(api.create_dockerfile(env["project"], env["user"], name="web", description="d"))
|
||||
df = result["dockerfile"]
|
||||
assert df["current_version"] == 1
|
||||
assert result["build"]["build_number"] == 1
|
||||
_drain_builds()
|
||||
build = store.get_build(result["build"]["uid"])
|
||||
assert build["status"] == store.BUILD_SUCCESS
|
||||
assert env["fake"].built_images and "web:1" in env["fake"].built_images[0]
|
||||
def test_create_instance_uses_shared_image(env):
|
||||
inst = run_async(api.create_instance(env["project"], name="inst", actor=("user", env["user"]["uid"])))
|
||||
assert inst["name"] == "inst"
|
||||
assert inst["owner_uid"] == "ctest-u1"
|
||||
spec = api.run_spec_for(inst, config.CONTAINER_IMAGE)
|
||||
assert spec.image == config.CONTAINER_IMAGE
|
||||
assert any(m.container == "/app" for m in spec.mounts)
|
||||
|
||||
|
||||
def test_unchanged_content_skips_build(env):
|
||||
result = run_async(api.create_dockerfile(env["project"], env["user"], name="web", content="FROM scratch\n"))
|
||||
df = store.get_dockerfile(result["dockerfile"]["uid"])
|
||||
same = run_async(api.save_version(df, env["user"], "FROM scratch\n"))
|
||||
assert same["changed"] is False and same["build"] is None
|
||||
changed = run_async(api.save_version(df, env["user"], "FROM scratch\nRUN echo hi\n"))
|
||||
assert changed["changed"] is True
|
||||
assert store.get_dockerfile(df["uid"])["current_version"] == 2
|
||||
assert changed["build"]["build_number"] == 2
|
||||
|
||||
|
||||
def test_build_failure_recorded(env):
|
||||
env["fake"].fail_build = True
|
||||
result = run_async(api.create_dockerfile(env["project"], env["user"], name="web"))
|
||||
_drain_builds()
|
||||
assert store.get_build(result["build"]["uid"])["status"] == store.BUILD_FAILED
|
||||
|
||||
|
||||
def test_duplicate_name_rejected(env):
|
||||
run_async(api.create_dockerfile(env["project"], env["user"], name="web"))
|
||||
def test_create_instance_requires_built_image(env):
|
||||
async def no_image(ref):
|
||||
return False
|
||||
env["fake"].image_exists = no_image
|
||||
with pytest.raises(api.ContainerError):
|
||||
run_async(api.create_dockerfile(env["project"], env["user"], name="web"))
|
||||
run_async(api.create_instance(env["project"], name="inst"))
|
||||
|
||||
|
||||
# ---------------- reconcile ----------------
|
||||
|
||||
def _ready_instance(env, **kwargs):
|
||||
result = run_async(api.create_dockerfile(env["project"], env["user"], name="web"))
|
||||
_drain_builds()
|
||||
df = store.get_dockerfile(result["dockerfile"]["uid"])
|
||||
build = store.get_build(result["build"]["uid"])
|
||||
return run_async(api.create_instance(env["project"], df, build, name=kwargs.pop("name", "inst"), **kwargs))
|
||||
return run_async(api.create_instance(env["project"], name=kwargs.pop("name", "inst"), **kwargs))
|
||||
|
||||
|
||||
def test_reconcile_launches_and_stops(env):
|
||||
@@ -153,7 +113,7 @@ def test_reconcile_launches_and_stops(env):
|
||||
|
||||
def test_reconcile_reaps_orphan(env):
|
||||
fake = env["fake"]
|
||||
run_async(fake.run(RunSpec(image="x:1", name="ghost", labels={INSTANCE_LABEL: "missing-uid"})))
|
||||
run_async(fake.run(RunSpec(image="ppy:latest", name="ghost", labels={INSTANCE_LABEL: "missing-uid"})))
|
||||
service = ContainerService()
|
||||
run_async(service.run_once())
|
||||
assert not run_async(fake.ps())
|
||||
@@ -173,26 +133,6 @@ def test_reconcile_removes_marked_instance(env):
|
||||
assert not run_async(env["fake"].ps())
|
||||
|
||||
|
||||
def test_build_status_not_flagged_as_error(env):
|
||||
from types import SimpleNamespace
|
||||
from devplacepy.services.devii.container import ContainerController
|
||||
from devplacepy.services.devii.agentic.loop import _is_error, _summary
|
||||
df = store.create_dockerfile("ctest-p1", env["user"], "svc", "", "")
|
||||
ver = store.create_version(df, "FROM scratch", env["user"])
|
||||
build = store.create_build(df, ver, 1, "svc:1", True)
|
||||
store.update_build(build["uid"], {"status": "building"})
|
||||
ctl = ContainerController(SimpleNamespace(username=None))
|
||||
result = run_async(ctl.dispatch("container_build_status", {"build_uid": build["uid"]}))
|
||||
assert _is_error(result) is False
|
||||
assert _summary(result) == "building"
|
||||
store.update_build(build["uid"], {"status": "failed", "error": "docker build exited 1"})
|
||||
failed = run_async(ctl.dispatch("container_build_status", {"build_uid": build["uid"]}))
|
||||
import json as _json
|
||||
assert _is_error(failed) is False
|
||||
assert "error" not in _json.loads(failed)
|
||||
assert _json.loads(failed)["build_error"] == "docker build exited 1"
|
||||
|
||||
|
||||
def test_schedule_fires(env):
|
||||
inst = _ready_instance(env, autostart=False)
|
||||
assert inst["desired_state"] == store.DESIRED_STOPPED
|
||||
@@ -281,18 +221,3 @@ def test_http_ingress_proxy(app_server):
|
||||
finally:
|
||||
httpd.shutdown()
|
||||
get_table("instances").delete(uid=uid)
|
||||
|
||||
|
||||
def test_http_admin_create_dockerfile(app_server, page, seeded_db):
|
||||
_promote_admin("alice_test")
|
||||
key = _api_key("alice_test")
|
||||
headers = {"X-API-KEY": key, "Accept": "application/json"}
|
||||
project = requests.post(f"{BASE_URL}/projects/create", headers=headers,
|
||||
data={"title": "WithAdmin", "description": "x", "project_type": "software", "status": "s"})
|
||||
slug = project.json()["data"]["slug"] or project.json()["data"]["uid"]
|
||||
r = requests.post(f"{BASE_URL}/projects/{slug}/containers/dockerfiles", headers=headers,
|
||||
data={"name": "svc", "description": "", "tags": ""})
|
||||
assert r.status_code == 200, r.text
|
||||
assert r.json()["data"]["dockerfile"]["name"] == "svc"
|
||||
data = requests.get(f"{BASE_URL}/projects/{slug}/containers/data", headers=headers).json()
|
||||
assert any(d["name"] == "svc" for d in data["dockerfiles"])
|
||||
|
||||
@@ -0,0 +1,208 @@
|
||||
import asyncio
|
||||
from datetime import datetime, timezone
|
||||
from pathlib import Path
|
||||
|
||||
import pytest
|
||||
|
||||
from devplacepy.database import init_db, get_table, refresh_snapshot
|
||||
from devplacepy import project_files
|
||||
from devplacepy.services.jobs import queue
|
||||
from devplacepy.services.jobs.fork_service import ForkService
|
||||
from tests.conftest import run_async
|
||||
|
||||
|
||||
@pytest.fixture(autouse=True)
|
||||
def _init_db():
|
||||
init_db()
|
||||
yield
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def fork_env(tmp_path, monkeypatch):
|
||||
monkeypatch.setattr("devplacepy.services.jobs.fork_service.STAGING_DIR", tmp_path / "staging")
|
||||
monkeypatch.setattr("devplacepy.project_files.PROJECT_FILES_DIR", tmp_path / "pf")
|
||||
yield tmp_path
|
||||
jobs = get_table("jobs")
|
||||
for row in list(jobs.find(kind="fork")):
|
||||
jobs.delete(uid=row["uid"])
|
||||
projects = get_table("projects")
|
||||
files = get_table("project_files")
|
||||
for project in list(projects.find()):
|
||||
if str(project.get("user_uid", "")).startswith("forktest-owner"):
|
||||
for node in list(files.find(project_uid=project["uid"])):
|
||||
files.delete(uid=node["uid"])
|
||||
projects.delete(uid=project["uid"])
|
||||
for node in list(files.find()):
|
||||
if str(node.get("project_uid", "")).startswith("forktest"):
|
||||
files.delete(uid=node["uid"])
|
||||
forks = get_table("project_forks")
|
||||
for relation in list(forks.find()):
|
||||
if str(relation.get("forked_by_uid", "")).startswith("forktest-owner"):
|
||||
forks.delete(uid=relation["uid"])
|
||||
users = get_table("users")
|
||||
for user in list(users.find()):
|
||||
if str(user.get("uid", "")).startswith("forktest-owner"):
|
||||
users.delete(uid=user["uid"])
|
||||
|
||||
|
||||
_counter = [0]
|
||||
|
||||
|
||||
def _make_source_project(*, is_private=False, binary=False):
|
||||
_counter[0] += 1
|
||||
pid = f"forktest-{_counter[0]}"
|
||||
owner_uid = f"forktest-owner-{_counter[0]}"
|
||||
user = {"uid": owner_uid, "username": f"forktester{_counter[0]}"}
|
||||
get_table("users").insert({"uid": owner_uid, "username": user["username"], "xp": 0, "level": 1})
|
||||
get_table("projects").insert({
|
||||
"uid": pid,
|
||||
"user_uid": owner_uid,
|
||||
"slug": f"{pid}-source",
|
||||
"title": "Source Project",
|
||||
"description": "the original",
|
||||
"project_type": "software",
|
||||
"platforms": "linux",
|
||||
"status": "Released",
|
||||
"is_private": 1 if is_private else 0,
|
||||
"read_only": 0,
|
||||
"stars": 0,
|
||||
"created_at": datetime.now(timezone.utc).isoformat(),
|
||||
})
|
||||
project_files.write_text_file(pid, user, "README.md", "# hello\nworld")
|
||||
project_files.write_text_file(pid, user, "src/app.py", "print(1)\n")
|
||||
if binary:
|
||||
project_files.store_upload(pid, user, "assets", "logo.bin", bytes(range(256)))
|
||||
return pid, owner_uid
|
||||
|
||||
|
||||
def _tree(directory):
|
||||
root = Path(directory)
|
||||
out = {}
|
||||
for path in sorted(root.rglob("*")):
|
||||
if path.is_file():
|
||||
out[path.relative_to(root).as_posix()] = path.read_bytes()
|
||||
return out
|
||||
|
||||
|
||||
def _process_fork_jobs():
|
||||
async def drive():
|
||||
svc = ForkService()
|
||||
for _ in range(400):
|
||||
await svc.run_once()
|
||||
refresh_snapshot()
|
||||
pending = [r for r in get_table("jobs").find(kind="fork") if r["status"] in ("pending", "running")]
|
||||
if not pending and not svc._inflight:
|
||||
return
|
||||
await asyncio.sleep(0.05)
|
||||
run_async(drive())
|
||||
|
||||
|
||||
def _enqueue(source_uid, owner_uid, title="My Fork"):
|
||||
return queue.enqueue(
|
||||
"fork",
|
||||
{"source_project_uid": source_uid, "title": title, "forked_by_uid": owner_uid},
|
||||
"user", owner_uid, title,
|
||||
)
|
||||
|
||||
|
||||
def test_enqueue_creates_pending_job(fork_env):
|
||||
uid = _enqueue("p", "forktest-owner-x")
|
||||
job = queue.get_job(uid)
|
||||
assert job["status"] == "pending"
|
||||
assert job["kind"] == "fork"
|
||||
assert job["payload"]["source_project_uid"] == "p"
|
||||
|
||||
|
||||
def test_process_creates_fork_with_files(fork_env, tmp_path):
|
||||
pid, owner_uid = _make_source_project(binary=True)
|
||||
uid = _enqueue(pid, owner_uid, title="Forked Copy")
|
||||
_process_fork_jobs()
|
||||
job = queue.get_job(uid)
|
||||
assert job["status"] == "done"
|
||||
result = job["result"]
|
||||
new_uid = result["project_uid"]
|
||||
forked = get_table("projects").find_one(uid=new_uid)
|
||||
assert forked is not None
|
||||
assert forked["user_uid"] == owner_uid
|
||||
assert forked["title"] == "Forked Copy"
|
||||
assert result["project_url"] == f"/projects/{forked['slug']}"
|
||||
|
||||
source_dir = tmp_path / "exp_source"
|
||||
fork_dir = tmp_path / "exp_fork"
|
||||
project_files.export_to_dir(pid, "", source_dir)
|
||||
project_files.export_to_dir(new_uid, "", fork_dir)
|
||||
assert _tree(fork_dir) == _tree(source_dir)
|
||||
|
||||
|
||||
def test_fork_relation_direction(fork_env):
|
||||
pid, owner_uid = _make_source_project()
|
||||
uid = _enqueue(pid, owner_uid)
|
||||
_process_fork_jobs()
|
||||
new_uid = queue.get_job(uid)["result"]["project_uid"]
|
||||
relation = get_table("project_forks").find_one(forked_project_uid=new_uid)
|
||||
assert relation is not None
|
||||
assert relation["source_project_uid"] == pid
|
||||
assert relation["forked_project_uid"] == new_uid
|
||||
assert relation["forked_by_uid"] == owner_uid
|
||||
|
||||
|
||||
def test_fork_copies_binary_file(fork_env, tmp_path):
|
||||
pid, owner_uid = _make_source_project(binary=True)
|
||||
uid = _enqueue(pid, owner_uid)
|
||||
_process_fork_jobs()
|
||||
new_uid = queue.get_job(uid)["result"]["project_uid"]
|
||||
fork_dir = tmp_path / "binfork"
|
||||
project_files.export_to_dir(new_uid, "", fork_dir)
|
||||
assert (fork_dir / "assets" / "logo.bin").read_bytes() == bytes(range(256))
|
||||
|
||||
|
||||
def test_fork_preserves_private_flag(fork_env):
|
||||
pid, owner_uid = _make_source_project(is_private=True)
|
||||
uid = _enqueue(pid, owner_uid)
|
||||
_process_fork_jobs()
|
||||
new_uid = queue.get_job(uid)["result"]["project_uid"]
|
||||
assert get_table("projects").find_one(uid=new_uid)["is_private"] == 1
|
||||
|
||||
|
||||
def test_missing_source_fails_job_without_orphan(fork_env):
|
||||
before = {p["uid"] for p in get_table("projects").find()}
|
||||
uid = _enqueue("does-not-exist", "forktest-owner-missing")
|
||||
_process_fork_jobs()
|
||||
job = queue.get_job(uid)
|
||||
assert job["status"] == "failed"
|
||||
assert "source project not found" in job["error"].lower()
|
||||
after = {p["uid"] for p in get_table("projects").find()}
|
||||
assert before == after
|
||||
|
||||
|
||||
def test_cleanup_keeps_project(fork_env):
|
||||
pid, owner_uid = _make_source_project()
|
||||
uid = _enqueue(pid, owner_uid)
|
||||
_process_fork_jobs()
|
||||
job = queue.get_job(uid)
|
||||
new_uid = job["result"]["project_uid"]
|
||||
ForkService().cleanup(job)
|
||||
assert get_table("projects").find_one(uid=new_uid) is not None
|
||||
|
||||
|
||||
def test_retention_sweep_keeps_project(fork_env):
|
||||
pid, owner_uid = _make_source_project()
|
||||
uid = _enqueue(pid, owner_uid)
|
||||
_process_fork_jobs()
|
||||
new_uid = queue.get_job(uid)["result"]["project_uid"]
|
||||
get_table("jobs").update({"uid": uid, "expires_at": "2000-01-01T00:00:00+00:00"}, ["uid"])
|
||||
svc = ForkService()
|
||||
run_async(svc.run_once())
|
||||
refresh_snapshot()
|
||||
assert queue.get_job(uid) is None
|
||||
assert get_table("projects").find_one(uid=new_uid) is not None
|
||||
|
||||
|
||||
def test_orphan_running_recovered_on_enable(fork_env):
|
||||
uid = _enqueue("p", "forktest-owner-o")
|
||||
get_table("jobs").update({"uid": uid, "status": "running", "started_at": "2020-01-01T00:00:00+00:00"}, ["uid"])
|
||||
svc = ForkService()
|
||||
run_async(svc.on_enable())
|
||||
job = queue.get_job(uid)
|
||||
assert job["status"] == "pending"
|
||||
assert job["retry_count"] == 1
|
||||
Reference in New Issue
Block a user