chainshift-dashboard / core /job_realtime.py
GitHub Action
Sync from GitHub
ef78361
Raw
History Blame Contribute Delete
7.72 kB
"""Sentiment Analysis Job ์‹ค์‹œ๊ฐ„ ๋ชจ๋‹ˆํ„ฐ๋ง (Phase 4.1)
Usage:
from job_realtime import get_active_jobs, get_job_by_id
# ์‹คํ–‰ ์ค‘์ธ Job ๋ชฉ๋ก
jobs = get_active_jobs(campaign_id=27)
# ํŠน์ • Job ์ƒ์„ธ
job = get_job_by_id("abc-123")
Note:
Streamlit์€ WebSocket์„ ์ง์ ‘ ์ง€์›ํ•˜์ง€ ์•Š์œผ๋ฏ€๋กœ,
st.rerun() ๋˜๋Š” st.cache_data(ttl=5)๋กœ ํด๋ง ๋ฐฉ์‹ ์‚ฌ์šฉ
"""
import os
from datetime import datetime, timezone
from functools import lru_cache
from pathlib import Path
from typing import Literal
import streamlit as st
from supabase import create_client, Client
from dotenv import load_dotenv
# Load environment variables (project root)
_project_root = Path(__file__).parent.parent.parent.parent
for _env_name in (".env.dev", ".env.prod"):
_env_path = _project_root / _env_name
if _env_path.exists():
load_dotenv(_env_path, override=True)
break
@lru_cache()
def _get_client() -> Client:
"""Supabase ํด๋ผ์ด์–ธํŠธ (์บ์‹œ๋จ)"""
url = os.environ.get("SUPABASE_URL", "")
key = os.environ.get("SUPABASE_SERVICE_KEY", "")
if not url or not key:
raise ValueError("Missing SUPABASE_URL or SUPABASE_SERVICE_KEY")
return create_client(url, key)
JobStatus = Literal["queued", "running", "completed", "failed", "cancelled"]
@st.cache_data(ttl=5)
def get_active_jobs(campaign_id: int | None = None, limit: int = 10) -> list[dict]:
"""์‹คํ–‰ ์ค‘์ด๊ฑฐ๋‚˜ ๋Œ€๊ธฐ ์ค‘์ธ Job ๋ชฉ๋ก ์กฐํšŒ
Args:
campaign_id: ํŠน์ • ์บ ํŽ˜์ธ๋งŒ ํ•„ํ„ฐ๋ง (None์ด๋ฉด ์ „์ฒด)
limit: ์ตœ๋Œ€ ๊ฒฐ๊ณผ ์ˆ˜
Returns:
Job ๋ชฉ๋ก (์ตœ์‹ ์ˆœ ์ •๋ ฌ)
"""
client = _get_client()
query = client.table("sentiment_analysis_jobs").select(
"id, campaign_id, status, progress, message, "
"total_answers, processed_answers, nudge_candidates, "
"created_at, started_at, completed_at, error_message"
).in_("status", ["queued", "running"])
if campaign_id:
query = query.eq("campaign_id", campaign_id)
query = query.order("created_at", desc=True).limit(limit)
result = query.execute()
return result.data or []
@st.cache_data(ttl=5)
def get_recent_jobs(
campaign_id: int | None = None,
status: JobStatus | None = None,
limit: int = 20,
) -> list[dict]:
"""์ตœ๊ทผ Job ๋ชฉ๋ก ์กฐํšŒ (ํžˆ์Šคํ† ๋ฆฌ์šฉ)
Args:
campaign_id: ํŠน์ • ์บ ํŽ˜์ธ๋งŒ ํ•„ํ„ฐ๋ง
status: ํŠน์ • ์ƒํƒœ๋งŒ ํ•„ํ„ฐ๋ง
limit: ์ตœ๋Œ€ ๊ฒฐ๊ณผ ์ˆ˜
Returns:
Job ๋ชฉ๋ก (์ตœ์‹ ์ˆœ ์ •๋ ฌ)
"""
client = _get_client()
query = client.table("sentiment_analysis_jobs").select(
"id, campaign_id, status, progress, message, "
"total_answers, processed_answers, nudge_candidates, "
"created_at, started_at, completed_at, error_message"
)
if campaign_id:
query = query.eq("campaign_id", campaign_id)
if status:
query = query.eq("status", status)
query = query.order("created_at", desc=True).limit(limit)
result = query.execute()
return result.data or []
def get_job_by_id(job_id: str) -> dict | None:
"""ํŠน์ • Job ์ƒ์„ธ ์กฐํšŒ
Args:
job_id: Job ID
Returns:
Job ์ •๋ณด ๋˜๋Š” None
"""
client = _get_client()
result = client.table("sentiment_analysis_jobs").select(
"id, campaign_id, status, progress, message, "
"total_answers, processed_answers, nudge_candidates, "
"created_at, started_at, completed_at, error_message"
).eq("id", job_id).execute()
return result.data[0] if result.data else None
def format_job_duration(job: dict) -> str:
"""Job ์†Œ์š” ์‹œ๊ฐ„ ํฌ๋งทํŒ…
Args:
job: Job ์ •๋ณด dict
Returns:
"2์‹œ๊ฐ„ 15๋ถ„" ํ˜•ํƒœ์˜ ๋ฌธ์ž์—ด
"""
started_at = job.get("started_at")
completed_at = job.get("completed_at")
if not started_at:
return "-"
try:
# ISO ๋ฌธ์ž์—ด ํŒŒ์‹ฑ
if isinstance(started_at, str):
start = datetime.fromisoformat(started_at.replace("Z", "+00:00"))
else:
start = started_at
if completed_at:
if isinstance(completed_at, str):
end = datetime.fromisoformat(completed_at.replace("Z", "+00:00"))
else:
end = completed_at
else:
end = datetime.now(timezone.utc)
seconds = int((end - start).total_seconds())
if seconds < 60:
return f"{seconds}์ดˆ"
elif seconds < 3600:
return f"{seconds // 60}๋ถ„"
else:
hours = seconds // 3600
minutes = (seconds % 3600) // 60
if minutes > 0:
return f"{hours}์‹œ๊ฐ„ {minutes}๋ถ„"
return f"{hours}์‹œ๊ฐ„"
except Exception:
return "-"
def get_status_emoji(status: str) -> str:
"""Job ์ƒํƒœ ์ด๋ชจ์ง€ ๋ฐ˜ํ™˜"""
return {
"queued": "๐ŸŸก",
"running": "๐Ÿ”ต",
"completed": "๐ŸŸข",
"failed": "๐Ÿ”ด",
"cancelled": "โšช",
}.get(status, "โšช")
def get_status_label(status: str) -> str:
"""Job ์ƒํƒœ ํ•œ๊ธ€ ๋ผ๋ฒจ ๋ฐ˜ํ™˜"""
return {
"queued": "๋Œ€๊ธฐ ์ค‘",
"running": "์‹คํ–‰ ์ค‘",
"processing": "์ฒ˜๋ฆฌ ์ค‘",
"completed": "์™„๋ฃŒ",
"failed": "์‹คํŒจ",
"cancelled": "์ทจ์†Œ๋จ",
}.get(status, status)
# ============================================================================
# Hierarchy Analysis Jobs (hierarchy_analysis_jobs table)
# ============================================================================
_HIERARCHY_FIELDS = (
"id, user_id, prompt, title, status, progress, current_step, "
"created_at, started_at, completed_at, error_message"
)
_HIERARCHY_STEP_LABELS: dict[str, str] = {
"prompt_enhancement": "ํ”„๋กฌํ”„ํŠธ ๋ถ„์„",
"query_generation": "์ฟผ๋ฆฌ ์ƒ์„ฑ",
"keyword_data": "ํ‚ค์›Œ๋“œ ์ˆ˜์ง‘",
"hierarchy_generation": "๊ณ„์ธต ๊ตฌ์กฐ ์ƒ์„ฑ",
"search_volume": "๊ฒ€์ƒ‰๋Ÿ‰ ์กฐํšŒ",
"question_generation": "์งˆ๋ฌธ ์ƒ์„ฑ",
"ratio_supplement": "๋น„์œจ ๋ณด์ •",
"finalization": "์ตœ์ข… ์ €์žฅ",
"done": "์™„๋ฃŒ",
}
def get_hierarchy_step_label(step: str | None) -> str:
"""Hierarchy ํŒŒ์ดํ”„๋ผ์ธ ์Šคํ… ํ•œ๊ธ€ ๋ผ๋ฒจ ๋ฐ˜ํ™˜"""
if not step:
return "๋Œ€๊ธฐ ์ค‘"
return _HIERARCHY_STEP_LABELS.get(step, step)
def get_active_hierarchy_jobs(user_id: str | None = None, limit: int = 5) -> list[dict]:
"""์‹คํ–‰ ์ค‘์ด๊ฑฐ๋‚˜ ๋Œ€๊ธฐ ์ค‘์ธ Hierarchy Job ๋ชฉ๋ก (Supabase ์ง์ ‘ ์ฟผ๋ฆฌ)
Args:
user_id: ์‚ฌ์šฉ์ž ID๋กœ ํ•„ํ„ฐ๋ง (None์ด๋ฉด ์ „์ฒด)
limit: ์ตœ๋Œ€ ๊ฒฐ๊ณผ ์ˆ˜
"""
client = _get_client()
query = (
client.table("hierarchy_analysis_jobs")
.select(_HIERARCHY_FIELDS)
.in_("status", ["queued", "processing"])
)
if user_id:
query = query.eq("user_id", user_id)
query = query.order("created_at", desc=True).limit(limit)
result = query.execute()
return result.data or []
def get_recent_hierarchy_jobs(
user_id: str | None = None,
status: str | None = None,
limit: int = 20,
) -> list[dict]:
"""์ตœ๊ทผ Hierarchy Job ์ด๋ ฅ ์กฐํšŒ
Args:
user_id: ์‚ฌ์šฉ์ž ID๋กœ ํ•„ํ„ฐ๋ง
status: ํŠน์ • ์ƒํƒœ๋งŒ ํ•„ํ„ฐ๋ง
limit: ์ตœ๋Œ€ ๊ฒฐ๊ณผ ์ˆ˜
"""
client = _get_client()
query = client.table("hierarchy_analysis_jobs").select(_HIERARCHY_FIELDS)
if user_id:
query = query.eq("user_id", user_id)
if status:
query = query.eq("status", status)
query = query.order("created_at", desc=True).limit(limit)
result = query.execute()
return result.data or []