26S26g2

QK·2026년 9월 27일
#!/usr/bin/env python3
"""
Confluence Docs Importer & vLLM Metadata Enricher for VectorDB
- Confluence REST API(단일 페이지 또는 Space 전체 CQL) 호출
- Storage Format(XHTML) -> 정제된 Markdown 변환 및 코드/명령어 보존 청킹
- B300 vLLM 고속 배치를 통한 RCA 특화 메타데이터 추출:
  - doc_category (POST_MORTEM, SOP, TROUBLESHOOTING, ARCHITECTURE, etc.)
  - symptoms_and_errors (발생 증상, 에러 코드, 알람명)
  - related_metrics (연관 Prometheus/MinIO 메트릭 명칭)
  - target_resources (Node, PVC, Pod, Drive 등 인프라 리소스)
  - root_cause_hypotheses (장애 보고서인 경우 언급된 실제 근본 원인 요약)
  - is_destructive_action (SOP 조치 시 파괴적 명령 포함 여부)
  - questions (HyDE 가상 질문)
- Vector DB 일괄 적재(Bulk Upsert)용 단일 통합 JSON 생성
"""

import os
import re
import json
import asyncio
import argparse
from typing import List, Dict, Any, Optional, Tuple
import requests
from bs4 import BeautifulSoup
from markdownify import markdownify as md
from openai import AsyncOpenAI
from langchain_text_splitters import MarkdownHeaderTextSplitter, RecursiveCharacterTextSplitter


# ==============================================================================
# 1. Confluence REST API 및 HTML -> Markdown 전처리
# ==============================================================================
def clean_confluence_storage_xhtml(xhtml_content: str) -> str:
    """Confluence 매크로(code, info panel, status 등)를 표준 HTML/Markdown으로 정규화합니다."""
    if not xhtml_content:
        return ""
    content = xhtml_content

    # 1. Code Block 매크로 변환 (<ac:structured-macro ac:name="code"> -> <pre><code>)
    code_pattern = re.compile(
        r'<ac:structured-macro[^>]*ac:name="code"[^>]*>.*?'
        r'<ac:plain-text-body><!\[CDATA\[(.*?)\]\]></ac:plain-text-body>.*?'
        r'</ac:structured-macro>',
        re.DOTALL
    )
    content = code_pattern.sub(r'<pre><code>\1</code></pre>', content)

    # 2. Panel 매크로 변환 (info, note, warning, tip -> <blockquote>)
    panel_pattern = re.compile(
        r'<ac:structured-macro[^>]*ac:name="(info|note|warning|tip)"[^>]*>.*?'
        r'<ac:rich-text-body>(.*?)</ac:rich-text-body>.*?'
        r'</ac:structured-macro>',
        re.DOTALL
    )
    content = panel_pattern.sub(r'<blockquote>\2</blockquote>', content)

    # 3. Status/Tag 매크로 변환
    status_pattern = re.compile(
        r'<ac:structured-macro[^>]*ac:name="status"[^>]*>.*?'
        r'<ac:parameter ac:name="title">(.*?)</ac:parameter>.*?'
        r'</ac:structured-macro>',
        re.DOTALL
    )
    content = status_pattern.sub(r'[\1]', content)

    # 4. 기타 Confluence 태그 제거
    content = re.sub(r'</?ac:[^>]*>', '', content)
    content = re.sub(r'</?ri:[^>]*>', '', content)
    content = re.sub(r'\s+xmlns(?::\w+)?="[^"]*"', '', content)
    return content


