"""Stage 3 실행기 — 프로그램 요약 → unit 별 로직 조각 추출 (docs/logic-chunk-design.md). python -m summarize.runner [--program X] [--limit N] [--fake] [--dry-run] 흐름 (프로그램마다): 1. 프로그램 요약(ProgramSummary) — 없거나 stale 이거나 프롬프트 버전이 지났으면 생성 2. unit(FORM/METHOD/FUNCTION/MODULE/EVENT) 하나씩 → 프로그램 요약을 문맥으로 붙여 LLM 이 로직 조각(LogicChunk) 을 골라낸다. 300줄 넘는 unit 은 파서 sub_chunks 창으로 나눠 호출. 3. 조각은 summarize/chunks.py 로 검증(줄 앵커) · 사실 채움(파서) 후 logic_chunk / chunk_fts 에 저장. 멱등성: unit.summary_status='done' 이고 prompt_version 이 같으면 스킵 (code_hash 가 바뀐 unit 은 loader 가 재적재 시 상태를 초기화한다). LLM Key 미확보 상태에서는 --fake 로 파이프라인만 검증. """ from __future__ import annotations import argparse import json import re import sqlite3 import sys import time from concurrent.futures import ThreadPoolExecutor, as_completed from datetime import datetime, timezone from config.settings import settings from index.db import bigrams, clean_comment, connect, loads, text_symbol_phrases from index.loader import refresh_program_fts from query.tools import _structure_summary from .chunks import ResolvedChunk, covered_lines, numbered_code, parser_hints, resolve_chunks from .llm_client import PendingResponse, create_llm from .schemas import CHUNK_KINDS, ProgramSummary, UnitExtraction PROMPT_VERSION = 2 ELIGIBLE = ("FORM", "METHOD", "FUNCTION", "MODULE", "EVENT") MAX_WINDOW_LINES = 300 # 이보다 긴 unit 은 창으로 나눠 LLM 에 보여준다 FALLBACK_WINDOW = 250 # sub_chunks 가 없을 때의 고정 창 크기 MAX_EVENT_CODE_LINES = 150 # 프로그램 요약에 넣는 이벤트 블록 코드 상한 SYSTEM_PROMPT = ( "당신은 SAP ABAP 시니어 개발자다. 코드를 업무 관점으로 한국어로 설명한다. " "SAP 표준 영어 용어(예: Goods Receipt, G/L Account)와 기술 객체명(테이블·FM·BAPI)은 함께 적는다. " "추측이 필요한 부분은 confidence 를 낮추고 unclear 에 기록한다. 출력은 JSON 만." ) PROGRAM_PROMPT = """[작업] program_summary 다음 ABAP 프로그램을 업무 관점으로 요약하라. ProgramSummary JSON 스키마로만 답하라. [프로그램] {program} — {title} [패키지] {devclass} {pkg_text} [선택화면] {selection} [텍스트 심볼] {text_symbols} [구조 사실 — 파서 추출] 주 흐름(이벤트 → PERFORM 체인): {main_flow} 테이블 read: {tables_read} 테이블 write: {tables_write} 외부 호출: {external_calls} 출력 형태: {output_type} [unit 목록] (유형 이름 — 헤더 주석) {unit_list} [이벤트 블록 코드 — 줄번호| 내용] {event_code} [출력 JSON 스키마 — 아래 키만 사용] {{"program": "{program}", "title_ko": "", "business_purpose_ko": "업무 목적 2~3문장", "business_purpose_en": "", "main_flow": ["처리 순서를 업무 언어로"], "selection_screen": [{{"name": "", "desc_ko": ""}}], "key_internal_tables": [{{"name": "", "filled_by": [], "consumed_by": [], "desc_ko": ""}}], "output_type": [], "business_tags": [], "sap_module": "FI/CO/MM/SD 등", "keywords_ko": ["한국어 검색 키워드"], "keywords_en": ["English search keywords"], "related_tcodes": [], "notes": [], "confidence": 0.0, "unclear": []}} """ UNIT_PROMPT = """[작업] extract_chunks 다음 ABAP 코드 단위에서 **업무적으로 의미 있는 로직 조각**을 골라내고 각각을 설명하라. UnitExtraction JSON 스키마로만 답하라. [프로그램] {program} — {title} [프로그램 요약] {program_purpose} [주 흐름] {main_flow} [핵심 내부테이블] {key_tables} [단위] unit_id: {unit_id} unit_type: {unit_type}, name: {name} signature: {signature} header_comment: {header_comment} {window_note} [파서 힌트 — DB 접근 · 호출 · 검증이 시작되는 줄] {hints} [텍스트 심볼] {text_symbols} [코드 — 줄번호| 내용] {code} [조각 선정 규칙] - 조각 = 하나의 업무 로직을 이루는 연속된 줄 범위. 예: 특정 테이블 SELECT 와 그 직후 결과 정리(SORT/DELETE/LOOP 집계) 까지가 한 로직이면 하나로 묶는다. BAPI/FM 호출은 파라미터 채우기 ~ 호출 ~ 결과 처리까지 한 조각. - 조각으로 만들지 않고 버리는 것: 단순 선언(DATA/TYPES), 변수 초기화(CLEAR/REFRESH/FREE), ALV 필드카탈로그·레이아웃· 컬럼 속성 설정, 화면 속성 LOOP AT SCREEN, 단순 이벤트 등록, 로그성 WRITE, 주석만 있는 구간. - 조각은 서로 겹치지 않는다. 보통 3~80줄. FORM/METHOD 전체를 통째로 하나의 조각으로 만들지 않는다 (내부에 로직이 하나뿐이면 그 로직 범위만). - 의미 있는 조각이 없으면 chunks 를 빈 배열로 둔다. 그래도 unit_purpose_ko 는 채운다. - line_start / line_end 는 위 코드에 붙은 줄번호 그대로. first_line 에는 line_start 줄의 코드를 그대로 복사한다. - kind 는 다음 중 하나: {kinds} - purpose_ko 는 "무엇을 왜 하는지" 한 문장(한국어). keywords_ko 는 업무 용어(입고, 반제, 잔액 …), keywords_en 은 SAP 표준 영어 용어, sap_objects 는 테이블·FM·BAPI·클래스·T-Code 이름. [출력 JSON 스키마] {{"unit_purpose_ko": "이 단위의 역할 한 문장", "chunks": [{{"line_start": 0, "line_end": 0, "first_line": "", "kind": "sql_select", "purpose_ko": "", "purpose_en": "", "keywords_ko": [], "keywords_en": [], "sap_objects": [], "confidence": 0.0}}], "unclear": []}} """ _SELSCREEN = re.compile(r"^\s*(PARAMETERS?|SELECT-OPTIONS|SELECTION-SCREEN\s+BEGIN\s+OF\s+BLOCK)\b(.*)$", re.I) _TEXT_SYM = re.compile(r"TEXT-(\w{3})", re.I) def _now() -> str: return datetime.now(timezone.utc).isoformat(timespec="seconds") def _j(v) -> str: return json.dumps(v, ensure_ascii=False) # ------------------------------------------------------------- 프로그램 요약 def _include_lines(con: sqlite3.Connection, program: str) -> dict[str, list[str]]: return { r["include"]: (r["code"] or "").split("\n") for r in con.execute("SELECT include, code FROM include WHERE program=?", (program,)) } def _text_symbols(p: sqlite3.Row) -> dict[str, str]: return {t["symbol"].upper(): t["text"] for t in (loads(p["text_symbols_json"]) or []) if t.get("symbol")} def _program_context(con: sqlite3.Connection, p: sqlite3.Row, lines_by_inc: dict[str, list[str]]) -> dict: structure = _structure_summary(con, p["name"]) selection = [] for inc_lines in lines_by_inc.values(): for ln in inc_lines: m = _SELSCREEN.match(ln) if m: selection.append(ln.strip()[:100]) units = con.execute( "SELECT * FROM unit WHERE program=? AND unit_type IN ('FORM','METHOD','FUNCTION','MODULE','EVENT') " "ORDER BY include, line_start", (p["name"],)).fetchall() unit_list = [ f"{u['unit_type']} {u['name']} ({u['include']} L{u['line_start']}-{u['line_end']})" + (f" — {clean_comment(u['header_comment'])}" if clean_comment(u["header_comment"]) else "") for u in units[:150] ] event_code, budget = [], MAX_EVENT_CODE_LINES for u in units: if u["unit_type"] != "EVENT" or budget <= 0: continue lines = lines_by_inc.get(u["include"], []) end = min(u["line_end"], u["line_start"] + budget - 1) event_code.append(numbered_code(lines, u["line_start"], end)) budget -= end - u["line_start"] + 1 ts = _text_symbols(p) return { "structure": structure, "selection": "; ".join(selection[:30]) or "-", "text_symbols": "; ".join(f"{k}={v}" for k, v in list(ts.items())[:40]) or "-", "unit_list": "\n".join(unit_list) or "-", "event_code": "\n".join(event_code) or "-", } def ensure_program_summary(con: sqlite3.Connection, llm, p: sqlite3.Row, lines_by_inc: dict[str, list[str]], dry_run: bool) -> tuple[dict, str]: """프로그램 요약을 (필요하면) 생성하고 (summary dict, 상태) 를 돌려준다. 상태: done | skipped | failed | dry | pending(file 백엔드에서 응답 대기) """ existing = loads(p["summary_json"]) or {} if p["summary_status"] == "done" and existing.get("prompt_version") == PROMPT_VERSION: return existing, "skipped" ctx = _program_context(con, p, lines_by_inc) st = ctx["structure"] prompt = PROGRAM_PROMPT.format( program=p["name"], title=p["title_ko"] or "-", devclass=p["devclass"] or "-", pkg_text=p["pkg_text"] or "", selection=ctx["selection"], text_symbols=ctx["text_symbols"], main_flow="\n".join(st["main_flow"]) or "-", tables_read=_j(st["tables_read"][:40]), tables_write=_j(st["tables_write"][:40]), external_calls=_j(st["external_calls"][:40]), output_type=_j(st["output_type"]), unit_list=ctx["unit_list"], event_code=ctx["event_code"], ) if dry_run: return existing, "dry" try: raw = llm.complete_json(SYSTEM_PROMPT, prompt) except PendingResponse: # file 백엔드 — 프롬프트만 내놓고 응답을 기다린다. 실패로 기록하지 않는다. return existing, "pending" except Exception as e: # noqa: BLE001 — HTTP 오류·JSON 아닌 응답. 이 프로그램만 failed 로 두고 계속 con.execute("UPDATE program SET summary_status='failed' WHERE name=?", (p["name"],)) con.commit() print(f"[FAIL] program summary {p['name']}: {type(e).__name__}: {e}") return existing, "failed" try: raw["program"] = p["name"] summary = ProgramSummary.model_validate(raw) summary.tables_read = st["tables_read"] summary.tables_write = st["tables_write"] summary.external_calls = st["external_calls"] summary.output_type = summary.output_type or st["output_type"] summary.title_ko = summary.title_ko or p["title_ko"] or "" summary.prompt_version = PROMPT_VERSION con.execute("UPDATE program SET summary_json=?, summary_status='done' WHERE name=?", (summary.model_dump_json(), p["name"])) con.commit() return summary.model_dump(), "done" except Exception as e: # noqa: BLE001 con.execute("UPDATE program SET summary_status='failed' WHERE name=?", (p["name"],)) con.commit() print(f"[FAIL] program summary {p['name']}: {type(e).__name__}: {e}") return existing, "failed" # ------------------------------------------------------------- unit 조각 추출 def _windows(u: sqlite3.Row) -> list[tuple[int, int]]: if u["loc"] <= MAX_WINDOW_LINES: return [(u["line_start"], u["line_end"])] subs = loads(u["sub_chunks_json"]) or [] wins = [(s["line_start"], s["line_end"]) for s in subs if s["line_end"] >= s["line_start"]] if not wins: wins = [(a, min(a + FALLBACK_WINDOW - 1, u["line_end"])) for a in range(u["line_start"], u["line_end"] + 1, FALLBACK_WINDOW)] return wins def _known_symbols(con: sqlite3.Connection, program: str, unit_id: str) -> set[str]: """조각의 테이블·호출을 파서로 다시 뽑을 때 쓰는 '심볼이라 DB 테이블이 아니다' 집합. `TABLES:`/`NODES:` 로 선언된 DDIC 작업영역은 제외한다 — 넣으면 `SELECT ... FROM zfit0060` 이 심볼로 오인돼 조각의 tables_read 가 비어버린다 (parser/run.py 의 같은 가드와 기준을 맞춘다). """ return {r["name"] for r in con.execute( "SELECT name FROM symbol WHERE program=? AND (scope='global' OR unit_id=?) " "AND kind NOT IN ('tables','nodes')", (program, unit_id))} def _store_chunks(con: sqlite3.Connection, u: sqlite3.Row, chunks: list[ResolvedChunk], text_symbols: dict[str, str] | None = None) -> None: for r in con.execute("SELECT chunk_id FROM logic_chunk WHERE unit_id=?", (u["unit_id"],)).fetchall(): con.execute("DELETE FROM chunk_fts WHERE chunk_id=?", (r["chunk_id"],)) con.execute("DELETE FROM logic_chunk WHERE unit_id=?", (u["unit_id"],)) now = _now() for c in chunks: chunk_id = f"{u['unit_id']}#C{c.seq}" con.execute( "INSERT INTO logic_chunk(chunk_id, program, include, unit_id, seq, line_start, line_end, code_hash, " "kind, purpose_ko, purpose_en, keywords_ko, keywords_en, sap_objects, tables_read, tables_write, " "calls, confidence, prompt_version, extracted_at) VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)", (chunk_id, u["program"], u["include"], u["unit_id"], c.seq, c.line_start, c.line_end, c.code_hash, c.kind, c.purpose_ko, c.purpose_en, _j(c.keywords_ko), _j(c.keywords_en), _j(c.sap_objects), _j(c.tables_read), _j(c.tables_write), _j(c.calls), c.confidence, PROMPT_VERSION, now), ) # 조각 코드에 쓰인 TEXT-nnn 의 한국어 원문을 색인에 넣는다 (수정사항 3번) — # LLM 설명과 무관하게 화면 문구로도 조각에 도달할 수 있어야 한다. ts_phrases = text_symbol_phrases(c.code, text_symbols or {}) keywords = " ".join(c.keywords_ko + c.keywords_en + [c.kind] + ts_phrases) objects = " ".join(c.sap_objects + c.tables_read + c.tables_write + c.calls) con.execute( "INSERT INTO chunk_fts(chunk_id, program, purpose, keywords, objects, bigrams) VALUES(?,?,?,?,?,?)", (chunk_id, u["program"], f"{c.purpose_ko} {c.purpose_en}".strip(), keywords, objects, bigrams(f"{c.purpose_ko} {' '.join(c.keywords_ko)} {' '.join(ts_phrases)}")), ) def build_unit_prompts(con: sqlite3.Connection, u: sqlite3.Row, program_summary: dict, title: str, lines: list[str], text_symbols: dict[str, str]) -> list[dict]: """unit 의 창별 프롬프트를 만든다 (DB 읽기 — 메인 스레드 전용). LLM 호출을 병렬화하려면 '프롬프트 조립(DB) → 호출(병렬) → 저장(DB)' 으로 갈라야 한다. sqlite 커넥션은 스레드 간 공유가 안 되기 때문이다. """ windows = _windows(u) out: list[dict] = [] for wi, (ws, we) in enumerate(windows, 1): code = numbered_code(lines, ws, we) used_ts = {m.group(1).upper() for m in _TEXT_SYM.finditer(code)} ts_lines = [f"TEXT-{k} = {text_symbols[k]}" for k in sorted(used_ts) if k in text_symbols] prompt = UNIT_PROMPT.format( program=u["program"], title=title or "-", program_purpose=program_summary.get("business_purpose_ko") or "(요약 없음 — 구조 사실만 참고)", main_flow=" → ".join(program_summary.get("main_flow") or [])[:600] or "-", key_tables=", ".join(f"{t.get('name')}({t.get('desc_ko', '')})" for t in (program_summary.get("key_internal_tables") or [])[:10]) or "-", unit_id=u["unit_id"], unit_type=u["unit_type"], name=u["name"], signature=u["signature"] or "-", header_comment=clean_comment(u["header_comment"]) or "-", window_note=(f"[창] 이 단위는 길어서 {len(windows)}개 창으로 나눠 보여준다. 지금은 {wi}번째 창 " f"(L{ws}-L{we}). 이 창 안의 줄만 조각으로 만들라." if len(windows) > 1 else ""), hints="\n".join(parser_hints(lines, ws, we)) or "-", text_symbols="\n".join(ts_lines) or "-", code=code, kinds=", ".join(CHUNK_KINDS), ) out.append({"window": (ws, we), "prompt": prompt}) return out def finish_unit(con: sqlite3.Connection, u: sqlite3.Row, calls: list[dict], lines: list[str], text_symbols: dict[str, str]) -> dict: """창별 LLM 응답 → 검증·저장 (DB 쓰기 — 메인 스레드 전용). calls 원소: {window, prompt, raw?, error?, pending?} 한 창이라도 pending 이면 아무것도 저장하지 않는다 (부분 저장 방지). """ if any(c.get("pending") for c in calls): return {"status": "pending", "chunks": 0, "dropped": 0} errors = [c["error"] for c in calls if c.get("error")] if errors: raise errors[0] known = _known_symbols(con, u["program"], u["unit_id"]) all_chunks: list[ResolvedChunk] = [] dropped_all: list[str] = [] unit_purpose = "" for c in calls: ws, we = c["window"] ext = UnitExtraction.model_validate(c["raw"]) unit_purpose = unit_purpose or ext.unit_purpose_ko.strip() chunks, dropped = resolve_chunks(ext.chunks, lines, ws, we, u["include"], known) all_chunks.extend(chunks) dropped_all.extend(dropped) all_chunks.sort(key=lambda c: (c.line_start, c.line_end)) for i, c in enumerate(all_chunks, 1): c.seq = i _store_chunks(con, u, all_chunks, text_symbols) covered = covered_lines(all_chunks) thin = { "purpose_ko": unit_purpose or clean_comment(u["header_comment"]), "chunk_count": len(all_chunks), "covered_lines": covered, "coverage": round(covered / u["loc"], 3) if u["loc"] else 0.0, "dropped": dropped_all[:10], } _mark_unit_done(con, u, thin, len(all_chunks)) con.commit() return {"status": "done", "chunks": len(all_chunks), "dropped": len(dropped_all)} def _mark_unit_done(con: sqlite3.Connection, u: sqlite3.Row, thin: dict, n_chunks: int) -> None: con.execute( "UPDATE unit SET summary_json=?, summary_status='done', prompt_version=?, chunk_count=?, " "summary_error=NULL WHERE unit_id=?", (_j(thin), PROMPT_VERSION, n_chunks, u["unit_id"]), ) text = " ".join(filter(None, [u["name"], u["signature"] or "", thin["purpose_ko"]])) con.execute("UPDATE unit_fts SET purpose=?, bigrams=? WHERE unit_id=?", (thin["purpose_ko"], bigrams(text), u["unit_id"])) def clone_unit_chunks(con: sqlite3.Connection, rep: sqlite3.Row, target: sqlite3.Row) -> int: """code_hash 가 같은 unit 으로 조각을 복제한다 (수정사항 6번). _BAK / _COPY 관행 때문에 같은 코드가 여러 프로그램에 그대로 들어 있다. 대표 unit 1건만 LLM 에 보내고 나머지는 여기서 만든다 (실측 샘플: 중복 그룹 108개, 호출 179회 절감). 코드는 같아도 include 안에서의 **줄 위치는 다르다** — line_start 차이만큼 평행이동한다. """ shift = target["line_start"] - rep["line_start"] src = con.execute("SELECT * FROM logic_chunk WHERE unit_id=? ORDER BY seq", (rep["unit_id"],)).fetchall() for r in con.execute("SELECT chunk_id FROM logic_chunk WHERE unit_id=?", (target["unit_id"],)).fetchall(): con.execute("DELETE FROM chunk_fts WHERE chunk_id=?", (r["chunk_id"],)) con.execute("DELETE FROM logic_chunk WHERE unit_id=?", (target["unit_id"],)) for r in src: chunk_id = f"{target['unit_id']}#C{r['seq']}" con.execute( "INSERT INTO logic_chunk(chunk_id, program, include, unit_id, seq, line_start, line_end, code_hash, " "kind, purpose_ko, purpose_en, keywords_ko, keywords_en, sap_objects, tables_read, tables_write, " "calls, confidence, prompt_version, extracted_at) VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)", (chunk_id, target["program"], target["include"], target["unit_id"], r["seq"], r["line_start"] + shift, r["line_end"] + shift, r["code_hash"], r["kind"], r["purpose_ko"], r["purpose_en"], r["keywords_ko"], r["keywords_en"], r["sap_objects"], r["tables_read"], r["tables_write"], r["calls"], r["confidence"], PROMPT_VERSION, _now()), ) f = con.execute("SELECT purpose, keywords, objects, bigrams FROM chunk_fts WHERE chunk_id=?", (r["chunk_id"],)).fetchone() if f: con.execute( "INSERT INTO chunk_fts(chunk_id, program, purpose, keywords, objects, bigrams) " "VALUES(?,?,?,?,?,?)", (chunk_id, target["program"], f["purpose"], f["keywords"], f["objects"], f["bigrams"]), ) thin = dict(loads(rep["summary_json"]) or {}) thin["cloned_from"] = rep["unit_id"] _mark_unit_done(con, target, thin, len(src)) return len(src) # ------------------------------------------------------------------ 배치 진입점 def _model_label(fake: bool, backend: str | None) -> str: """llm_usage_log·대시보드에 남길 '무엇이 요약을 만들었나' 라벨. 환경변수의 LLM_MODEL 을 그대로 쓰면 안 된다 — file 백엔드로 사람·에이전트가 채운 결과를 쓰지도 않은 모델 이름으로 기록하게 된다(실측: file 백엔드 실행이 'z-ai/glm-5.2:free' 로 남았다). """ name = backend or ("fake" if fake else "api") if name == "fake": return "fake" if name == "file": return "file(사람·에이전트 응답)" return settings.llm_model def _llm_call(llm, system: str, prompt: str) -> dict: """스레드에서 도는 부분 — DB 를 건드리지 않는다. 예외는 값으로 돌려준다.""" try: return {"raw": llm.complete_json(system, prompt)} except PendingResponse as e: return {"pending": True, "job_id": e.job_id} except Exception as e: # noqa: BLE001 — 호출 실패는 unit 단위로 기록하고 계속 return {"error": e} def _progress(i: int, n: int, t0: float, c: dict) -> None: """호출 하나가 끝날 때마다 한 줄 — 긴 프로그램에서 멈춘 것처럼 보이지 않게 (stderr).""" state = "실패" if "error" in c else ("대기" if c.get("pending") else "ok") elapsed = time.time() - t0 eta = elapsed / i * (n - i) if i else 0 print(f" [{i}/{n}] {state} {elapsed:.0f}s 경과, 약 {eta:.0f}s 남음", file=sys.stderr, flush=True) def _run_calls_parallel(llm, calls: list[dict], concurrency: int, on_done=None) -> None: """calls 각 원소의 'prompt' 를 호출하고 결과를 그 자리에 채운다 (in-place). on_done(call) 은 호출 하나가 끝날 때마다 **메인 스레드에서** 불린다 — 여기서 DB 에 저장하면 중간에 끊어도 그때까지 끝난 unit 은 남는다 (고객사 요청 2026-09-21). Ctrl+C 는 남은 호출을 취소하고 바로 올라간다 (진행 중인 호출은 응답 상한까지 기다릴 수 있다). """ if not calls: return t0 = time.time() n = len(calls) if concurrency <= 1 or n == 1: for i, c in enumerate(calls, 1): c.update(_llm_call(llm, SYSTEM_PROMPT, c["prompt"])) _progress(i, n, t0, c) if on_done: on_done(c) return ex = ThreadPoolExecutor(max_workers=min(concurrency, n)) try: futures = {ex.submit(_llm_call, llm, SYSTEM_PROMPT, c["prompt"]): c for c in calls} for i, fut in enumerate(as_completed(futures), 1): c = futures[fut] c.update(fut.result()) _progress(i, n, t0, c) if on_done: on_done(c) except KeyboardInterrupt: print("\n중단 요청 — 남은 호출을 취소한다 (이미 끝난 unit 은 저장됨)", file=sys.stderr, flush=True) ex.shutdown(wait=False, cancel_futures=True) raise ex.shutdown(wait=True) def _dedupe_plan(con: sqlite3.Connection, pending: list[sqlite3.Row]) -> tuple[list[sqlite3.Row], dict]: """(LLM 에 보낼 대표 unit 목록, {대표 unit_id: [복제 대상 unit...]}) (수정사항 6번) 이미 done 이고 조각이 있는 unit 과 code_hash 가 같으면 LLM 호출 없이 바로 복제한다. **배치 전체를 한 번에 넘겨야 한다.** 프로그램별로 나눠 호출하면 공용 인클루드 (ZFICOM/ZFIALV 처럼 17개 프로그램에 그대로 복사된 코드)가 서로 다른 호출에 흩어져 중복이 잡히지 않는다 — 실측에서 절감 0회가 나왔다. 프로그램 안의 중복은 드물고 절감분은 거의 전부 프로그램을 가로지르는 중복이다. """ reps: list[sqlite3.Row] = [] followers: dict[str, list[sqlite3.Row]] = {} rep_by_hash: dict[str, sqlite3.Row] = {} for u in pending: h = u["code_hash"] if not h: reps.append(u) continue if h in rep_by_hash: followers.setdefault(rep_by_hash[h]["unit_id"], []).append(u) continue # 이번 배치 밖에서 이미 추출된 동일 코드 unit 이 있으면 그걸 대표로 쓴다 (호출 0회) done = con.execute( "SELECT * FROM unit WHERE code_hash=? AND unit_id!=? AND summary_status='done' " "AND prompt_version=? AND chunk_count>0 LIMIT 1", (h, u["unit_id"], PROMPT_VERSION), ).fetchone() if done: followers.setdefault(done["unit_id"], []).append(u) rep_by_hash[h] = done continue rep_by_hash[h] = u reps.append(u) return reps, followers def summarize_units(program: str | None, limit: int | None, fake: bool = False, dry_run: bool = False, trigger: str = "manual", backend: str | None = None, concurrency: int | None = None, dedupe: bool = True) -> dict: """프로그램(들)의 요약 + unit 조각 추출을 실행한다. 이름은 API 호환을 위해 유지.""" con = connect() llm = create_llm(fake=fake, backend=backend) conc = max(1, concurrency or settings.llm_concurrency) stats = {"programs": 0, "program_summaries": 0, "done": 0, "skipped": 0, "failed": 0, "cloned": 0, "awaiting_response": 0, "chunks": 0, "dropped": 0, "llm_calls_saved": 0, "elapsed_s": 0} t0 = time.time() run_id: int | None = None try: sql = ("SELECT p.*, COALESCE(k.text_ko,'') AS pkg_text FROM program p " "LEFT JOIN package k ON k.devclass=p.devclass WHERE p.has_source=1") params: list = [] if program: sql += " AND p.name=?" params.append(program.upper()) programs = con.execute(sql + " ORDER BY p.name", params).fetchall() # 대상 unit 수집 (leaf 부터 — topo 역순) pending: list[sqlite3.Row] = [] for p in programs: rows = con.execute( "SELECT u.*, t.ord FROM unit u LEFT JOIN topo t ON t.unit_id=u.unit_id " "WHERE u.program=? AND u.unit_type IN ('FORM','METHOD','FUNCTION','MODULE','EVENT') " "ORDER BY COALESCE(t.ord, 9999) DESC", (p["name"],)).fetchall() for u in rows: if u["summary_status"] == "done" and u["prompt_version"] == PROMPT_VERSION: stats["skipped"] += 1 else: pending.append(u) if limit: pending = pending[:limit] # 주의: "pending" = 이번 실행의 대상 unit 수(기존 의미, API/대시보드가 사용). # file 백엔드의 '응답 대기' 는 별 키인 "awaiting_response" 다. stats["pending"] = len(pending) # 중복 제거는 **배치 전체**를 대상으로 한 번만 세운다 (프로그램별로 나누면 공용 # 인클루드의 교차 중복이 잡히지 않는다). LLM 은 대표 unit 만 보고, 나머지는 복제한다. reps, followers = _dedupe_plan(con, pending) if dedupe else (pending, {}) stats["llm_calls_saved"] = sum(len(v) for v in followers.values()) pending_by_prog: dict[str, list[sqlite3.Row]] = {} for u in reps: pending_by_prog.setdefault(u["program"], []).append(u) needs_summary = {p["name"] for p in programs if not ( p["summary_status"] == "done" and (loads(p["summary_json"]) or {}).get("prompt_version") == PROMPT_VERSION)} if (pending or needs_summary) and not dry_run: cur = con.execute( "INSERT INTO llm_usage_log(ts, program, model, trigger_by, status, total, " "calls, prompt_tokens, completion_tokens, cost_usd, done, failed, skipped, elapsed_s, chunks) " "VALUES(datetime('now','localtime'),?,?,?,'running',?,0,0,0,0.0,0,0,?,0,0)", (program or "*", _model_label(fake, backend), trigger, len(pending), stats["skipped"]), ) run_id = cur.lastrowid con.commit() for p in programs: units = pending_by_prog.get(p["name"], []) if not units and p["name"] not in needs_summary: continue stats["programs"] += 1 lines_by_inc = _include_lines(con, p["name"]) summary, st = ensure_program_summary(con, llm, p, lines_by_inc, dry_run) if st == "done": stats["program_summaries"] += 1 elif st == "pending": # 프로그램 요약은 unit 프롬프트의 **문맥**이다 (2단계 설계). 요약이 아직 없는데 # unit 프롬프트를 내놓으면 "(요약 없음)" 으로 만들어진 프롬프트를 받게 되고, # 요약이 채워진 다음 패스에서 프롬프트 내용이 달라져 전부 다시 답해야 한다. # 그래서 요약이 대기 중이면 이 프로그램의 unit 은 이번 패스에서 건너뛴다. stats["awaiting_response"] += 1 stats["blocked_units"] = stats.get("blocked_units", 0) + len(units) continue text_symbols = _text_symbols(p) if dry_run: continue # 1) 프롬프트 조립 (DB 읽기, 메인 스레드) — units 는 이미 대표 unit 만 들어 있다 jobs: list[tuple[sqlite3.Row, list[dict]]] = [] for u in units: prompts = build_unit_prompts(con, u, summary, p["title_ko"] or "", lines_by_inc.get(u["include"], []), text_symbols) jobs.append((u, [dict(x) for x in prompts])) # 2) LLM 호출 (병렬, DB 접근 없음) — 수정사항 10번 # + 3) 검증·저장: unit 의 호출이 모두 끝나는 즉시 메인 스레드에서 저장 (중간에 끊어도 보존) all_calls = [c for _, calls in jobs for c in calls] print(f"[{p['name']}] unit {len(units)}건 → LLM 호출 {len(all_calls)}건 (동시 {conc})", file=sys.stderr, flush=True) def _save_unit(u: sqlite3.Row, calls: list[dict]) -> None: lines = lines_by_inc.get(u["include"], []) try: r = finish_unit(con, u, calls, lines, text_symbols) except Exception as e: # noqa: BLE001 err = f"{type(e).__name__}: {e}"[:500] con.execute("UPDATE unit SET summary_status='failed', summary_error=? WHERE unit_id=?", (err, u["unit_id"])) con.commit() stats["failed"] += 1 print(f"[FAIL] {u['unit_id']}: {err}") return if r["status"] == "pending": stats["awaiting_response"] += 1 return stats["done"] += 1 stats["chunks"] += r["chunks"] stats["dropped"] += r["dropped"] con.commit() owner = {id(c): (u, calls) for u, calls in jobs for c in calls} remaining = {u["unit_id"]: len(calls) for u, calls in jobs} def _on_call_done(c: dict) -> None: u, calls = owner[id(c)] remaining[u["unit_id"]] -= 1 if remaining[u["unit_id"]] == 0: # 창이 여러 개인 unit 은 마지막 창까지 기다린다 _save_unit(u, calls) _run_calls_parallel(llm, all_calls, conc, on_done=_on_call_done) if run_id is not None: _update_run(con, run_id, stats, llm, t0, status="running") # ---- 복제 패스: code_hash 가 같은 나머지 unit 에 조각을 복제한다 (LLM 호출 없음) ---- # 대표와 복제 대상이 서로 다른 프로그램일 수 있어 프로그램 루프가 끝난 뒤에 돈다. touched_programs = {p["name"] for p in programs if pending_by_prog.get(p["name"])} if not dry_run: for rep_id, targets in followers.items(): rep_row = con.execute("SELECT * FROM unit WHERE unit_id=?", (rep_id,)).fetchone() if not rep_row or rep_row["summary_status"] != "done": continue # 대표가 아직 안 끝났다(응답 대기·실패) — 다음 실행에서 복제된다 for tgt in targets: stats["chunks"] += clone_unit_chunks(con, rep_row, tgt) stats["cloned"] += 1 touched_programs.add(tgt["program"]) con.commit() # 요약·조각이 생긴 프로그램의 색인을 다시 만든다 (수정사항 1번) — 단건 경로 for name in sorted(touched_programs): refresh_program_fts(con, name) con.commit() finally: if run_id is not None: stats["elapsed_s"] = round(time.time() - t0, 1) _update_run(con, run_id, stats, llm, t0, status="done") con.close() if hasattr(llm, "usage"): stats["llm_usage"] = llm.usage stats["elapsed_s"] = round(time.time() - t0, 1) return stats def _update_run(con, run_id: int, stats: dict, llm, t0: float, status: str) -> None: u = getattr(llm, "usage", None) or {} con.execute( "UPDATE llm_usage_log SET status=?, calls=?, prompt_tokens=?, completion_tokens=?, " "cost_usd=?, done=?, failed=?, skipped=?, elapsed_s=?, chunks=? WHERE id=?", (status, u.get("calls", 0), u.get("prompt_tokens", 0), u.get("completion_tokens", 0), u.get("cost_usd", 0.0), stats["done"], stats["failed"], stats["skipped"], round(time.time() - t0, 1), stats["chunks"], run_id), ) con.commit() def main() -> None: ap = argparse.ArgumentParser( description="Stage 3 — 프로그램 요약 + 로직 조각 추출", epilog="LLM 키가 없으면: --llm file 로 프롬프트를 파일로 내놓고 " "python -m summarize.jobs 로 응답을 채운 뒤 같은 명령을 다시 실행한다.", ) ap.add_argument("--program", default=None) ap.add_argument("--limit", type=int, default=None, help="처리할 unit 수 상한") ap.add_argument("--llm", choices=["api", "file", "fake"], default=None, help="api=환경변수 키로 호출 / file=프롬프트 파일 큐(키 불필요) / fake=더미") ap.add_argument("--fake", action="store_true", help="--llm fake 의 하위호환 별칭") ap.add_argument("--concurrency", type=int, default=None, help=f"동시 LLM 호출 수 (기본 LLM_CONCURRENCY={settings.llm_concurrency})") ap.add_argument("--no-dedupe", action="store_true", help="code_hash 가 같은 unit 도 각각 LLM 에 보낸다 (기본은 대표 1건만)") ap.add_argument("--dry-run", action="store_true") args = ap.parse_args() stats = summarize_units( args.program, args.limit, fake=args.fake, dry_run=args.dry_run, backend=args.llm, concurrency=args.concurrency, dedupe=not args.no_dedupe, ) print(json.dumps(stats, ensure_ascii=False)) if stats.get("pending"): print(f"\n응답 대기 {stats['awaiting_response']}건 — 프롬프트: {settings.data_llm_jobs}\n" f" python -m summarize.jobs list --pending\n" f" python -m summarize.jobs show \n" f" python -m summarize.jobs answer --file <응답.json>\n" f" → 응답을 채운 뒤 같은 명령을 다시 실행하면 적재된다.") if stats.get("blocked_units"): print(f"\nunit {stats['blocked_units']}건은 프로그램 요약을 문맥으로 쓰므로 대기 중이다 — " f"먼저 program_summary 작업에 답하고 다시 실행하면 unit 프롬프트가 나온다.") if __name__ == "__main__": main()