Spaces:
Sleeping
Sleeping
| """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 | |
| 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"] | |
| 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 [] | |
| 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 [] | |