def get_confluence_pages(
    base_url: str,
    page_id: Optional[str] = None,
    space_key: Optional[str] = None,
    auth_token: Optional[str] = None,
    username: Optional[str] = None,
    password: Optional[str] = None,
    verify_ssl: bool = True,
    limit: int = 50
) -> List[Dict[str, Any]]:
    """단일 페이지 또는 Space 내 페이지 목록을 REST API로 조회합니다."""
    headers = {"Accept": "application/json"}
    auth = None
    if auth_token:
        headers["Authorization"] = f"Bearer {auth_token}"
    elif username and password:
        auth = (username, password)

    results = []
    base = base_url.rstrip("/")

    if page_id:
        url = f"{base}/rest/api/content/{page_id}?expand=body.storage,version,space,history"
        resp = requests.get(url, headers=headers, auth=auth, verify=verify_ssl, timeout=15)
        resp.raise_for_status()
        results.append(resp.json())
    elif space_key:
        url = f"{base}/rest/api/content?spaceKey={space_key}&type=page&expand=body.storage,version,space,history&limit={limit}"
        resp = requests.get(url, headers=headers, auth=auth, verify=verify_ssl, timeout=20)
        resp.raise_for_status()
        data = resp.json()
        results.extend(data.get("results", []))
    else:
        raise ValueError("page_id 또는 space_key 중 하나는 반드시 지정해야 합니다.")

    return results


# ==============================================================================
# 2. 코드 보존형 계층 청킹
# ==============================================================================
def split_markdown_preserving_code(
    md_text: str, 
    page_meta: Dict[str, Any],
    default_component: str
) -> List[Dict[str, Any]]:
    """Markdown 헤더 계층을 보존하고 코드/로그 블록이 중간에 끊기지 않도록 청킹합니다."""
    headers_to_split = [
        ("#", "h1"),
        ("##", "h2"),
        ("###", "h3"),
        ("####", "h4"),
    ]
    header_splitter = MarkdownHeaderTextSplitter(
        headers_to_split_on=headers_to_split, 
        strip_headers=False
    )
    sections = header_splitter.split_text(md_text)

    text_splitter = RecursiveCharacterTextSplitter(
        chunk_size=1200,
        chunk_overlap=150,
        separators=[
            "\n```\n",
            "\n```bash\n",
            "\n```yaml\n",
            "\n```json\n",
            "\n\n",
            "\n",
            " "
        ],
        keep_separator=True
    )

    page_id = page_meta.get("id", "")
    page_title = page_meta.get("title", f"confluence_{page_id}")
    source_url = page_meta.get("source_url", "")
    chunks = []
    chunk_idx = 0

    for section in sections:
        sub_docs = text_splitter.split_documents([section])
        headers = [section.metadata.get(h) for h in ["h1", "h2", "h3", "h4"] if section.metadata.get(h)]
        breadcrumb = " > ".join(headers) if headers else page_title

        for doc in sub_docs:
            chunk_idx += 1
            chunks.append({
                "chunk_id": f"confluence_{page_id}_{chunk_idx:04d}",
                "page_id": page_id,
                "title": page_title,
                "component": default_component,
                "source_url": source_url,
                "breadcrumb": breadcrumb,
                "content": doc.page_content,
                "header_metadata": section.metadata
            })

    return chunks


