|
# 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)]
|