Files

205 lines
7.1 KiB
Python

"""LLM 작업 큐 CLI — LLM 키 없이 요약/조각 추출을 돌리는 입구.
`--llm file` 로 러너를 돌리면 프롬프트가 `data/llm_jobs/` 에 파일로 쌓인다.
사람이든 코딩 에이전트든 그 프롬프트를 읽고 응답 JSON 을 채워 넣으면, 러너를 다시 돌릴 때
그 응답으로 파이프라인이 이어진다.
# 1) 프롬프트 내놓기 (LLM 호출 0회)
python -m summarize.runner --llm file --program ZFIR10070 --limit 3
# 2) 대기 목록 보기 / 프롬프트 읽기
python -m summarize.jobs list
python -m summarize.jobs show <job_id>
# 3) 응답 채우기 (스키마 검증 후 저장)
python -m summarize.jobs answer <job_id> --file answer.json
python -m summarize.jobs answer <job_id> --stdin < answer.json
# 4) 같은 명령을 다시 — 이번엔 응답을 읽어 조각을 적재한다
python -m summarize.runner --llm file --program ZFIR10070 --limit 3
`answer` 는 저장 전에 프롬프트의 [작업] 종류에 맞는 pydantic 스키마로 검증한다 —
잘못된 JSON 을 큐에 넣어두고 나중에 러너에서 실패하는 걸 막는다.
"""
from __future__ import annotations
import argparse
import json
import re
import sys
from pathlib import Path
from config.settings import settings
from .llm_client import JOB_INDEX
from .schemas import ProgramSummary, UnitExtraction
TASK_SCHEMA = {"program_summary": ProgramSummary, "extract_chunks": UnitExtraction}
def jobs_dir(override: str | None = None) -> Path:
return Path(override) if override else Path(settings.data_llm_jobs)
def _index_rows(d: Path) -> list[dict]:
idx = d / JOB_INDEX
if not idx.exists():
return []
rows = []
for line in idx.read_text(encoding="utf-8").splitlines():
if line.strip():
try:
rows.append(json.loads(line))
except json.JSONDecodeError:
continue
return rows
def _status(d: Path, row: dict) -> str:
return "answered" if (d / f"{row['job_id']}.response.json").exists() else "pending"
def task_of(d: Path, jid: str) -> str:
"""프롬프트 파일에서 [작업] 종류를 읽는다."""
p = d / f"{jid}.prompt.md"
if not p.exists():
return ""
m = re.search(r"^\[작업\]\s*(\S+)", p.read_text(encoding="utf-8"), re.M)
return m.group(1) if m else ""
def cmd_list(args) -> int:
d = jobs_dir(args.dir)
rows = _index_rows(d)
if not rows:
print(f"작업이 없습니다: {d}\n먼저 실행: python -m summarize.runner --llm file --program <PROG> --limit 3")
return 0
shown = 0
for r in rows:
st = _status(d, r)
if args.pending and st != "pending":
continue
print(f"{r['job_id']} {st:9s} {r.get('task',''):16s} {r.get('program','')} "
f"{r.get('unit_id','') or '(프로그램 요약)'}"
+ (f" 창{r['window']}" if r.get("window") else ""))
shown += 1
n_pending = sum(1 for r in rows if _status(d, r) == "pending")
print(f"\n{len(rows)}건 (표시 {shown}) — 대기 {n_pending}, 응답완료 {len(rows) - n_pending}")
return 0
def cmd_show(args) -> int:
d = jobs_dir(args.dir)
p = d / f"{args.job_id}.prompt.md"
if not p.exists():
print(f"프롬프트 없음: {p}", file=sys.stderr)
return 1
sys.stdout.write(p.read_text(encoding="utf-8"))
return 0
def cmd_answer(args) -> int:
d = jobs_dir(args.dir)
prompt = d / f"{args.job_id}.prompt.md"
if not prompt.exists():
print(f"프롬프트 없음: {prompt}", file=sys.stderr)
return 1
if args.stdin:
raw = sys.stdin.read()
else:
try:
raw = Path(args.file).read_text(encoding="utf-8")
except OSError as e:
print(f"응답 파일을 읽을 수 없습니다: {e}", file=sys.stderr)
return 1
try:
data = json.loads(raw)
except json.JSONDecodeError as e:
s, e2 = raw.find("{"), raw.rfind("}")
if s < 0 or e2 <= s:
print(f"JSON 파싱 실패: {e}", file=sys.stderr)
return 1
data = json.loads(raw[s : e2 + 1])
task = task_of(d, args.job_id)
schema = TASK_SCHEMA.get(task)
if schema and not args.no_validate:
try:
schema.model_validate(data)
except Exception as e: # noqa: BLE001
print(f"[스키마 불일치] task={task}\n{e}", file=sys.stderr)
print("\n무시하고 저장하려면 --no-validate", file=sys.stderr)
return 1
out = d / f"{args.job_id}.response.json"
out.write_text(json.dumps(data, ensure_ascii=False, indent=2), encoding="utf-8")
print(f"저장: {out} (task={task or '?'}, 검증={'skip' if args.no_validate else 'ok'})")
n_pending = sum(1 for r in _index_rows(d) if _status(d, r) == "pending")
print(f"남은 대기: {n_pending}건 — 전부 채우면 러너를 같은 옵션으로 다시 실행하세요")
return 0
def cmd_stats(args) -> int:
d = jobs_dir(args.dir)
rows = _index_rows(d)
by_task: dict[str, list[int]] = {}
for r in rows:
t = r.get("task") or "?"
b = by_task.setdefault(t, [0, 0])
b[0 if _status(d, r) == "pending" else 1] += 1
print(json.dumps({
"dir": str(d),
"total": len(rows),
"pending": sum(b[0] for b in by_task.values()),
"answered": sum(b[1] for b in by_task.values()),
"by_task": {k: {"pending": v[0], "answered": v[1]} for k, v in sorted(by_task.items())},
}, ensure_ascii=False, indent=2))
return 0
def cmd_clear(args) -> int:
d = jobs_dir(args.dir)
n = 0
for p in list(d.glob("*.prompt.md")) + list(d.glob("*.response.json")) + [d / JOB_INDEX]:
if p.exists() and (not args.answered_only or p.name.endswith(".response.json")):
p.unlink()
n += 1
print(f"삭제 {n}개 파일 ({d})")
return 0
def main() -> None:
ap = argparse.ArgumentParser(description="LLM 작업 큐 (키 없이 요약 돌리기)")
ap.add_argument("--dir", default=None, help=f"작업 디렉토리 (기본 {settings.data_llm_jobs})")
sub = ap.add_subparsers(dest="cmd", required=True)
p = sub.add_parser("list", help="작업 목록")
p.add_argument("--pending", action="store_true", help="대기 중인 것만")
p.set_defaults(func=cmd_list)
p = sub.add_parser("show", help="프롬프트 출력")
p.add_argument("job_id")
p.set_defaults(func=cmd_show)
p = sub.add_parser("answer", help="응답 JSON 저장 (스키마 검증)")
p.add_argument("job_id")
g = p.add_mutually_exclusive_group(required=True)
g.add_argument("--file", help="응답 JSON 파일")
g.add_argument("--stdin", action="store_true", help="표준입력에서 읽기")
p.add_argument("--no-validate", action="store_true")
p.set_defaults(func=cmd_answer)
p = sub.add_parser("stats", help="큐 통계")
p.set_defaults(func=cmd_stats)
p = sub.add_parser("clear", help="큐 비우기")
p.add_argument("--answered-only", action="store_true")
p.set_defaults(func=cmd_clear)
args = ap.parse_args()
sys.exit(args.func(args))
if __name__ == "__main__":
main()