# ==============================================================================
# 3. vLLM 확장 메타데이터 추출 엔진 (SOP & 장애 이력 특화)
# ==============================================================================
PROMPT_TEMPLATE = """당신은 대규모 클라우드 인프라 및 쿠버네티스/스토리지 SRE/DevOps 전문가입니다.
주어진 사내 Confluence 문서(장애 이력 보고서 또는 SOP/운영 매뉴얼) 청크를 정밀 분석하여, 장애 분석 시스템(Incident RCA Agent)과 Vector DB 검색에 필요한 메타데이터를 정형 JSON 형태로 추출하십시오.

[문서 제목 및 헤더 경로]
제목: {title}
경로: {breadcrumb}

[문서 내용]
{content}

[추출 규칙 및 필드 설명]
1. summary: 이 청크의 핵심 내용(장애 원인, 조치 절차 등) 1줄 요약
2. doc_category: 다음 중 가장 적합한 카테고리 하나만 선택 ["POST_MORTEM", "SOP", "TROUBLESHOOTING", "ARCHITECTURE", "CONFIGURATION", "GENERAL"]
3. questions: 운영자가 장애 발생 시 또는 트러블슈팅 중 검색할 법한 구체적 가상 질문 2~3개
4. symptoms_and_errors: 본문에 언급된 장애 증상, 알람명, 에러 로그 메시지, HTTP 상태 코드 목록 (예: ["Node Disk Latency High", "503 Slow Down", "FailedMount"])
5. related_metrics: 본문 내용과 직결되는 Prometheus/MinIO 메트릭 명칭 목록 (예: ["minio_node_drive_total_duration_microseconds", "go_goroutines"])
6. target_resources: 연관된 쿠버네티스/인프라 리소스 목록 (예: ["Node", "DirectPVDrive", "Pod", "Service", "PVC"])
7. root_cause_hypotheses: 본문이 장애 이력(POST_MORTEM)인 경우 서술된 실제 근본 원인 목록 (SOP인 경우 빈 배열)
8. commands: 본문에 등장하는 핵심 CLI 명령어 또는 점검 커맨드 목록 (예: ["mc admin drive offline ...", "kubectl describe pod ..."])
9. is_destructive_action: 본문에 포함된 조치 커맨드가 재기동, 드라이브 포맷, 오프라인 격리 등 서비스 영향이 있는 작업인지 여부 (true 또는 false)

반드시 마크다운 코드블록이나 불필요한 설명 없이 순수 JSON 문자열만 출력하십시오.
{{
  "summary": "1줄 요약",
  "doc_category": "POST_MORTEM",
  "questions": ["질문 1", "질문 2"],
  "symptoms_and_errors": ["증상 1", "에러 2"],
  "related_metrics": ["metric_name_1"],
  "target_resources": ["Node", "Pod"],
  "root_cause_hypotheses": ["특정 NVMe I/O 행으로 인한 분산락 경합"],
  "commands": ["명령어 1"],
  "is_destructive_action": false
}}"""

def clean_json_response(raw_text: str) -> Dict[str, Any]:
    """vLLM 응답에서 마크다운 코드블록을 제거하고 순수 JSON 객체로 파싱합니다."""
    text = raw_text.strip()
    if text.startswith("```"):
        text = re.sub(r"^```[a-zA-Z]*\n", "", text)
        text = re.sub(r"\n```$", "", text).strip()
    return json.loads(text)


async def enrich_confluence_chunk(
    client: AsyncOpenAI, 
    model: str, 
    chunk: Dict[str, Any], 
    semaphore: asyncio.Semaphore
) -> Dict[str, Any]:
    """vLLM 비동기 호출을 통해 청크를 보강합니다."""
    async with semaphore:
        prompt = PROMPT_TEMPLATE.format(
            title=chunk["title"],
            breadcrumb=chunk["breadcrumb"],
            content=chunk["content"]
        )
        try:
            response = await client.chat.completions.create(
                model=model,
                messages=[{"role": "user", "content": prompt}],
                temperature=0.1,
                max_tokens=650,
                response_format={"type": "json_object"}
            )
            enrichment_data = clean_json_response(response.choices[0].message.content)
        except Exception:
            enrichment_data = {
                "summary": chunk["breadcrumb"],
                "doc_category": "SOP",
                "questions": [],
                "symptoms_and_errors": [],
                "related_metrics": [],
                "target_resources": [],
                "root_cause_hypotheses": [],
                "commands": [],
                "is_destructive_action": False
            }

        # 1) Vector 임베딩 검색용 상단 컨텍스트 헤더 구성
        q_lines = "\n".join([f"- Q: {q}" for q in enrichment_data.get("questions", [])])
        symptoms_str = ", ".join(enrichment_data.get("symptoms_and_errors", [])) or "None"
        resources_str = ", ".join(enrichment_data.get("target_resources", [])) or "General"
        causes_str = ", ".join(enrichment_data.get("root_cause_hypotheses", [])) or "N/A"
        category = enrichment_data.get("doc_category", "SOP")

        header_block = (
            f"[Source: Confluence | Title: {chunk['title']}]\n"
            f"[Path: {chunk['breadcrumb']}]\n"
            f"[Category: {category}] [Resources: {resources_str}]\n"
            f"[Summary: {enrichment_data.get('summary', '')}]\n"
            f"[Symptoms/Alerts: {symptoms_str}]\n"
            f"[Known Root Cause: {causes_str}]"
        )
        if q_lines:
            header_block += f"\n{q_lines}"

        enriched_content = f"{header_block}\n\n{chunk['content']}"

        # 2) Vector DB Upsert용 정규화 스키마 반환
        return {
            "id": chunk["chunk_id"],
            "page_content": enriched_content.strip(),
            "metadata": {
                "doc_type": "confluence",
                "page_id": chunk["page_id"],
                "title": chunk["title"],
                "component": chunk["component"],
                "doc_category": category,
                "source_url": chunk["source_url"],
                "breadcrumb": chunk["breadcrumb"],
                "summary": enrichment_data.get("summary", ""),
                "questions": enrichment_data.get("questions", []),
                "symptoms_and_errors": enrichment_data.get("symptoms_and_errors", []),
                "related_metrics": enrichment_data.get("related_metrics", []),
                "target_resources": enrichment_data.get("target_resources", []),
                "root_cause_hypotheses": enrichment_data.get("root_cause_hypotheses", []),
                "commands": enrichment_data.get("commands", []),
                "is_destructive_action": bool(enrichment_data.get("is_destructive_action", False)),
                **chunk["header_metadata"]
            }
        }


