"""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 []