"""
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
def clean_confluence_storage_xhtml(xhtml_content: str) -> str:
"""Confluence 매크로(code, info panel, status 등)를 표준 HTML/Markdown으로 정규화합니다."""
if not xhtml_content:
return ""
content = xhtml_content
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)
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)
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)
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
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
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
}
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']}"
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"]
}
}
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)
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 = []
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})")
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()
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
tasks = [enrich_confluence_chunk(client, args.model, c, semaphore) for c in initial_chunks]
enriched_chunks = await asyncio.gather(*tasks)
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)
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 생성 파이프라인")
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="대상 시스템 컴포넌트 명칭")
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()