# ==============================================================================
# 4. 전체 파이프라인 제어
# ==============================================================================
async def run_pipeline(args):
    os.makedirs(args.output_dir, exist_ok=True)
    raw_md_dir = os.path.join(args.output_dir, "raw_confluence_md")
    os.makedirs(raw_md_dir, exist_ok=True)

    print("=" * 65)
    print(f"[*] Confluence 서버 접속: {args.url}")
    print(f"[*] 대상 컴포넌트: {args.component}")
    print(f"[*] vLLM 엔드포인트: {args.vllm_url} (Model: {args.model})")
    print("=" * 65)

    # 1. Confluence REST API 호출
    pages = get_confluence_pages(
        base_url=args.url,
        page_id=args.page_id,
        space_key=args.space,
        auth_token=args.token,
        username=args.user,
        password=args.password,
        verify_ssl=not args.no_verify_ssl,
        limit=args.limit
    )

    if not pages:
        print("[!] 조회된 Confluence 페이지가 없습니다.")
        return

    print(f"[*] 총 {len(pages)}개의 Confluence 페이지를 성공적으로 가져왔습니다.")

    client = AsyncOpenAI(base_url=args.vllm_url, api_key="EMPTY")
    semaphore = asyncio.Semaphore(args.concurrency)
    all_chunks_for_vectordb = []

    # 2. 각 페이지별 변환, 청킹, vLLM 보강
    for idx, page in enumerate(pages, 1):
        p_id = page.get("id", "")
        p_title = page.get("title", f"page_{p_id}")
        source_url = f"{args.url.rstrip('/')}/pages/viewpage.action?pageId={p_id}"
        print(f"\n[{idx}/{len(pages)}] 처리 중: {p_title} (ID: {p_id})")

        # HTML -> Markdown 변환
        raw_storage = page.get("body", {}).get("storage", {}).get("value", "")
        clean_html = clean_confluence_storage_xhtml(raw_storage)
        clean_markdown = md(
            clean_html,
            heading_style="ATX",
            autolinks=True,
            code_language="bash",
            strip=["script", "style"]
        ).strip()

        # 원본 Markdown 로컬 백업 저장
        safe_title = re.sub(r'[\\/*?:"<>| ]', '_', p_title).strip()
        md_save_path = os.path.join(raw_md_dir, f"{p_id}_{safe_title}.md")
        with open(md_save_path, "w", encoding="utf-8") as f:
            f.write(f"# {p_title}\n\n{clean_markdown}")

        page_meta = {"id": p_id, "title": p_title, "source_url": source_url}
        initial_chunks = split_markdown_preserving_code(clean_markdown, page_meta, args.component)
        print(f"    - 분할된 청크 수: {len(initial_chunks)}개 (vLLM 분석 중...)")

        if not initial_chunks:
            continue

        # vLLM 고속 비동기 배치 보강
        tasks = [enrich_confluence_chunk(client, args.model, c, semaphore) for c in initial_chunks]
        enriched_chunks = await asyncio.gather(*tasks)

        # 개별 파일별 중간 JSON 저장
        per_page_json = os.path.join(args.output_dir, f"confluence_{p_id}_{safe_title}.json")
        with open(per_page_json, "w", encoding="utf-8") as f:
            json.dump(enriched_chunks, f, ensure_ascii=False, indent=2)

        all_chunks_for_vectordb.extend(enriched_chunks)

    # 3. Vector DB Bulk Upsert용 최종 단일 통합 JSON 생성
    final_output_path = os.path.join(args.output_dir, args.final_json)
    with open(final_output_path, "w", encoding="utf-8") as f:
        json.dump(all_chunks_for_vectordb, f, ensure_ascii=False, indent=2)

    print("\n" + "=" * 65)
    print(f"[✓] Confluence 문서 수집 및 vLLM 메타데이터 보강 완료!")
    print(f"[*] 총 생성된 Vector DB 청크 수: {len(all_chunks_for_vectordb)}개")
    print(f"[*] 원본 Markdown 백업 위치: {os.path.abspath(raw_md_dir)}")
    print(f"[*] 최종 단일 통합 JSON 파일 : {os.path.abspath(final_output_path)}")
    print("=" * 65)


