"""hyungi_Document_Server — FastAPI 엔트리포인트""" from contextlib import asynccontextmanager from fastapi import FastAPI, Request from fastapi.responses import RedirectResponse from sqlalchemy import func, select, text from api.audio import router as audio_router from api.internal_study import router as internal_study_router from api.internal_worker import router as internal_worker_router from api.auth import router as auth_router from api.briefing import router as briefing_router from api.config import router as config_router from api.dashboard import router as dashboard_router from api.digest import router as digest_router from api.document_notes import router as document_notes_router from api.document_reads import router as document_reads_router from api.documents import router as documents_router from api.eid_chat import router as eid_chat_router from api.events import router as events_router from api.library import router as library_router from api.memos import router as memos_router from api.news import router as news_router from api.queue_overview import router as queue_overview_router from api.search import router as search_router from api.setup import router as setup_router from api.study_question_progress import router as study_question_progress_router from api.study_questions import router as study_questions_router from api.study_sessions import router as study_sessions_router from api.study_topics import router as study_topics_router from api.study_reminders import router as study_reminders_router from api.study_cards import router as study_cards_router from api.video import router as video_router from core.config import settings from core.database import async_session, engine, init_db from models.user import User @asynccontextmanager async def lifespan(app: FastAPI): """앱 시작/종료 시 실행되는 lifespan 핸들러""" import asyncio from apscheduler.schedulers.asyncio import AsyncIOScheduler from apscheduler.triggers.cron import CronTrigger from zoneinfo import ZoneInfo KST = ZoneInfo("Asia/Seoul") from services.search.query_analyzer import prewarm_analyzer from workers.briefing_worker import run as morning_briefing_run from workers.daily_digest import run as daily_digest_run from workers.dedup_reconcile import run as dedup_reconcile_run from workers.digest_worker import run as global_digest_run from workers.file_watcher import watch_inbox from workers.mailplus_archive import run as mailplus_run from workers.news_collector import run as news_collector_run from workers.fulltext_worker import reconcile_unresolved as fulltext_reconcile_run from workers.kosha_collector import run as kosha_collector_run from workers.csb_collector import run as csb_collector_run from workers.api_standards_collector import run as api_standards_run from workers.ccps_collector import run as ccps_collector_run from workers.queue_consumer import consume_queue, consume_fast_queue, consume_markdown_queue from workers.study_queue_consumer import consume_study_queue from workers.study_session_queue_consumer import consume_study_session_queue from workers.study_memo_card_jobs_consumer import consume_study_memo_card_queue from workers.study_card_enqueue import run as study_card_enqueue_run from workers.study_reminder import run as study_reminder_run from workers.study_weakness import run as study_weakness_run from workers.study_question_embed_worker import ( refresh_stale_related as study_q_related_refresh, run as study_q_embed_run, ) from workers.tier_backfill import run as tier_backfill_run from workers.upload_cleanup import cleanup_orphan_uploads # 시작: DB 연결 확인 await init_db() # NAS 마운트 확인 (NFS 미마운트 시 로컬 빈 디렉토리에 쓰는 것 방지) from pathlib import Path nas_check = Path(settings.nas_mount_path) / "PKM" if not nas_check.is_dir(): raise RuntimeError( f"NAS 마운트 확인 실패: {nas_check} 디렉토리 없음. " f"NFS 마운트 상태를 확인하세요." ) # APScheduler: 백그라운드 작업 scheduler = AsyncIOScheduler(timezone="Asia/Seoul") # 상시 실행 scheduler.add_job(consume_queue, "interval", minutes=1, id="queue_consumer") # PR-DocSrv-Markdown-Consumer-Split-1: markdown(marker) 전용 consumer. # 대형 PDF split 변환(수십 분)이 메인 consume_queue 를 점유해 전 파이프라인을 # stall 시키던 문제 제거. max_instances=1(기본) 으로 동시 marker 변환 2건은 방지. scheduler.add_job(consume_markdown_queue, "interval", minutes=1, id="markdown_consumer") # 2026-06-12 fast-consumer split: embed/chunk(건당 <1s)를 LLM 사이클에서 분리 — # classify(~190s×3)가 사이클을 점유해 벡터 적재가 굶던 구조 캡 해소 (markdown 선례). scheduler.add_job(consume_fast_queue, "interval", minutes=1, id="fast_queue_consumer") scheduler.add_job(watch_inbox, "interval", minutes=5, id="file_watcher") scheduler.add_job(cleanup_orphan_uploads, "interval", minutes=10, id="upload_cleanup") # PR-4: study_questions 자동 임베딩 (status='none/failed/stale' 행을 batch=10 처리). # 별도 큐 테이블 없이 status 자체가 큐. backfill 도 cron 이 'none' 행을 자연스럽게 처리. scheduler.add_job(study_q_embed_run, "interval", minutes=1, id="study_q_embed") # PR-12-A 후속: related-types 캐시 stale 행 재계산. 임베딩 워커와 분리한 별도 cron. # 새 문제 ready / 같은 토픽 invalidation / 임계값 변경 시 NULL 마킹된 행을 batch=20 처리. scheduler.add_job(study_q_related_refresh, "interval", minutes=1, id="study_q_related_refresh") # Phase 4-A: study_question_jobs 처리 — wrong/unsure AI 풀이 prefetch. # MLX gate 직렬화 + BATCH_SIZE=1 로 GPU 부하 통제. STALE_MINUTES=10 자체 복구. scheduler.add_job(consume_study_queue, "interval", minutes=1, id="study_queue_consumer") # Phase 4-B v1: study_quiz_session_jobs 처리 — 세션 단위 자유 마크다운 분석. # 4-A 와 같은 MLX gate 공유 — 4-A 처리 중이면 직렬 대기. scheduler.add_job(consume_study_session_queue, "interval", minutes=1, id="study_session_queue_consumer") # 공부 암기노트 Phase 1: card_extract 큐 consumer + 버전키 폴러(study_card_enqueue). # 별 테이블/별 consumer 로 기존 study queue 와 격리. settings.study_card_extract_enabled 게이트. scheduler.add_job(consume_study_memo_card_queue, "interval", minutes=1, id="study_memo_card_consumer") scheduler.add_job(study_card_enqueue_run, "interval", minutes=1, id="study_card_enqueue") # PR-B 레거시 tier 백필 — 30분 주기로 호출되지만 KST 00:00~06:00 시간대만 실제 enqueue. # safety > law > manual 우선순위로 25건씩. 6720 레거시 → 야간당 ~150건 → 약 45일 소화. scheduler.add_job(tier_backfill_run, "interval", minutes=30, id="tier_backfill") # 일일 스케줄 (KST) # law_monitor 스케줄 제거 (safety-library-1 B-1 PR①, 2026-06-13) — 매일 버전 체인 밖 # 레거시 스냅샷을 증식하던 유일 경로 차단. 파일은 강등 보존(1사이클 관찰 후 삭제), # 대체 = statute_collector (스케줄 등록은 PR② 잡 코드와 함께 — R8-B1). scheduler.add_job(mailplus_run, CronTrigger(hour=7, timezone=KST), id="mailplus_morning") scheduler.add_job(mailplus_run, CronTrigger(hour=18, timezone=KST), id="mailplus_evening") scheduler.add_job(daily_digest_run, CronTrigger(hour=20, timezone=KST), id="daily_digest") scheduler.add_job(global_digest_run, CronTrigger(hour=4, minute=0, timezone=KST), id="global_digest") scheduler.add_job(morning_briefing_run, CronTrigger(hour=5, minute=10, timezone=KST), id="morning_briefing") # 공부 암기노트 Phase 1: 공부중 토픽 due 요약 알람 재료 (09/13/19 KST). LLM 0. scheduler.add_job(study_reminder_run, CronTrigger(hour="9,13,19", timezone=KST), id="study_reminder") # 이드 W3-2: 공부중 토픽 약점 derived 스냅샷 (nightly 04:30 KST, LLM 0). study_diagnosis 표면 source. scheduler.add_job(study_weakness_run, CronTrigger(hour=4, minute=30, timezone=KST), id="study_weakness") scheduler.add_job(news_collector_run, "interval", hours=6, id="news_collector") # crawl-24x7 A-2 안전망: fulltext 영구 실패(3회 소진) 문서를 RSS 요약 기준으로 # 후속 enqueue (silent skip 누적 방지). 03:40 = dedup_reconcile(03:30) 직후 비충돌 슬롯. scheduler.add_job(fulltext_reconcile_run, CronTrigger(hour=3, minute=40, timezone=KST), id="fulltext_reconcile") # plan ds-s1-backend-1 B-4: dedup 컬럼(duplicate_of/duplicate_count) 야간 절대 재계산. # soft-delete 잔여 드리프트 정리(멱등, 드리프트 없으면 no-op). cron 03:30 (다른 잡과 비충돌). scheduler.add_job(dedup_reconcile_run, CronTrigger(hour=3, minute=30, timezone=KST), id="dedup_reconcile") # crawl-24x7 C-2: KOSHA 재해사례 diff + GUIDE 점진 백필 (daily, 새벽 잡들과 비충돌 슬롯). scheduler.add_job(kosha_collector_run, CronTrigger(hour=6, minute=40, timezone=KST), id="kosha_collector") # 사이클 3 C-2 잔여: CSB sitemap lastmod diff (weekly 월, cap 40 + 워터마크 점진 백필). scheduler.add_job(csb_collector_run, CronTrigger(day_of_week="mon", hour=6, minute=50, timezone=KST), id="csb_collector") # 사이클 3 C-4: API 표준 공지 목록 diff (monthly — 월 1~2건 공지 페이스). scheduler.add_job(api_standards_run, CronTrigger(day=5, hour=7, minute=5, timezone=KST), id="api_standards_collector") # 사이클 3 C-2 잔여: CCPS Beacon 월간 PDF (playwright 익명 경유 — WAF 차단 시 health 로 가시화). scheduler.add_job(ccps_collector_run, CronTrigger(day=5, hour=7, minute=20, timezone=KST), id="ccps_collector") scheduler.start() # Phase 2.1 (async 구조): QueryAnalyzer prewarm. # 대표 쿼리 15~20개를 background task로 분석해 cache 적재. # 첫 사용자 요청부터 cache hit rate 70~80% 목표. # 논블로킹 — startup을 막지 않음. MLX 부하 완화 위해 delay_between=0.5. prewarm_task = asyncio.create_task(prewarm_analyzer()) prewarm_task.add_done_callback( lambda t: t.exception() and None # 예외는 query_analyzer 내부에서 로깅 ) yield # 종료: 스케줄러 → DB 순서로 정리 scheduler.shutdown(wait=False) await engine.dispose() app = FastAPI( title="hyungi_Document_Server", description="Self-hosted PKM 웹 애플리케이션 API", version="2.0.0", lifespan=lifespan, ) # ─── 라우터 등록 ─── app.include_router(setup_router, prefix="/api/setup", tags=["setup"]) app.include_router(config_router, prefix="/api/config", tags=["config"]) app.include_router(auth_router, prefix="/api/auth", tags=["auth"]) app.include_router(documents_router, prefix="/api/documents", tags=["documents"]) # 회독 카운트 — /api/documents/{id}/read* 경로. documents_router 와 prefix 같아 충돌 없음. app.include_router(document_reads_router, prefix="/api/documents", tags=["document-reads"]) app.include_router(document_notes_router, prefix="/api/documents", tags=["document-notes"]) app.include_router(search_router, prefix="/api/search", tags=["search"]) # 이드 채팅 표면 (D-1) — POST /api/eid/chat. SSE 스트리밍, EidAIClient.call_stream 봉쇄 경유. app.include_router(eid_chat_router, prefix="/api/eid", tags=["eid-chat"]) app.include_router(memos_router, prefix="/api/memos", tags=["memos"]) app.include_router(events_router, prefix="/api/events", tags=["events"]) app.include_router(dashboard_router, prefix="/api/dashboard", tags=["dashboard"]) app.include_router(library_router, prefix="/api/library", tags=["library"]) app.include_router(news_router, prefix="/api/news", tags=["news"]) # 처리 머신 보드 (plan ds-processing-ui-6an) — GET /api/queue/overview app.include_router(queue_overview_router, prefix="/api/queue", tags=["queue"]) app.include_router(digest_router, prefix="/api/digest", tags=["digest"]) app.include_router(briefing_router, prefix="/api/briefing", tags=["briefing"]) app.include_router(audio_router, prefix="/api/audio", tags=["audio"]) app.include_router(internal_study_router, prefix="/internal/study", tags=["internal-study"]) app.include_router(internal_worker_router, prefix="/internal/worker", tags=["internal-worker"]) app.include_router(video_router, prefix="/api/video", tags=["video"]) app.include_router(study_sessions_router, prefix="/api/study-sessions", tags=["study-sessions"]) app.include_router(study_topics_router, prefix="/api/study-topics", tags=["study-topics"]) # study_questions: 라우터 안에서 /study-topics/{id}/questions 와 /study-questions/{id} 두 줄기를 모두 정의하므로 prefix=/api 로 등록 app.include_router(study_questions_router, prefix="/api", tags=["study-questions"]) app.include_router(study_reminders_router, prefix="/api/study-reminders", tags=["study-reminders"]) app.include_router(study_cards_router, prefix="/api/study-cards", tags=["study-cards"]) # Phase 1: 학습 진행 상태 (review-complete + review-queue). prefix=/api/study-topics 안에 정의됨. app.include_router(study_question_progress_router, prefix="/api", tags=["study-progress"]) # TODO: Phase 5에서 추가 # app.include_router(tasks.router, prefix="/api/tasks", tags=["tasks"]) # app.include_router(export.router, prefix="/api/export", tags=["export"]) # ─── 셋업 미들웨어: 유저 0명이면 /setup으로 리다이렉트 ─── SETUP_BYPASS_PREFIXES = ( "/api/setup", "/api/config", "/setup", "/health", "/docs", "/openapi.json", "/redoc", ) @app.middleware("http") async def setup_redirect_middleware(request: Request, call_next): path = request.url.path # 바이패스 경로는 항상 통과 if any(path.startswith(p) for p in SETUP_BYPASS_PREFIXES): return await call_next(request) # 유저 존재 여부 확인 try: async with async_session() as session: result = await session.execute(select(func.count(User.id))) user_count = result.scalar() if user_count == 0: return RedirectResponse(url="/setup") except Exception: pass # DB 연결 실패 시 통과 (health에서 확인 가능) return await call_next(request) # ─── 셋업 페이지 라우트 (API가 아닌 HTML 페이지) ─── @app.get("/setup") async def setup_page_redirect(request: Request): """셋업 위자드 페이지로 포워딩""" from api.setup import setup_page from core.database import get_session async for session in get_session(): return await setup_page(request, session) @app.get("/health") async def health_check(): """헬스체크 — DB 연결 상태 포함""" db_ok = False try: async with engine.connect() as conn: await conn.execute(text("SELECT 1")) db_ok = True except Exception: pass return { "status": "ok" if db_ok else "degraded", "version": "2.0.0", "database": "connected" if db_ok else "disconnected", }