90 lines
3.0 KiB
Python
90 lines
3.0 KiB
Python
"""Хранилище фоновых заданий (создание кластера) в памяти процесса.
|
|
|
|
При перезапуске контейнера история заданий обнуляется — это ожидаемо для dev-среды.
|
|
|
|
Автор: Сергей Антропов
|
|
Сайт: https://devops.org.ru
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import logging
|
|
import uuid
|
|
from dataclasses import dataclass, field
|
|
from datetime import datetime, timezone
|
|
from typing import Any, Literal
|
|
|
|
logger = logging.getLogger("kind_k8s.job_store")
|
|
|
|
JobStatus = Literal["queued", "running", "success", "failed"]
|
|
|
|
|
|
@dataclass
|
|
class JobRecord:
|
|
"""Описание одного задания."""
|
|
|
|
job_id: str
|
|
kind: str
|
|
status: JobStatus
|
|
cluster_name: str | None
|
|
created_at_utc: str
|
|
message: str | None = None
|
|
result: dict[str, Any] | None = None
|
|
|
|
|
|
class JobStore:
|
|
"""Потокобезопасное (asyncio) хранилище заданий."""
|
|
|
|
def __init__(self) -> None:
|
|
self._jobs: dict[str, JobRecord] = {}
|
|
self._lock = asyncio.Lock()
|
|
|
|
async def create_job(self, kind: str, *, cluster_name: str | None) -> JobRecord:
|
|
"""Зарегистрировать задание в статусе ``queued``."""
|
|
jid = uuid.uuid4().hex
|
|
now = datetime.now(timezone.utc).isoformat()
|
|
rec = JobRecord(
|
|
job_id=jid,
|
|
kind=kind,
|
|
status="queued",
|
|
cluster_name=cluster_name,
|
|
created_at_utc=now,
|
|
)
|
|
async with self._lock:
|
|
self._jobs[jid] = rec
|
|
logger.info("Создано задание %s kind=%s cluster=%s", jid, kind, cluster_name)
|
|
return rec
|
|
|
|
async def set_running(self, job_id: str) -> None:
|
|
async with self._lock:
|
|
if job_id in self._jobs:
|
|
self._jobs[job_id].status = "running"
|
|
self._jobs[job_id].message = None
|
|
|
|
async def set_success(self, job_id: str, *, result: dict[str, Any] | None = None, message: str | None = None) -> None:
|
|
async with self._lock:
|
|
if job_id in self._jobs:
|
|
self._jobs[job_id].status = "success"
|
|
self._jobs[job_id].result = result
|
|
self._jobs[job_id].message = message
|
|
|
|
async def set_failed(self, job_id: str, message: str) -> None:
|
|
async with self._lock:
|
|
if job_id in self._jobs:
|
|
self._jobs[job_id].status = "failed"
|
|
self._jobs[job_id].message = message
|
|
logger.warning("Задание %s завершилось ошибкой: %s", job_id, message)
|
|
|
|
async def get(self, job_id: str) -> JobRecord | None:
|
|
async with self._lock:
|
|
return self._jobs.get(job_id)
|
|
|
|
def snapshot_all(self) -> list[JobRecord]:
|
|
"""Снимок всех заданий (для отладки; без блокировки — eventual consistency)."""
|
|
return list(self._jobs.values())
|
|
|
|
|
|
# Синглтон на процесс uvicorn
|
|
job_store = JobStore()
|