def main():
    parser = argparse.ArgumentParser(description="Confluence 문서 API 조회 -> MD 변환 -> vLLM 기반 메타데이터 보강 -> 통합 JSON 생성 파이프라인")
    # Confluence API 옵션
    parser.add_argument("--url", default=os.getenv("CONFLUENCE_URL", "https://confluence.internal.net"), help="Confluence Base URL")
    parser.add_argument("--page-id", help="단일 페이지 ID (지정 시 1개 페이지만 처리)")
    parser.add_argument("--space", help="Confluence Space Key (지정 시 해당 Space 전체 페이지 조회)")
    parser.add_argument("--token", default=os.getenv("CONFLUENCE_TOKEN"), help="Confluence Personal Access Token (PAT)")
    parser.add_argument("--user", default=os.getenv("CONFLUENCE_USER"), help="Basic Auth 사용자명")
    parser.add_argument("--password", default=os.getenv("CONFLUENCE_PASSWORD"), help="Basic Auth 비밀번호/API토큰")
    parser.add_argument("--limit", type=int, default=100, help="Space 조회 시 최대 페이지 수")
    parser.add_argument("--no-verify-ssl", action="store_true", help="사내 사설 인증서 SSL 검증 비활성화")
    parser.add_argument("--component", default="minio-aistor", help="대상 시스템 컴포넌트 명칭")

    # 출력 및 vLLM 옵션
    parser.add_argument("--output-dir", default="./confluence_vectordb_dataset", help="출력 저장 디렉터리")
    parser.add_argument("--final-json", default="bulk_confluence_chunks.json", help="최종 단일 통합 파일명")
    parser.add_argument("--vllm-url", default="http://localhost:8000/v1", help="vLLM API 엔드포인트")
    parser.add_argument("--model", default="/data/models/Qwen2.5-32B-Instruct", help="vLLM에 서빙 중인 모델 이름/경로")
    parser.add_argument("--concurrency", type=int, default=48, help="vLLM 동시 요청 수 (B300 기준 48~64 권장)")

    args = parser.parse_args()
    if not args.page_id and not args.space:
        parser.error("--page-id 또는 --space 중 하나는 반드시 지정해야 합니다.")

    asyncio.run(run_pipeline(args))


if __name__ == "__main__":
    main()
profile
engineer

0개의 댓글