feat: add zip file extraction support with error handling for invalid archives
This commit is contained in:
@@ -133,6 +133,45 @@ def cmd_attachments_prune(args):
|
||||
print(f"Pruned {len(orphans)} orphan attachment(s) older than {args.hours}h")
|
||||
|
||||
|
||||
def _remove_zip_artifacts(job):
|
||||
import shutil
|
||||
from pathlib import Path
|
||||
from devplacepy.services.jobs.zip_service import STAGING_DIR
|
||||
local_path = (job.get("result") or {}).get("local_path")
|
||||
if local_path:
|
||||
Path(local_path).unlink(missing_ok=True)
|
||||
shutil.rmtree(STAGING_DIR / job["uid"], ignore_errors=True)
|
||||
|
||||
|
||||
def cmd_zips_prune(args):
|
||||
from datetime import datetime, timezone
|
||||
from devplacepy.services.jobs import queue
|
||||
now = datetime.now(timezone.utc)
|
||||
removed = 0
|
||||
for job in queue.list_jobs(kind="zip", status=queue.DONE):
|
||||
expires_at = job.get("expires_at")
|
||||
if not expires_at:
|
||||
continue
|
||||
try:
|
||||
expiry = datetime.fromisoformat(expires_at)
|
||||
except (ValueError, TypeError):
|
||||
continue
|
||||
if expiry < now:
|
||||
_remove_zip_artifacts(job)
|
||||
get_table("jobs").delete(uid=job["uid"])
|
||||
removed += 1
|
||||
print(f"Pruned {removed} expired zip job(s)")
|
||||
|
||||
|
||||
def cmd_zips_clear(args):
|
||||
from devplacepy.services.jobs import queue
|
||||
jobs = queue.list_jobs(kind="zip")
|
||||
for job in jobs:
|
||||
_remove_zip_artifacts(job)
|
||||
get_table("jobs").delete(uid=job["uid"])
|
||||
print(f"Cleared {len(jobs)} zip job(s) and their archives")
|
||||
|
||||
|
||||
def main():
|
||||
parser = argparse.ArgumentParser(description="DevPlace admin CLI")
|
||||
sub = parser.add_subparsers(title="commands", dest="command")
|
||||
@@ -181,6 +220,13 @@ def main():
|
||||
devii_reset.add_argument("--all", action="store_true", help="Reset every quota (users and guests)")
|
||||
devii_reset.set_defaults(func=cmd_devii_reset_quota)
|
||||
|
||||
zips = sub.add_parser("zips", help="Zip archive job management")
|
||||
zips_sub = zips.add_subparsers(title="action", dest="action")
|
||||
zips_prune = zips_sub.add_parser("prune", help="Delete expired zip archives and their job rows")
|
||||
zips_prune.set_defaults(func=cmd_zips_prune)
|
||||
zips_clear = zips_sub.add_parser("clear", help="Delete every zip archive and job row")
|
||||
zips_clear.set_defaults(func=cmd_zips_clear)
|
||||
|
||||
args = parser.parse_args()
|
||||
if hasattr(args, "func"):
|
||||
args.func(args)
|
||||
|
||||
@@ -166,6 +166,9 @@ def init_db():
|
||||
_index(db, "gateway_usage_ledger", "idx_gw_usage_owner_time", ["owner_kind", "owner_id", "created_at"])
|
||||
_index(db, "gateway_usage_ledger", "idx_gw_usage_endpoint_time", ["endpoint", "created_at"])
|
||||
_index(db, "gateway_concurrency_samples", "idx_gw_conc_time", ["created_at"])
|
||||
_index(db, "jobs", "idx_jobs_kind_status", ["kind", "status"])
|
||||
_index(db, "jobs", "idx_jobs_owner", ["owner_kind", "owner_id"])
|
||||
_index(db, "jobs", "idx_jobs_expires", ["expires_at"])
|
||||
|
||||
if "news" in tables:
|
||||
for article in db["news"].find(status=None):
|
||||
|
||||
@@ -540,6 +540,50 @@ four ways to sign requests.
|
||||
destructive=True,
|
||||
params=[field("project_slug", "path", "string", True, "PROJECT_SLUG", "Slug or UID of the project.")],
|
||||
),
|
||||
endpoint(
|
||||
id="projects-zip",
|
||||
method="POST",
|
||||
path="/projects/{project_slug}/zip",
|
||||
title="Queue a project zip",
|
||||
summary="Start a background job that archives the whole project. Returns the job uid and status URL to poll.",
|
||||
auth="public",
|
||||
params=[field("project_slug", "path", "string", True, "PROJECT_SLUG", "Slug or UID of the project.")],
|
||||
sample_response={"uid": "ZIP_JOB_UID", "status_url": "/zips/ZIP_JOB_UID"},
|
||||
),
|
||||
endpoint(
|
||||
id="zips-status",
|
||||
method="GET",
|
||||
path="/zips/{uid}",
|
||||
title="Zip job status",
|
||||
summary="Poll a zip job. While pending or running download_url is null; once done it points at the archive.",
|
||||
auth="public",
|
||||
params=[field("uid", "path", "string", True, "ZIP_JOB_UID", "Zip job uid returned when the job was queued.")],
|
||||
sample_response={
|
||||
"uid": "ZIP_JOB_UID",
|
||||
"kind": "zip",
|
||||
"status": "done",
|
||||
"preferred_name": "my-project",
|
||||
"download_url": "/zips/ZIP_JOB_UID/download",
|
||||
"error": None,
|
||||
"bytes_in": 20480,
|
||||
"bytes_out": 8192,
|
||||
"item_count": 12,
|
||||
"file_count": 10,
|
||||
"dir_count": 2,
|
||||
"created_at": "2026-06-09T10:00:00+00:00",
|
||||
"completed_at": "2026-06-09T10:00:03+00:00",
|
||||
},
|
||||
),
|
||||
endpoint(
|
||||
id="zips-download",
|
||||
method="GET",
|
||||
path="/zips/{uid}/download",
|
||||
title="Download a zip archive",
|
||||
summary="Stream the finished archive as application/zip. Each access extends the retention window.",
|
||||
auth="public",
|
||||
interactive=True,
|
||||
params=[field("uid", "path", "string", True, "ZIP_JOB_UID", "Zip job uid of a finished job.")],
|
||||
),
|
||||
endpoint(
|
||||
id="gists-list",
|
||||
method="GET",
|
||||
@@ -1051,6 +1095,19 @@ sign requests.
|
||||
],
|
||||
sample_response={"ok": True, "redirect": "/projects/PROJECT_SLUG/files", "data": {"path": "src/old.py"}},
|
||||
),
|
||||
endpoint(
|
||||
id="project-files-zip",
|
||||
method="POST",
|
||||
path="/projects/{project_slug}/files/zip",
|
||||
title="Queue a zip of files",
|
||||
summary="Archive the whole tree, or a subtree via the path query. Returns the job uid and status URL to poll with /zips/{uid}.",
|
||||
auth="public",
|
||||
params=[
|
||||
field("project_slug", "path", "string", True, "PROJECT_SLUG", "Project slug or uid."),
|
||||
field("path", "query", "string", False, "src", "Relative file or directory to archive; empty for the whole project."),
|
||||
],
|
||||
sample_response={"uid": "ZIP_JOB_UID", "status_url": "/zips/ZIP_JOB_UID"},
|
||||
),
|
||||
],
|
||||
},
|
||||
{
|
||||
|
||||
+4
-1
@@ -17,12 +17,13 @@ from devplacepy.schemas import ValidationErrorOut
|
||||
from fastapi.responses import JSONResponse
|
||||
from devplacepy.utils import get_current_user, time_ago
|
||||
from devplacepy.seo import base_seo_context, site_url, website_schema
|
||||
from devplacepy.routers import auth, feed, posts, comments, projects, project_files, profile, messages, notifications, votes, avatar, follow, admin, seo, bugs, news, gists, services as services_router, uploads, push, leaderboard, reactions, bookmarks, polls, docs, openai_gateway, devii
|
||||
from devplacepy.routers import auth, feed, posts, comments, projects, project_files, profile, messages, notifications, votes, avatar, follow, admin, seo, bugs, news, gists, services as services_router, uploads, push, leaderboard, reactions, bookmarks, polls, docs, openai_gateway, devii, zips
|
||||
from devplacepy.services.manager import service_manager
|
||||
from devplacepy.services.news import NewsService
|
||||
from devplacepy.services.bot import BotsService
|
||||
from devplacepy.services.openai_gateway import GatewayService
|
||||
from devplacepy.services.devii import DeviiService
|
||||
from devplacepy.services.jobs.zip_service import ZipService
|
||||
|
||||
logging.basicConfig(
|
||||
level=logging.INFO,
|
||||
@@ -178,6 +179,7 @@ app.include_router(gists.router, prefix="/gists")
|
||||
app.include_router(news.router, prefix="/news")
|
||||
app.include_router(services_router.router, prefix="/admin/services")
|
||||
app.include_router(uploads.router, prefix="/uploads")
|
||||
app.include_router(zips.router, prefix="/zips")
|
||||
|
||||
|
||||
@app.middleware("http")
|
||||
@@ -244,6 +246,7 @@ async def startup():
|
||||
service_manager.register(BotsService())
|
||||
service_manager.register(GatewayService())
|
||||
service_manager.register(DeviiService())
|
||||
service_manager.register(ZipService())
|
||||
if not os.environ.get("DEVPLACE_DISABLE_SERVICES"):
|
||||
if acquire_service_lock():
|
||||
service_manager.set_lock_owner(True)
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
# retoor <retoor@molodetz.nl>
|
||||
|
||||
import logging
|
||||
import shutil
|
||||
from datetime import datetime, timezone
|
||||
from pathlib import Path
|
||||
|
||||
@@ -356,6 +357,34 @@ def read_file(project_uid: str, raw_path: str) -> dict:
|
||||
return payload
|
||||
|
||||
|
||||
def export_to_dir(project_uid: str, subpath: str, dest_dir) -> int:
|
||||
dest = Path(dest_dir).resolve()
|
||||
dest.mkdir(parents=True, exist_ok=True)
|
||||
if subpath:
|
||||
base = normalize_path(subpath)
|
||||
rows = _descendants(project_uid, base)
|
||||
strip = _parent_of(base)
|
||||
else:
|
||||
rows = list(_table().find(project_uid=project_uid))
|
||||
strip = ""
|
||||
written = 0
|
||||
for row in sorted(rows, key=lambda r: r["path"]):
|
||||
relative = row["path"][len(strip):].lstrip("/") if strip else row["path"]
|
||||
target = (dest / relative).resolve()
|
||||
if target != dest and not target.is_relative_to(dest):
|
||||
raise ProjectFileError("path escapes export directory")
|
||||
if row["type"] == "dir":
|
||||
target.mkdir(parents=True, exist_ok=True)
|
||||
continue
|
||||
target.parent.mkdir(parents=True, exist_ok=True)
|
||||
if row.get("is_binary") and row.get("stored_name") and row.get("directory"):
|
||||
shutil.copyfile(PROJECT_FILES_DIR / row["directory"] / row["stored_name"], target)
|
||||
else:
|
||||
target.write_text(row.get("content") or "", encoding="utf-8")
|
||||
written += 1
|
||||
return written
|
||||
|
||||
|
||||
def delete_all_project_files(project_uid: str) -> None:
|
||||
if "project_files" not in db.tables:
|
||||
return
|
||||
|
||||
@@ -49,13 +49,14 @@ DOCS_PAGES = [
|
||||
{"slug": "devii-security", "title": "Security and limits", "kind": "prose", "admin": True, "section": SECTION_DEVII},
|
||||
{"slug": "devii-config", "title": "Configuration and CLI", "kind": "prose", "admin": True, "section": SECTION_DEVII},
|
||||
# Services - the background service framework and every service in detail (admins only)
|
||||
{"slug": "services", "title": "Services overview", "kind": "prose", "admin": True, "section": SECTION_SERVICES},
|
||||
{"slug": "services-overview", "title": "Services overview", "kind": "prose", "admin": True, "section": SECTION_SERVICES},
|
||||
{"slug": "services-framework", "title": "Service framework", "kind": "prose", "admin": True, "section": SECTION_SERVICES},
|
||||
{"slug": "services-data", "title": "Reading service data", "kind": "prose", "admin": True, "section": SECTION_SERVICES},
|
||||
{"slug": "services-gateway", "title": "GatewayService (AI)", "kind": "prose", "admin": True, "section": SECTION_SERVICES},
|
||||
{"slug": "services-devii", "title": "DeviiService", "kind": "prose", "admin": True, "section": SECTION_SERVICES},
|
||||
{"slug": "services-news", "title": "NewsService", "kind": "prose", "admin": True, "section": SECTION_SERVICES},
|
||||
{"slug": "services-bots", "title": "BotsService", "kind": "prose", "admin": True, "section": SECTION_SERVICES},
|
||||
{"slug": "services-zip", "title": "ZipService", "kind": "prose", "admin": True, "section": SECTION_SERVICES},
|
||||
# Architecture - platform design, structure, and development process (admins only)
|
||||
{"slug": "architecture", "title": "Architecture overview", "kind": "prose", "admin": True, "section": SECTION_ARCH},
|
||||
{"slug": "architecture-backend", "title": "Backend and request pipeline", "kind": "prose", "admin": True, "section": SECTION_ARCH},
|
||||
@@ -63,6 +64,7 @@ DOCS_PAGES = [
|
||||
{"slug": "architecture-styling", "title": "Styling and templates", "kind": "prose", "admin": True, "section": SECTION_ARCH},
|
||||
{"slug": "architecture-conventions", "title": "Conventions and rules", "kind": "prose", "admin": True, "section": SECTION_ARCH},
|
||||
{"slug": "architecture-workflow", "title": "Development workflow", "kind": "prose", "admin": True, "section": SECTION_ARCH},
|
||||
{"slug": "architecture-jobs", "title": "Async job services", "kind": "prose", "admin": True, "section": SECTION_ARCH},
|
||||
# Testing - test framework, load testing, and make targets (admins only)
|
||||
{"slug": "testing", "title": "Testing overview", "kind": "prose", "admin": True, "section": SECTION_TESTING},
|
||||
{"slug": "testing-framework", "title": "Test framework and rules", "kind": "prose", "admin": True, "section": SECTION_TESTING},
|
||||
|
||||
@@ -16,6 +16,7 @@ from devplacepy.models import (
|
||||
)
|
||||
from devplacepy.responses import respond, action_result, wants_json, json_error
|
||||
from devplacepy.schemas import ProjectFilesOut
|
||||
from devplacepy.services.jobs import queue
|
||||
from devplacepy.seo import base_seo_context
|
||||
from devplacepy.utils import get_current_user, require_user, not_found
|
||||
from devplacepy import project_files
|
||||
@@ -78,6 +79,29 @@ async def project_file_raw(request: Request, project_slug: str, path: str):
|
||||
return json_error(404, str(exc))
|
||||
|
||||
|
||||
@router.post("/{project_slug}/files/zip")
|
||||
async def project_files_zip(request: Request, project_slug: str, path: str = ""):
|
||||
project = _load_project(project_slug)
|
||||
user = get_current_user(request)
|
||||
subpath = ""
|
||||
name = project.get("title") or "project"
|
||||
if path:
|
||||
try:
|
||||
subpath = project_files.normalize_path(path)
|
||||
except ProjectFileError as exc:
|
||||
return json_error(400, str(exc))
|
||||
if project_files.get_node(project["uid"], subpath) is None:
|
||||
return json_error(404, "Path not found")
|
||||
name = subpath.rsplit("/", 1)[-1]
|
||||
owner_kind, owner_id = ("user", user["uid"]) if user else ("guest", request.client.host or "anon")
|
||||
uid = queue.enqueue(
|
||||
"zip",
|
||||
{"source": {"type": "project_tree", "project_uid": project["uid"], "path": subpath}},
|
||||
owner_kind, owner_id, name,
|
||||
)
|
||||
return JSONResponse({"uid": uid, "status_url": f"/zips/{uid}"})
|
||||
|
||||
|
||||
@router.post("/{project_slug}/files/write")
|
||||
async def project_file_write(request: Request, project_slug: str, data: Annotated[ProjectFileWriteForm, Form()]):
|
||||
user = require_user(request)
|
||||
|
||||
@@ -3,8 +3,9 @@ from typing import Annotated
|
||||
from sqlalchemy import or_
|
||||
from fastapi import APIRouter, Request, Form
|
||||
from devplacepy.models import ProjectForm
|
||||
from fastapi.responses import HTMLResponse, RedirectResponse
|
||||
from devplacepy.database import get_table, get_users_by_uids, get_site_stats, get_user_votes, paginate
|
||||
from fastapi.responses import HTMLResponse, RedirectResponse, JSONResponse
|
||||
from devplacepy.database import get_table, get_users_by_uids, get_site_stats, get_user_votes, paginate, resolve_by_slug
|
||||
from devplacepy.services.jobs import queue
|
||||
from devplacepy.content import load_detail, delete_content_item, create_content_item, detail_context, canonical_redirect, first_image_url
|
||||
from devplacepy.templating import templates
|
||||
from devplacepy.utils import get_current_user, require_user, not_found, XP_PROJECT
|
||||
@@ -113,6 +114,22 @@ async def project_detail(request: Request, project_slug: str):
|
||||
}), model=ProjectDetailOut)
|
||||
|
||||
|
||||
@router.post("/{project_slug}/zip")
|
||||
async def zip_project(request: Request, project_slug: str):
|
||||
project = resolve_by_slug(get_table("projects"), project_slug)
|
||||
if not project:
|
||||
raise not_found("Project not found")
|
||||
user = get_current_user(request)
|
||||
owner_kind, owner_id = ("user", user["uid"]) if user else ("guest", request.client.host or "anon")
|
||||
uid = queue.enqueue(
|
||||
"zip",
|
||||
{"source": {"type": "project_tree", "project_uid": project["uid"], "path": ""}},
|
||||
owner_kind, owner_id,
|
||||
project.get("title") or "project",
|
||||
)
|
||||
return JSONResponse({"uid": uid, "status_url": f"/zips/{uid}"})
|
||||
|
||||
|
||||
@router.post("/delete/{project_slug}")
|
||||
async def delete_project(request: Request, project_slug: str):
|
||||
user = require_user(request)
|
||||
|
||||
@@ -0,0 +1,59 @@
|
||||
# retoor <retoor@molodetz.nl>
|
||||
|
||||
import logging
|
||||
from pathlib import Path
|
||||
|
||||
from fastapi import APIRouter, Request
|
||||
from fastapi.responses import FileResponse, JSONResponse
|
||||
|
||||
from devplacepy.services.jobs import queue
|
||||
from devplacepy.services.jobs.zip_service import ZIPS_DIR
|
||||
from devplacepy.schemas import ZipJobOut
|
||||
from devplacepy.utils import not_found
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
router = APIRouter()
|
||||
|
||||
RETENTION_EXTEND_SECONDS = 7 * 24 * 60 * 60
|
||||
|
||||
|
||||
def _status_payload(job: dict) -> dict:
|
||||
result = job.get("result", {})
|
||||
done = job.get("status") == queue.DONE
|
||||
return {
|
||||
"uid": job.get("uid", ""),
|
||||
"kind": job.get("kind", ""),
|
||||
"status": job.get("status", ""),
|
||||
"preferred_name": job.get("preferred_name"),
|
||||
"download_url": result.get("download_url") if done else None,
|
||||
"error": job.get("error") or None,
|
||||
"bytes_in": int(job.get("bytes_in") or 0),
|
||||
"bytes_out": int(job.get("bytes_out") or 0),
|
||||
"item_count": int(job.get("item_count") or 0),
|
||||
"file_count": int(result.get("file_count") or 0),
|
||||
"dir_count": int(result.get("dir_count") or 0),
|
||||
"created_at": job.get("created_at"),
|
||||
"completed_at": job.get("completed_at") or None,
|
||||
}
|
||||
|
||||
|
||||
@router.get("/{uid}")
|
||||
async def zip_status(request: Request, uid: str):
|
||||
job = queue.get_job(uid)
|
||||
if not job or job.get("kind") != "zip":
|
||||
raise not_found("Job not found")
|
||||
return JSONResponse(ZipJobOut.model_validate(_status_payload(job)).model_dump(mode="json"))
|
||||
|
||||
|
||||
@router.get("/{uid}/download")
|
||||
async def zip_download(request: Request, uid: str):
|
||||
job = queue.get_job(uid)
|
||||
if not job or job.get("kind") != "zip" or job.get("status") != queue.DONE:
|
||||
raise not_found("Archive not available")
|
||||
result = job.get("result", {})
|
||||
local_path = result.get("local_path", "")
|
||||
resolved = Path(local_path).resolve()
|
||||
if not local_path or not resolved.is_relative_to(ZIPS_DIR.resolve()) or not resolved.is_file():
|
||||
raise not_found("Archive not available")
|
||||
queue.touch_job(uid, RETENTION_EXTEND_SECONDS)
|
||||
return FileResponse(resolved, filename=result.get("final_name", "download.zip"), media_type="application/zip")
|
||||
@@ -34,6 +34,22 @@ class ValidationErrorOut(_Out):
|
||||
messages: list = []
|
||||
|
||||
|
||||
class ZipJobOut(_Out):
|
||||
uid: str = ""
|
||||
kind: str = ""
|
||||
status: str = ""
|
||||
preferred_name: Optional[str] = None
|
||||
download_url: Optional[str] = None
|
||||
error: Optional[str] = None
|
||||
bytes_in: int = 0
|
||||
bytes_out: int = 0
|
||||
item_count: int = 0
|
||||
file_count: int = 0
|
||||
dir_count: int = 0
|
||||
created_at: Optional[str] = None
|
||||
completed_at: Optional[str] = None
|
||||
|
||||
|
||||
# ---------- shared resource projections ----------
|
||||
|
||||
class UserOut(_Out):
|
||||
|
||||
@@ -280,6 +280,30 @@ ACTIONS: tuple[Action, ...] = (
|
||||
body("path", "Relative path to delete.", required=True),
|
||||
),
|
||||
),
|
||||
Action(
|
||||
name="zip_project",
|
||||
method="POST",
|
||||
path="/projects/{project_slug}/zip",
|
||||
summary="Build a downloadable zip archive of a whole project",
|
||||
description=(
|
||||
"Queues a background zip job and returns {uid, status_url}. Poll the status_url with "
|
||||
"zip_status until status is 'done', then give the user the download_url."
|
||||
),
|
||||
params=(path("project_slug", "Project slug or uid."),),
|
||||
requires_auth=False,
|
||||
),
|
||||
Action(
|
||||
name="zip_status",
|
||||
method="GET",
|
||||
path="/zips/{uid}",
|
||||
summary="Check a zip job and obtain its download link once finished",
|
||||
description=(
|
||||
"Returns the job status and stats. When status is 'done', download_url points at the "
|
||||
"ready archive; while 'pending' or 'running', poll again shortly."
|
||||
),
|
||||
params=(path("uid", "Zip job uid returned by zip_project."),),
|
||||
requires_auth=False,
|
||||
),
|
||||
Action(
|
||||
name="search_users",
|
||||
method="GET",
|
||||
|
||||
@@ -0,0 +1 @@
|
||||
# retoor <retoor@molodetz.nl>
|
||||
@@ -0,0 +1,246 @@
|
||||
# retoor <retoor@molodetz.nl>
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
import logging
|
||||
from datetime import datetime, timedelta, timezone
|
||||
|
||||
from devplacepy.database import get_table, get_int_setting
|
||||
from devplacepy.services.base import BaseService, ConfigField
|
||||
from devplacepy.services.jobs import queue
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
DEFAULT_RETENTION_SECONDS = 7 * 24 * 60 * 60
|
||||
DEFAULT_MAX_CONCURRENT = 2
|
||||
DEFAULT_JOB_TIMEOUT_SECONDS = 600
|
||||
|
||||
|
||||
class JobService(BaseService):
|
||||
kind = ""
|
||||
min_interval = 1
|
||||
max_retries = 3
|
||||
|
||||
def __init__(self, name: str, interval_seconds: int = 2):
|
||||
super().__init__(name=name, interval_seconds=interval_seconds)
|
||||
self.kind = self.kind or name
|
||||
self.retention_key = f"{name}_retention_seconds"
|
||||
self.concurrency_key = f"{name}_max_concurrent"
|
||||
self.timeout_key = f"{name}_job_timeout_seconds"
|
||||
self.config_fields = [
|
||||
ConfigField(self.retention_key, "Artifact retention (seconds)", type="int",
|
||||
default=DEFAULT_RETENTION_SECONDS, minimum=60,
|
||||
help="Finished artifacts are deleted once unused for this long.",
|
||||
group="Jobs"),
|
||||
ConfigField(self.concurrency_key, "Max concurrent jobs", type="int",
|
||||
default=DEFAULT_MAX_CONCURRENT, minimum=1, maximum=16,
|
||||
help="Number of jobs processed in parallel.",
|
||||
group="Jobs"),
|
||||
ConfigField(self.timeout_key, "Job timeout (seconds)", type="int",
|
||||
default=DEFAULT_JOB_TIMEOUT_SECONDS, minimum=10,
|
||||
help="A running job with no live task after this long is retried.",
|
||||
group="Jobs"),
|
||||
]
|
||||
self._inflight: dict = {}
|
||||
|
||||
def retention_seconds(self) -> int:
|
||||
return max(60, get_int_setting(self.retention_key, DEFAULT_RETENTION_SECONDS))
|
||||
|
||||
def max_concurrent(self) -> int:
|
||||
return max(1, get_int_setting(self.concurrency_key, DEFAULT_MAX_CONCURRENT))
|
||||
|
||||
def job_timeout_seconds(self) -> int:
|
||||
return max(10, get_int_setting(self.timeout_key, DEFAULT_JOB_TIMEOUT_SECONDS))
|
||||
|
||||
async def process(self, job: dict) -> dict:
|
||||
raise NotImplementedError
|
||||
|
||||
def cleanup(self, job: dict) -> None:
|
||||
pass
|
||||
|
||||
async def on_enable(self) -> None:
|
||||
self._recover_orphans(min_age_seconds=0)
|
||||
|
||||
async def on_disable(self) -> None:
|
||||
for entry in list(self._inflight.values()):
|
||||
entry["task"].cancel()
|
||||
self._inflight.clear()
|
||||
|
||||
async def run_once(self) -> None:
|
||||
self._reap()
|
||||
self._recover_orphans(min_age_seconds=self.job_timeout_seconds())
|
||||
self._refill()
|
||||
self._sweep_expired()
|
||||
|
||||
def _reap(self) -> None:
|
||||
for uid in list(self._inflight):
|
||||
entry = self._inflight[uid]
|
||||
task = entry["task"]
|
||||
if not task.done():
|
||||
continue
|
||||
del self._inflight[uid]
|
||||
duration_ms = int((datetime.now(timezone.utc) - entry["started"]).total_seconds() * 1000)
|
||||
exc = task.exception() if not task.cancelled() else asyncio.CancelledError()
|
||||
if exc is not None:
|
||||
self._finish_failed(uid, str(exc) or exc.__class__.__name__, duration_ms)
|
||||
else:
|
||||
self._finish_done(uid, task.result() or {}, duration_ms)
|
||||
|
||||
def _finish_done(self, uid: str, result_data: dict, duration_ms: int) -> None:
|
||||
now = datetime.now(timezone.utc)
|
||||
get_table("jobs").update({
|
||||
"uid": uid,
|
||||
"status": queue.DONE,
|
||||
"result": json.dumps(result_data),
|
||||
"error": "",
|
||||
"completed_at": now.isoformat(),
|
||||
"updated_at": now.isoformat(),
|
||||
"duration_ms": duration_ms,
|
||||
"last_accessed_at": now.isoformat(),
|
||||
"expires_at": (now + timedelta(seconds=self.retention_seconds())).isoformat(),
|
||||
"bytes_in": int(result_data.get("bytes_in", 0)),
|
||||
"bytes_out": int(result_data.get("bytes_out", 0)),
|
||||
"item_count": int(result_data.get("item_count", 0)),
|
||||
}, ["uid"])
|
||||
self.log(f"Job {uid} done in {duration_ms}ms")
|
||||
|
||||
def _finish_failed(self, uid: str, error: str, duration_ms: int) -> None:
|
||||
now = datetime.now(timezone.utc)
|
||||
get_table("jobs").update({
|
||||
"uid": uid,
|
||||
"status": queue.FAILED,
|
||||
"error": error[:2000],
|
||||
"completed_at": now.isoformat(),
|
||||
"updated_at": now.isoformat(),
|
||||
"duration_ms": duration_ms,
|
||||
}, ["uid"])
|
||||
self.log(f"Job {uid} failed: {error}")
|
||||
|
||||
def _refill(self) -> None:
|
||||
capacity = self.max_concurrent() - len(self._inflight)
|
||||
if capacity <= 0:
|
||||
return
|
||||
table = get_table("jobs")
|
||||
pending = list(table.find(kind=self.kind, status=queue.PENDING, order_by=["uid"], _limit=capacity))
|
||||
for row in pending:
|
||||
uid = row["uid"]
|
||||
if uid in self._inflight:
|
||||
continue
|
||||
now = datetime.now(timezone.utc)
|
||||
table.update({
|
||||
"uid": uid,
|
||||
"status": queue.RUNNING,
|
||||
"started_at": now.isoformat(),
|
||||
"updated_at": now.isoformat(),
|
||||
"error": "",
|
||||
}, ["uid"])
|
||||
job = queue.get_job(uid)
|
||||
self._inflight[uid] = {"task": asyncio.create_task(self._run_job(job)), "started": now}
|
||||
self.log(f"Job {uid} started")
|
||||
|
||||
async def _run_job(self, job: dict) -> dict:
|
||||
return await self.process(job)
|
||||
|
||||
def _recover_orphans(self, min_age_seconds: int) -> None:
|
||||
table = get_table("jobs")
|
||||
now = datetime.now(timezone.utc)
|
||||
for row in table.find(kind=self.kind, status=queue.RUNNING):
|
||||
uid = row["uid"]
|
||||
if uid in self._inflight:
|
||||
continue
|
||||
if min_age_seconds > 0 and not self._older_than(row.get("started_at"), now, min_age_seconds):
|
||||
continue
|
||||
retry_count = int(row.get("retry_count") or 0)
|
||||
if retry_count >= self.max_retries:
|
||||
self._finish_failed(uid, "exceeded retry limit after orphan recovery", 0)
|
||||
continue
|
||||
table.update({
|
||||
"uid": uid,
|
||||
"status": queue.PENDING,
|
||||
"retry_count": retry_count + 1,
|
||||
"started_at": "",
|
||||
"updated_at": now.isoformat(),
|
||||
}, ["uid"])
|
||||
self.log(f"Recovered orphaned job {uid} (retry {retry_count + 1})")
|
||||
|
||||
def _sweep_expired(self) -> None:
|
||||
table = get_table("jobs")
|
||||
now = datetime.now(timezone.utc)
|
||||
for row in table.find(kind=self.kind, status=queue.DONE):
|
||||
expires_at = row.get("expires_at")
|
||||
if not expires_at or not self._is_past(expires_at, now):
|
||||
continue
|
||||
job = queue.get_job(row["uid"])
|
||||
try:
|
||||
self.cleanup(job)
|
||||
except Exception as exc:
|
||||
self.log(f"Cleanup failed for {row['uid']}: {exc}")
|
||||
table.delete(uid=row["uid"])
|
||||
self.log(f"Pruned expired job {row['uid']}")
|
||||
|
||||
def _older_than(self, stamp: str, now: datetime, seconds: int) -> bool:
|
||||
moment = self._parse(stamp)
|
||||
if moment is None:
|
||||
return True
|
||||
return (now - moment).total_seconds() >= seconds
|
||||
|
||||
def _is_past(self, stamp: str, now: datetime) -> bool:
|
||||
moment = self._parse(stamp)
|
||||
return moment is not None and moment < now
|
||||
|
||||
def _parse(self, stamp: str):
|
||||
if not stamp:
|
||||
return None
|
||||
try:
|
||||
return datetime.fromisoformat(stamp)
|
||||
except (ValueError, TypeError):
|
||||
return None
|
||||
|
||||
def collect_metrics(self) -> dict:
|
||||
rows = queue.list_jobs(kind=self.kind)
|
||||
by_status = {queue.PENDING: 0, queue.RUNNING: 0, queue.DONE: 0, queue.FAILED: 0}
|
||||
bytes_in = bytes_out = items = 0
|
||||
durations = []
|
||||
for row in rows:
|
||||
by_status[row.get("status", "")] = by_status.get(row.get("status", ""), 0) + 1
|
||||
bytes_in += int(row.get("bytes_in") or 0)
|
||||
bytes_out += int(row.get("bytes_out") or 0)
|
||||
items += int(row.get("item_count") or 0)
|
||||
if row.get("status") == queue.DONE and row.get("duration_ms"):
|
||||
durations.append(int(row["duration_ms"]))
|
||||
avg_ms = int(sum(durations) / len(durations)) if durations else 0
|
||||
stats = [
|
||||
{"label": "Total jobs", "value": len(rows)},
|
||||
{"label": "Pending", "value": by_status[queue.PENDING]},
|
||||
{"label": "Running", "value": by_status[queue.RUNNING]},
|
||||
{"label": "Done", "value": by_status[queue.DONE]},
|
||||
{"label": "Failed", "value": by_status[queue.FAILED]},
|
||||
{"label": "Items zipped", "value": items},
|
||||
{"label": "Bytes in", "value": _human_bytes(bytes_in)},
|
||||
{"label": "Bytes out", "value": _human_bytes(bytes_out)},
|
||||
{"label": "Avg duration", "value": f"{avg_ms} ms"},
|
||||
]
|
||||
recent = sorted(rows, key=lambda r: r.get("uid", ""), reverse=True)[:10]
|
||||
table = {
|
||||
"columns": ["Job", "Status", "Name", "Out", "Items"],
|
||||
"rows": [
|
||||
[
|
||||
row.get("uid", "")[:8],
|
||||
row.get("status", ""),
|
||||
(row.get("preferred_name") or "")[:32],
|
||||
_human_bytes(int(row.get("bytes_out") or 0)),
|
||||
int(row.get("item_count") or 0),
|
||||
]
|
||||
for row in recent
|
||||
],
|
||||
}
|
||||
return {"stats": stats, "table": table}
|
||||
|
||||
|
||||
def _human_bytes(size: int) -> str:
|
||||
value = float(size)
|
||||
for unit in ("B", "KB", "MB", "GB", "TB"):
|
||||
if value < 1024 or unit == "TB":
|
||||
return f"{value:.0f} {unit}" if unit == "B" else f"{value:.1f} {unit}"
|
||||
value /= 1024
|
||||
return f"{value:.1f} TB"
|
||||
@@ -0,0 +1,103 @@
|
||||
# retoor <retoor@molodetz.nl>
|
||||
|
||||
import json
|
||||
import logging
|
||||
from datetime import datetime, timedelta, timezone
|
||||
|
||||
from devplacepy.database import db, get_table
|
||||
from devplacepy.utils import generate_uid
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
PENDING = "pending"
|
||||
RUNNING = "running"
|
||||
DONE = "done"
|
||||
FAILED = "failed"
|
||||
|
||||
|
||||
def _now() -> str:
|
||||
return datetime.now(timezone.utc).isoformat()
|
||||
|
||||
|
||||
def _table():
|
||||
return get_table("jobs")
|
||||
|
||||
|
||||
def _decode(field: str) -> dict:
|
||||
if not field:
|
||||
return {}
|
||||
try:
|
||||
return json.loads(field)
|
||||
except (ValueError, TypeError):
|
||||
return {}
|
||||
|
||||
|
||||
def _hydrate(row: dict) -> dict:
|
||||
row = dict(row)
|
||||
row["payload"] = _decode(row.get("payload"))
|
||||
row["result"] = _decode(row.get("result"))
|
||||
return row
|
||||
|
||||
|
||||
def enqueue(kind: str, payload: dict, owner_kind: str, owner_id: str, preferred_name: str = "") -> str:
|
||||
uid = generate_uid()
|
||||
now = _now()
|
||||
_table().insert({
|
||||
"uid": uid,
|
||||
"kind": kind,
|
||||
"status": PENDING,
|
||||
"owner_kind": owner_kind,
|
||||
"owner_id": owner_id,
|
||||
"preferred_name": preferred_name,
|
||||
"payload": json.dumps(payload),
|
||||
"result": "",
|
||||
"error": "",
|
||||
"retry_count": 0,
|
||||
"created_at": now,
|
||||
"started_at": "",
|
||||
"completed_at": "",
|
||||
"updated_at": now,
|
||||
"duration_ms": 0,
|
||||
"last_accessed_at": "",
|
||||
"expires_at": "",
|
||||
"bytes_in": 0,
|
||||
"bytes_out": 0,
|
||||
"item_count": 0,
|
||||
})
|
||||
logger.info("Enqueued %s job %s for %s/%s", kind, uid, owner_kind, owner_id)
|
||||
return uid
|
||||
|
||||
|
||||
def get_job(uid: str) -> dict | None:
|
||||
if "jobs" not in db.tables:
|
||||
return None
|
||||
row = _table().find_one(uid=uid)
|
||||
return _hydrate(row) if row else None
|
||||
|
||||
|
||||
def touch_job(uid: str, extend_seconds: int) -> None:
|
||||
if "jobs" not in db.tables:
|
||||
return
|
||||
row = _table().find_one(uid=uid)
|
||||
if not row:
|
||||
return
|
||||
now = datetime.now(timezone.utc)
|
||||
_table().update({
|
||||
"uid": uid,
|
||||
"last_accessed_at": now.isoformat(),
|
||||
"expires_at": (now + timedelta(seconds=extend_seconds)).isoformat(),
|
||||
"updated_at": now.isoformat(),
|
||||
}, ["uid"])
|
||||
|
||||
|
||||
def list_jobs(kind: str = None, status: str = None, owner: tuple = None) -> list[dict]:
|
||||
if "jobs" not in db.tables:
|
||||
return []
|
||||
filters = {}
|
||||
if kind:
|
||||
filters["kind"] = kind
|
||||
if status:
|
||||
filters["status"] = status
|
||||
if owner:
|
||||
filters["owner_kind"], filters["owner_id"] = owner
|
||||
return [_hydrate(row) for row in _table().find(**filters)]
|
||||
@@ -0,0 +1,96 @@
|
||||
# retoor <retoor@molodetz.nl>
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import shutil
|
||||
import sys
|
||||
from pathlib import Path
|
||||
|
||||
from devplacepy.config import BASE_DIR, STATIC_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"
|
||||
WORKER_MODULE = "devplacepy.services.jobs.zip_worker"
|
||||
|
||||
|
||||
def _shard(uid: str) -> str:
|
||||
return f"{uid[:2]}/{uid[2:4]}"
|
||||
|
||||
|
||||
class ZipService(JobService):
|
||||
kind = "zip"
|
||||
title = "Zip"
|
||||
description = (
|
||||
"Builds downloadable zip archives in a subprocess off the request path, tracks each "
|
||||
"job with full statistics, and deletes archives once they go unused for the retention "
|
||||
"window."
|
||||
)
|
||||
|
||||
def __init__(self):
|
||||
super().__init__(name="zip", interval_seconds=2)
|
||||
|
||||
async def process(self, job: dict) -> dict:
|
||||
source = job["payload"].get("source", {})
|
||||
uid = job["uid"]
|
||||
staging = STAGING_DIR / uid
|
||||
try:
|
||||
item_count = await asyncio.to_thread(self._materialize, source, staging)
|
||||
tmp_dir = ZIPS_DIR / _shard(uid)
|
||||
tmp_dir.mkdir(parents=True, exist_ok=True)
|
||||
tmp_zip = tmp_dir / f"{generate_uid()}.zip"
|
||||
stats = await self._run_worker(staging, tmp_zip)
|
||||
final_name = self._final_name(job.get("preferred_name", ""), stats["crc32"])
|
||||
final_path = tmp_dir / final_name
|
||||
os.replace(tmp_zip, final_path)
|
||||
finally:
|
||||
await asyncio.to_thread(shutil.rmtree, staging, ignore_errors=True)
|
||||
|
||||
return {
|
||||
"download_url": f"/zips/{uid}/download",
|
||||
"local_path": str(final_path),
|
||||
"final_name": final_name,
|
||||
"crc32": f"{stats['crc32']:08x}",
|
||||
"file_count": stats["file_count"],
|
||||
"dir_count": stats["dir_count"],
|
||||
"bytes_in": stats["bytes_in"],
|
||||
"bytes_out": stats["bytes_out"],
|
||||
"item_count": item_count,
|
||||
}
|
||||
|
||||
def _materialize(self, source: dict, staging: Path) -> int:
|
||||
if source.get("type") == "project_tree":
|
||||
return project_files.export_to_dir(source["project_uid"], source.get("path", ""), staging)
|
||||
raise ValueError(f"unsupported zip source: {source.get('type')}")
|
||||
|
||||
async def _run_worker(self, source_dir: Path, output_path: Path) -> dict:
|
||||
proc = await asyncio.create_subprocess_exec(
|
||||
sys.executable, "-m", WORKER_MODULE, str(source_dir), str(output_path),
|
||||
cwd=str(BASE_DIR),
|
||||
stdout=asyncio.subprocess.PIPE,
|
||||
stderr=asyncio.subprocess.PIPE,
|
||||
)
|
||||
out, err = await proc.communicate()
|
||||
if proc.returncode != 0:
|
||||
raise RuntimeError(f"zip worker exited {proc.returncode}: {err.decode('utf-8', 'replace')[:500]}")
|
||||
return json.loads(out.decode("utf-8"))
|
||||
|
||||
def _final_name(self, preferred_name: str, crc32: int) -> str:
|
||||
name = preferred_name or "download"
|
||||
if name.lower().endswith(".zip"):
|
||||
name = name[:-4]
|
||||
slug = slugify(name) or "download"
|
||||
return f"{crc32:08x}.{slug}.zip"
|
||||
|
||||
def cleanup(self, job: dict) -> None:
|
||||
result = job.get("result", {})
|
||||
local_path = result.get("local_path")
|
||||
if local_path:
|
||||
Path(local_path).unlink(missing_ok=True)
|
||||
shutil.rmtree(STAGING_DIR / job["uid"], ignore_errors=True)
|
||||
@@ -0,0 +1,65 @@
|
||||
# retoor <retoor@molodetz.nl>
|
||||
|
||||
import json
|
||||
import os
|
||||
import shutil
|
||||
import sys
|
||||
import zipfile
|
||||
import zlib
|
||||
|
||||
FIXED_DATE = (1980, 1, 1, 0, 0, 0)
|
||||
|
||||
|
||||
def _build(source_dir: str, output_path: str) -> dict:
|
||||
bytes_in = 0
|
||||
file_count = 0
|
||||
dir_count = 0
|
||||
with zipfile.ZipFile(output_path, "w", zipfile.ZIP_DEFLATED) as archive:
|
||||
for root, dirs, files in os.walk(source_dir):
|
||||
dirs.sort()
|
||||
for name in dirs:
|
||||
dir_count += 1
|
||||
arcname = os.path.relpath(os.path.join(root, name), source_dir) + "/"
|
||||
info = zipfile.ZipInfo(arcname, date_time=FIXED_DATE)
|
||||
info.external_attr = (0o040755 << 16) | 0x10
|
||||
archive.writestr(info, b"")
|
||||
for name in sorted(files):
|
||||
full = os.path.join(root, name)
|
||||
arcname = os.path.relpath(full, source_dir)
|
||||
bytes_in += os.path.getsize(full)
|
||||
info = zipfile.ZipInfo(arcname, date_time=FIXED_DATE)
|
||||
info.compress_type = zipfile.ZIP_DEFLATED
|
||||
info.external_attr = 0o644 << 16
|
||||
with open(full, "rb") as source, archive.open(info, "w") as target:
|
||||
shutil.copyfileobj(source, target)
|
||||
file_count += 1
|
||||
|
||||
crc = 0
|
||||
with open(output_path, "rb") as handle:
|
||||
for chunk in iter(lambda: handle.read(1024 * 1024), b""):
|
||||
crc = zlib.crc32(chunk, crc)
|
||||
|
||||
return {
|
||||
"bytes_in": bytes_in,
|
||||
"bytes_out": os.path.getsize(output_path),
|
||||
"file_count": file_count,
|
||||
"dir_count": dir_count,
|
||||
"crc32": crc & 0xFFFFFFFF,
|
||||
}
|
||||
|
||||
|
||||
def main(argv: list) -> int:
|
||||
if len(argv) != 3:
|
||||
sys.stderr.write("usage: zip_worker <source_dir> <output_zip>\n")
|
||||
return 2
|
||||
source_dir, output_path = argv[1], argv[2]
|
||||
if not os.path.isdir(source_dir):
|
||||
sys.stderr.write(f"source directory not found: {source_dir}\n")
|
||||
return 3
|
||||
stats = _build(source_dir, output_path)
|
||||
sys.stdout.write(json.dumps(stats))
|
||||
return 0
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
raise SystemExit(main(sys.argv))
|
||||
@@ -19,6 +19,7 @@ import { ApiKeyManager } from "./ApiKeyManager.js";
|
||||
import { UserAiUsage } from "./UserAiUsage.js";
|
||||
import { Tabs } from "./Tabs.js";
|
||||
import { DeviiTerminal } from "./DeviiTerminal.js";
|
||||
import { ZipDownloader } from "./ZipDownloader.js";
|
||||
|
||||
class Application {
|
||||
constructor() {
|
||||
@@ -46,6 +47,7 @@ class Application {
|
||||
this.userAiUsage = new UserAiUsage();
|
||||
this.tabs = new Tabs();
|
||||
this.devii = new DeviiTerminal();
|
||||
this.zipDownloader = new ZipDownloader();
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -358,6 +358,7 @@ export class ProjectFiles {
|
||||
items.push({ label: "Upload here", icon: "⬆", onSelect: () => this.uploadTo(dir) });
|
||||
}
|
||||
items.push({ separator: true });
|
||||
items.push({ label: "Download as zip", icon: "📦", onSelect: () => window.app.zipDownloader.start(`${this.base}/zip?path=${encodeURIComponent(path)}`) });
|
||||
items.push({ label: "Rename / move", icon: "✏", onSelect: () => this.rename(path) });
|
||||
items.push({ label: "Delete", icon: "🗑", danger: true, onSelect: () => this.remove(path) });
|
||||
return items;
|
||||
|
||||
@@ -0,0 +1,71 @@
|
||||
// retoor <retoor@molodetz.nl>
|
||||
|
||||
import { Http } from "./Http.js";
|
||||
|
||||
export class ZipDownloader {
|
||||
constructor() {
|
||||
this.pollMs = 1500;
|
||||
this.maxPolls = 200;
|
||||
this.active = new Set();
|
||||
document.addEventListener("click", (event) => {
|
||||
const trigger = event.target.closest("[data-zip-download]");
|
||||
if (!trigger) return;
|
||||
event.preventDefault();
|
||||
this.start(trigger.dataset.zipDownload);
|
||||
});
|
||||
}
|
||||
|
||||
async start(enqueueUrl) {
|
||||
if (this.active.has(enqueueUrl)) return;
|
||||
this.active.add(enqueueUrl);
|
||||
window.app.toast.show("Preparing download...", { type: "info", ms: 4000 });
|
||||
try {
|
||||
const job = await Http.postJson(enqueueUrl, {});
|
||||
await this.poll(job.status_url);
|
||||
} catch (err) {
|
||||
window.app.toast.show("Could not start the download", { type: "error" });
|
||||
} finally {
|
||||
this.active.delete(enqueueUrl);
|
||||
}
|
||||
}
|
||||
|
||||
poll(statusUrl) {
|
||||
return new Promise((resolve) => {
|
||||
let attempts = 0;
|
||||
const timer = setInterval(async () => {
|
||||
attempts += 1;
|
||||
if (attempts > this.maxPolls) {
|
||||
clearInterval(timer);
|
||||
window.app.toast.show("Download timed out", { type: "error" });
|
||||
resolve();
|
||||
return;
|
||||
}
|
||||
let status;
|
||||
try {
|
||||
status = await Http.getJson(statusUrl);
|
||||
} catch {
|
||||
return;
|
||||
}
|
||||
if (status.status === "done") {
|
||||
clearInterval(timer);
|
||||
this.trigger(status.download_url);
|
||||
window.app.toast.show("Download ready", { type: "success" });
|
||||
resolve();
|
||||
} else if (status.status === "failed") {
|
||||
clearInterval(timer);
|
||||
window.app.toast.show("Zip failed: " + (status.error || "unknown error"), { type: "error" });
|
||||
resolve();
|
||||
}
|
||||
}, this.pollMs);
|
||||
});
|
||||
}
|
||||
|
||||
trigger(url) {
|
||||
const link = document.createElement("a");
|
||||
link.href = url;
|
||||
link.download = "";
|
||||
document.body.appendChild(link);
|
||||
link.click();
|
||||
link.remove();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,62 @@
|
||||
<div class="docs-content" data-render>
|
||||
# Async job services
|
||||
|
||||
The async job framework (`services/jobs/`) is the standard way to run heavy, blocking work off the request path and hand the caller a result URL when it is finished. It builds on the [service framework](/docs/services-framework.html): a `JobService` is a `BaseService` whose `run_once` drives a queue of database-backed jobs. [ZipService](/docs/services-zip.html) is the first consumer.
|
||||
|
||||
> Audience: administrators and maintainers.
|
||||
|
||||
## When to use it
|
||||
|
||||
Use a job service whenever a request would otherwise block on slow work: archiving, exporting, rendering, batch processing, or anything that should survive a page reload and be polled. The request enqueues a row and returns immediately; a background worker does the work; the client polls for completion.
|
||||
|
||||
## The jobs table
|
||||
|
||||
All job kinds share one `jobs` table, discriminated by a `kind` column. The row is the queue. Columns: `uid` (a sortable uuid7, also the job id), `kind`, `status` (`pending`, `running`, `done`, `failed`), `owner_kind`/`owner_id` (attribution only), `preferred_name`, `payload` and `result` (JSON strings), `error`, `retry_count`, the timestamps `created_at`/`started_at`/`completed_at`/`updated_at`, `duration_ms`, `last_accessed_at`/`expires_at` (retention), and the common stat columns `bytes_in`/`bytes_out`/`item_count`. Kind-specific fields live inside the `result` JSON. Indexes cover `(kind, status)`, `(owner_kind, owner_id)`, and `expires_at`.
|
||||
|
||||
## Enqueue from any worker
|
||||
|
||||
`services/jobs/queue.py` is pure database access and is safe to call from any worker's request handler:
|
||||
|
||||
- `enqueue(kind, payload, owner_kind, owner_id, preferred_name)` inserts a `pending` row and returns its uid.
|
||||
- `get_job(uid)` returns the hydrated job (payload and result decoded).
|
||||
- `touch_job(uid, extend_seconds)` bumps `last_accessed_at` and extends `expires_at`; the download path calls this so an actively-used artifact is not pruned.
|
||||
- `list_jobs(kind, status, owner)` backs metrics and the CLI.
|
||||
|
||||
## Single-owner processing
|
||||
|
||||
Processing happens only in the worker that owns the service lock, because only that worker calls `supervise()` and therefore only it runs `run_once()`. There is exactly one processor across the deployment, so no cross-worker locking or atomic claim is needed. Status reads and downloads work from any worker because they read the shared database and the shared `static` filesystem.
|
||||
|
||||
`JobService.run_once` does four things each tick:
|
||||
|
||||
1. **Reap** finished in-flight tasks, writing `result`, timing, and final status.
|
||||
2. **Recover orphans**: any `running` row not in this process's in-flight map (necessarily left by a previous owner) is reset to `pending` with a bounded `retry_count`, or marked `failed` past the limit.
|
||||
3. **Refill**: claim the oldest `pending` jobs up to the concurrency limit (uuid7 sorts in creation order, giving FIFO), mark them `running`, and spawn an `asyncio` task per job.
|
||||
4. **Sweep**: delete jobs whose `expires_at` has passed, calling the subclass `cleanup` hook to remove the artifact first.
|
||||
|
||||
In-flight tasks span ticks; they are stored in an instance map and reaped when done, so the loop returns quickly each tick. A per-tick timeout guard recovers a `running` row whose task vanished.
|
||||
|
||||
## Writing a new kind
|
||||
|
||||
Subclass `JobService`, set `kind`, and implement two methods:
|
||||
|
||||
```python
|
||||
class MyJobService(JobService):
|
||||
kind = "mykind"
|
||||
title = "My jobs"
|
||||
|
||||
async def process(self, job):
|
||||
# do the work; return a dict stored as result JSON.
|
||||
# include bytes_in / bytes_out / item_count for the shared stats.
|
||||
return {"download_url": f"/.../{job['uid']}/download", "local_path": "...", "item_count": n}
|
||||
|
||||
def cleanup(self, job):
|
||||
# remove the artifact named in job["result"].
|
||||
...
|
||||
```
|
||||
|
||||
Register it in `main.py` startup with `service_manager.register(MyJobService())`. It inherits the retention, concurrency, and timeout config fields (keyed `{name}_retention_seconds`, `{name}_max_concurrent`, `{name}_job_timeout_seconds`), the aggregate metrics card, crash recovery, and automatic pruning. Add enqueue endpoints that own their authorization, a status route, and a download route, then document the kind under the Services section.
|
||||
|
||||
## Retention is built in
|
||||
|
||||
There is no separate reaper service. Each job service prunes its own expired artifacts every tick through its `cleanup` hook, so deletion is a standard capability every kind gets for free. Retention is admin-configurable per service and defaults to seven days; downloading an artifact extends its expiry.
|
||||
</div>
|
||||
@@ -0,0 +1,51 @@
|
||||
<div class="docs-content" data-render>
|
||||
# ZipService
|
||||
|
||||
`ZipService` (`services/jobs/zip_service.py`) builds downloadable zip archives off the request path and hands the caller a URL when the archive is ready. It is the first consumer of the generic [async job framework](/docs/architecture-jobs.html); read that page for the lifecycle, single-owner model, and retention that this service inherits.
|
||||
|
||||
> Audience: administrators and maintainers.
|
||||
|
||||
## What it does
|
||||
|
||||
A caller enqueues a job describing what to archive. The lock-owning worker materializes the source onto a temporary directory, runs a standalone Python subprocess to compress it, computes a CRC32 over the output, names the file from that checksum, and records full statistics. Callers poll the job until it is `done`, then download the archive. Archives that go unused for the retention window are deleted automatically.
|
||||
|
||||
## Source contract
|
||||
|
||||
A zip job payload carries a `source` object with a `type`. The only implemented type is `project_tree`:
|
||||
|
||||
```
|
||||
{"source": {"type": "project_tree", "project_uid": "...", "path": "src"}}
|
||||
```
|
||||
|
||||
An empty `path` archives the whole project; a non-empty `path` archives that file or directory subtree. `ZipService._materialize` dispatches on the type and calls `project_files.export_to_dir`, which reconstructs the virtual project filesystem (inline text rows plus on-disk binary blobs) onto a real staging directory under `static/uploads/zip_staging/{job_uid}`, with a traversal guard on every written path. Materialization runs inside the worker via `asyncio.to_thread`, never in the request handler.
|
||||
|
||||
## Subprocess compression
|
||||
|
||||
`ZipService._run_worker` spawns `python -m devplacepy.services.jobs.zip_worker <staging_dir> <tmp_zip>` with `asyncio.create_subprocess_exec`. The worker is stdlib-only (no system `zip` dependency), walks the tree (sorted) into a `ZIP_DEFLATED` archive with a fixed entry timestamp so the output is **deterministic** (identical content yields identical bytes), computes a streaming `zlib.crc32` over the result, and prints a JSON stats line: `bytes_in`, `bytes_out`, `file_count`, `dir_count`, `crc32`. The service parses that JSON; a non-zero exit raises and the job is marked `failed`. Because the archive is deterministic, the CRC32 acts as a content identity: re-archiving the same content produces the same filename, and the final write simply overwrites the previous archive.
|
||||
|
||||
## Naming and storage
|
||||
|
||||
The final name is `{crc32}.{slug}.zip`, where the slug comes from the caller's preferred name with any trailing `.zip` stripped and the rest slugified. The archive lives under `static/uploads/zips/{uid[:2]}/{uid[2:4]}/`. The temporary `{uuid4}.zip` is renamed into place with `os.replace`, so re-running a job simply overwrites the previous archive. The staging directory is removed afterwards.
|
||||
|
||||
## Endpoints
|
||||
|
||||
- `POST /projects/{project_slug}/zip` queues a whole-project archive. Authorization is checked here against the project; the response is `{uid, status_url}`.
|
||||
- `POST /projects/{project_slug}/files/zip?path=...` queues a subtree archive (empty `path` means the whole project).
|
||||
- `GET /zips/{uid}` returns the job status and statistics (`ZipJobOut`). `download_url` is null until the job is `done`.
|
||||
- `GET /zips/{uid}/download` streams the finished archive as `application/zip` and extends the retention window on each access.
|
||||
|
||||
The job uid (a sortable uuid7) is the capability token: anyone with the link can download, which matches the fact that projects are publicly viewable. The owner recorded on the job is used only for attribution and metrics, never for access control.
|
||||
|
||||
## Frontend
|
||||
|
||||
`static/js/ZipDownloader.js` (instantiated as `app.zipDownloader`) drives the whole flow: it POSTs to the enqueue URL, polls `GET /zips/{uid}` on an interval, and triggers a hidden download link once the job is `done`, with toasts for progress and errors. Any element with a `data-zip-download` attribute is wired automatically (the "Download zip" button on the project detail page and the "Download as zip" button on the files page); the files context menu calls `app.zipDownloader.start(url)` directly for a single file or folder.
|
||||
|
||||
## Configuration, metrics, and quotas
|
||||
|
||||
On `/admin/services` the service exposes the inherited job fields: artifact retention (default 7 days), maximum concurrent jobs (default 2), and job timeout. Its metrics card shows totals per status, total items zipped, bytes in and out, average duration, and a table of recent jobs.
|
||||
|
||||
## CLI
|
||||
|
||||
- `devplace zips prune` deletes expired archives and their job rows.
|
||||
- `devplace zips clear` deletes every zip archive and job row.
|
||||
</div>
|
||||
@@ -54,6 +54,7 @@
|
||||
|
||||
<div class="project-detail-actions">
|
||||
<a href="/projects/{{ project['slug'] or project['uid'] }}/files" class="project-star-btn">📁 Files</a>
|
||||
<button type="button" class="project-star-btn" data-zip-download="/projects/{{ project['slug'] or project['uid'] }}/zip">📦 Download zip</button>
|
||||
<button type="button" class="project-star-btn" data-share="/projects/{{ project['slug'] or project['uid'] }}">🔗 Share</button>
|
||||
{% if user %}
|
||||
{% set _type = "project" %}{% set _uid = project['uid'] %}{% set _my_vote = my_vote %}{% set _count = star_count %}{% set _btn_class = "project-star-btn" %}{% include "_star_vote.html" %}
|
||||
|
||||
@@ -17,6 +17,7 @@
|
||||
<button type="button" class="btn btn-secondary" id="pf-delete">Delete</button>
|
||||
<button type="button" class="btn btn-primary" id="pf-save">Save</button>
|
||||
{% endif %}
|
||||
<button type="button" class="btn btn-secondary" data-zip-download="/projects/{{ project['slug'] or project['uid'] }}/files/zip">Download as zip</button>
|
||||
<span class="pf-status" id="pf-status"></span>
|
||||
</div>
|
||||
</div>
|
||||
|
||||
Reference in New Issue
Block a user