diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 2ddff42..3ffd301 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -31,10 +31,24 @@ jobs: - name: Install run: | python -m pip install --upgrade pip - pip install -e ".[dev]" + pip install -e ".[dev,agent]" - name: Ruff run: | ruff check . ruff format --check . || echo "::warning::ruff format differences (run 'ruff format .')" - name: Pytest run: pytest -q + - name: Example flow (offline) + run: make flow CASE=cases/_example_night_clinic + - name: Upload flow artifacts + if: always() + uses: actions/upload-artifact@v4 + with: + name: example-flow + path: | + cases/_example_night_clinic/flow_log.json + cases/_example_night_clinic/run_manifest.json + cases/_example_night_clinic/report.md + cases/_example_night_clinic/figures/ + if-no-files-found: warn + retention-days: 14 diff --git a/Makefile b/Makefile index 418f7e6..e8eb75f 100644 --- a/Makefile +++ b/Makefile @@ -1,4 +1,4 @@ -.PHONY: install check lint test format app activity +.PHONY: install check lint test format app activity flow PY ?= python @@ -19,5 +19,10 @@ format: app: streamlit run app/streamlit_app.py +# Run the 6-step agent flow on a case (demo: skips the plan pre-registration gate). +CASE ?= cases/_example_night_clinic +flow: + $(PY) -m core.agent $(CASE) --allow-uncommitted + activity: $(PY) scripts/weekly_activity.py --days 7 diff --git a/README.md b/README.md index 9a69322..67a8548 100644 --- a/README.md +++ b/README.md @@ -24,7 +24,7 @@ cases/ _example_*/ 참고용 예시 케이스 <조-주제>/ 조별 케이스 (plan.yaml, fetch.py, estimate.py, report.md, figures/) app/streamlit_app.py 케이스 브라우저 (API 키 없이 실행) -docs/ops/ GitHub 온보딩, 모니터링 가이드 +docs/ops/ 조별 운영 가이드, GitHub 온보딩, 모니터링 docs/strategy/ 문제 정의·전략 문서 scripts/ 운영 스크립트 (weekly_activity.py 등) tests/ 테스트 @@ -60,7 +60,7 @@ cp -r cases/_template cases/group3-youth-rent # 폴더명: <조>-<주제>, 소 3. **추정** — `estimate.py`: `core.estimators`로 효과 추정 + 반증(placebo 등) → `figures/*.png` 4. **리포트** — `report.md`: 결과·한계·정책 시사점. `make app`에서 바로 보입니다. -자세한 절차: [`cases/_template/README.md`](cases/_template/README.md), 협업 규칙: [`CONTRIBUTING.md`](CONTRIBUTING.md), GitHub가 처음이라면: [`docs/ops/github-onboarding.md`](docs/ops/github-onboarding.md) +자세한 절차: [`cases/_template/README.md`](cases/_template/README.md), 협업 규칙: [`CONTRIBUTING.md`](CONTRIBUTING.md), GitHub가 처음이라면: [`docs/ops/github-onboarding.md`](docs/ops/github-onboarding.md), 조별 운영: [`docs/ops/group-guide.md`](docs/ops/group-guide.md) ## 7주 로드맵 @@ -74,6 +74,25 @@ cp -r cases/_template cases/group3-youth-rent # 폴더명: <조>-<주제>, 소 | 6 | 리포트·플랫폼 | `report.md`, Streamlit 반영 | | 7 | 발표·회고·공개 정리 | 최종 PR 머지, 릴리스 태그 | +## Flow — 6단계 에이전트 흐름 + +`core.agent`가 케이스 하나를 아래 6단계로 실행하고 `cases/<케이스>/flow_log.json`에 단계별 결과를 남깁니다. 앱의 **Flow** 페이지(`make app` → 사이드바 Flow)에서 실행하거나 저장된 로그를 볼 수 있습니다. + +| 단계 | 하는 일 | 주차 | +|---|---|---| +| ① 문제 정의 | `plan.yaml` 검증 + 사전 등록(커밋) 확인 — 추정 전에 확정 | W3 | +| ② 데이터 수집 | 공공데이터 수집·출처/라이선스 기록 | W2 | +| ③ 지표 구조화 | 패널 구성·품질 점검 | W3 | +| ④ 효과 추정 | DiD/이벤트 스터디/ITS + 반증 → 식별됨·조건부·식별 불가 | W4–5 | +| ⑤ 과잉해석 가드 | 결론 보류 규칙 + 인과 단정 표현 검사 | W6 | +| ⑥ 리포트 | `report.md`·그림·재현 기록 | W7 | + +```bash +make flow CASE=cases/_example_night_clinic # = python -m core.agent <케이스> --allow-uncommitted +``` + +`--allow-uncommitted`는 데모용입니다. 실제 분석은 `plan.yaml`을 먼저 커밋한 뒤 플래그 없이 실행하세요. LLM 서술은 `.env`에 키를 넣고 `--llm`으로 켭니다(`uv pip install -e ".[agent]"`). + ## 라이선스 - **코드**: MIT ([LICENSE](LICENSE)) — © 가짜연구소 Causal Inference Team diff --git a/app/flow_view.py b/app/flow_view.py new file mode 100644 index 0000000..238cbae --- /dev/null +++ b/app/flow_view.py @@ -0,0 +1,277 @@ +"""Flow view: normalize and render core.agent's flow_log.json (schema_version "1"). + +The adapter (`normalize_flow_log`) is the only place that knows the log layout, so a +schema change in core/agent/state.py stays a local fix here. +""" + +from __future__ import annotations + +import json +import os +from pathlib import Path +from typing import Any + +import yaml + +ROOT = Path(__file__).resolve().parents[1] +FLOW_LOG = "flow_log.json" + +# Mirrors core.agent.state.STEPS (copied so the app renders logs even without core). +STEPS: list[tuple[str, str]] = [ + ("define_problem", "① 문제 정의"), + ("collect", "② 데이터 수집"), + ("structure_metrics", "③ 지표 구조화·품질 점검"), + ("estimate", "④ 효과 추정"), + ("guard", "⑤ 과잉해석 가드"), + ("report", "⑥ 리포트·재현 기록"), +] +WEEKS = { # program week each step is taught (README "Flow" table) + "define_problem": "W3", + "collect": "W2", + "structure_metrics": "W3", + "estimate": "W4–5", + "guard": "W6", + "report": "W7", +} +STATUS_BADGE = { + "ok": ("OK", "green"), + "warn": ("WARN", "orange"), + "failed": ("FAILED", "red"), + "needs_human": ("NEEDS HUMAN", "violet"), + "blocked": ("BLOCKED", "red"), + "running": ("RUNNING", "blue"), + "skipped": ("SKIPPED", "gray"), +} +VERDICT_LABEL = { + "identified": ("식별됨", "green"), + "conditional": ("조건부", "orange"), + "not_identified": ("식별 불가", "red"), +} +LLM_KEY_ENVS = ("OPENAI_API_KEY", "ANTHROPIC_API_KEY", "OPENAI_BASE_URL") + + +def cases_dir() -> Path: + """Cases root; PEA_CASES_DIR overrides it (used by tests).""" + return Path(os.environ.get("PEA_CASES_DIR") or ROOT / "cases") + + +def list_cases(root: Path | None = None) -> list[Path]: + root = root or cases_dir() + if not root.is_dir(): + return [] + return sorted(p for p in root.iterdir() if p.is_dir() and (p / "plan.yaml").exists()) + + +def llm_available() -> bool: + return any(os.environ.get(k) for k in LLM_KEY_ENVS) + + +def badge(status: str | None) -> str: + text, color = STATUS_BADGE.get(status or "skipped", (str(status).upper(), "gray")) + return f":{color}-background[{text}]" + + +def verdict_badge(verdict: str | None) -> str: + if not verdict: + return ":gray-background[판정 없음]" + text, color = VERDICT_LABEL.get(verdict, (verdict, "gray")) + return f":{color}-background[**{text}**]" + + +def normalize_flow_log(raw: dict[str, Any]) -> dict[str, Any]: + """Return a log with exactly six steps in order; missing steps become 'skipped'.""" + by_name = {s.get("step"): s for s in raw.get("steps") or [] if isinstance(s, dict)} + steps = [] + for i, (name, title) in enumerate(STEPS, start=1): + s = dict(by_name.get(name) or {}) + s.setdefault("index", i) + s.setdefault("step", name) + s.setdefault("title", title) + s.setdefault("status", "skipped") + s.setdefault("message", "") + s["artifacts"] = s.get("artifacts") or {} + steps.append(s) + return { + "schema_version": str(raw.get("schema_version", "1")), + "case_dir": raw.get("case_dir"), + "question": raw.get("question"), + "status": raw.get("status"), + "verdict": raw.get("verdict"), + "started_at": raw.get("started_at"), + "ended_at": raw.get("ended_at"), + "steps": steps, + "quality": raw.get("quality") or [], + "results": raw.get("results") or [], + "guard": raw.get("guard") or {}, + "report_path": raw.get("report_path"), + } + + +def load_flow_log(case: Path) -> dict[str, Any] | None: + path = case / FLOW_LOG + if not path.exists(): + return None + return normalize_flow_log(json.loads(path.read_text(encoding="utf-8"))) + + +def load_plan_dict(case: Path) -> dict[str, Any]: + try: + data = yaml.safe_load((case / "plan.yaml").read_text(encoding="utf-8")) + except (OSError, yaml.YAMLError): + return {} + return data if isinstance(data, dict) else {} + + +def run_flow_for(case: Path, allow_uncommitted: bool, use_llm: bool) -> dict[str, Any]: + """Run core.agent's flow and return the normalized log.""" + from core.agent import run_flow # imported lazily: the app works without the agent + + state = run_flow(case, allow_uncommitted=allow_uncommitted, use_llm=use_llm) + return normalize_flow_log(state.to_log() if hasattr(state, "to_log") else state) + + +def _resolve(case: Path, rel: str | None) -> Path | None: + if not rel: + return None + p = Path(rel) + if not p.is_absolute(): + p = case / p + return p if p.exists() else None + + +# ---------- rendering ---------- + + +def _plan_summary(st: Any, plan: dict[str, Any], step: dict[str, Any]) -> None: + st.markdown(f"**질문** {plan.get('question', '-')}") + t = plan.get("treatment") or {} + c = plan.get("control") or {} + est = plan.get("estimator") or {} + st.markdown( + f"- 처치: {t.get('definition', '-') if isinstance(t, dict) else t}\n" + f"- 대조: {c.get('definition', '-') if isinstance(c, dict) else c}\n" + f"- 추정: `{est.get('method', '-') if isinstance(est, dict) else est}`\n" + f"- 사전 등록(커밋): {'예' if step['artifacts'].get('committed') else '아니오'}" + ) + + +def _collect(st: Any, plan: dict[str, Any], step: dict[str, Any]) -> None: + a = step["artifacts"] + rows = [ + { + "데이터": s.get("name"), + "제공": s.get("provider"), + "라이선스": s.get("license"), + "URL": s.get("url") or "", + } + for s in plan.get("data_sources") or [] + if isinstance(s, dict) + ] + if rows: + st.dataframe(rows, hide_index=True, width="stretch") + if a.get("licenses"): + st.caption("라이선스: " + ", ".join(map(str, a["licenses"]))) + if a.get("method"): + st.caption(f"수집 방식: {a['method']}") + + +def _quality(st: Any, log: dict[str, Any], step: dict[str, Any]) -> None: + if log["quality"]: + rows = [ + { + "상태": q.get("status"), + "점검": q.get("check"), + "값": str(q.get("value")), + "메시지": q.get("message", ""), + } + for q in log["quality"] + ] + st.dataframe(rows, hide_index=True, width="stretch") + a = step["artifacts"] + if a: + st.caption(" · ".join(f"{k}={v}" for k, v in a.items())) + + +def _fmt(x: Any) -> str: + return f"{x:.3f}" if isinstance(x, int | float) else "-" + + +def _estimate(st: Any, log: dict[str, Any]) -> None: + for r in log["results"]: + st.markdown( + f"**{r.get('method')}** · `{r.get('outcome')}` → {verdict_badge(r.get('verdict'))}" + ) + cols = st.columns(3) + cols[0].metric("추정치", _fmt(r.get("estimate"))) + cols[1].metric("95% CI", f"[{_fmt(r.get('ci_low'))}, {_fmt(r.get('ci_high'))}]") + cols[2].metric("p", _fmt(r.get("p_value"))) + checks = r.get("assumptions_checked") or {} + for name, chk in checks.items(): + if isinstance(chk, dict) and "passed" in chk: + mark = "통과" if chk["passed"] else "실패" + st.caption(f"{name}: {mark} (p={_fmt(chk.get('p_value'))})") + for w in r.get("warnings") or []: + st.warning(w) + + +def _guard(st: Any, log: dict[str, Any]) -> None: + g = log["guard"] + violations = g.get("violations") or [] + if violations: + st.error(f"과잉해석 표현 {len(violations)}건 발견") + for v in violations: + st.markdown(f"- **{v.get('phrase')}** — {v.get('reason')} \n > {v.get('sentence')}") + elif g: + st.success("과잉해석 표현 없음") + if g.get("narrative"): + st.info(g["narrative"]) + st.caption(f"서술 출처: {g.get('narrative_source', '-')}") + + +def _report(st: Any, case: Path, log: dict[str, Any], step: dict[str, Any]) -> None: + a = step["artifacts"] + # report_path may be absolute on the machine that ran the flow; resolve by name instead. + report_name = Path(log["report_path"]).name if log.get("report_path") else None + report = _resolve(case, a.get("report")) or _resolve(case, report_name) + for fig in a.get("figures") or []: + p = _resolve(case, fig) + if p: + st.image(str(p), caption=p.name) + if report: + with st.expander("report.md 미리보기", expanded=False): + st.markdown(report.read_text(encoding="utf-8")) + + +def render_flow(st: Any, case: Path, log: dict[str, Any], plan: dict[str, Any]) -> None: + if plan.get("synthetic_data"): + st.error("## ⚠ SYNTHETIC DATA\n합성 데이터입니다. 결과는 실제 정책 효과가 아닙니다.") + + top = st.columns(3) + top[0].markdown(f"**전체 상태** {badge(log.get('status'))}") + top[1].markdown(f"**판정** {verdict_badge(log.get('verdict'))}") + top[2].caption(f"{log.get('started_at') or '-'} → {log.get('ended_at') or '-'}") + + for step in log["steps"]: + name = step["step"] + dur = step.get("duration_s") + header = f"{step['title']} ({WEEKS.get(name, '')}) · {STATUS_BADGE.get(step['status'], (step['status'],))[0]}" + expanded = step["status"] not in {"ok", "skipped"} or name in {"estimate", "guard"} + with st.expander(header, expanded=expanded): + st.markdown( + f"{badge(step['status'])} {step.get('message', '')}" + + (f" \n:gray[{dur:.2f}s]" if isinstance(dur, int | float) else "") + ) + if step["status"] == "skipped": + continue + if name == "define_problem": + _plan_summary(st, plan, step) + elif name == "collect": + _collect(st, plan, step) + elif name == "structure_metrics": + _quality(st, log, step) + elif name == "estimate": + _estimate(st, log) + elif name == "guard": + _guard(st, log) + elif name == "report": + _report(st, case, log, step) diff --git a/app/pages/1_Flow.py b/app/pages/1_Flow.py new file mode 100644 index 0000000..951dd0f --- /dev/null +++ b/app/pages/1_Flow.py @@ -0,0 +1,61 @@ +"""Flow page: run the 6-step agent flow on a case, or render its existing flow_log.json.""" + +from __future__ import annotations + +import sys +from pathlib import Path + +import streamlit as st + +APP_DIR = Path(__file__).resolve().parents[1] +for p in (APP_DIR, APP_DIR.parent): # app/ for flow_view, repo root for core + if str(p) not in sys.path: + sys.path.insert(0, str(p)) + +import flow_view as fv # noqa: E402 + +st.set_page_config(page_title="Flow · 정책 효과 분석", layout="wide") +st.title("Flow — 6단계 분석 흐름") +st.caption("문제 정의 → 수집 → 지표 구조화 → 추정 → 과잉해석 가드 → 리포트") + +cases = fv.list_cases() +if not cases: + st.info("plan.yaml이 있는 케이스가 없습니다. `cases/_template`을 복사해 시작하세요.") + st.stop() + +case = st.sidebar.selectbox("케이스", cases, format_func=lambda p: p.name) +allow_uncommitted = st.sidebar.checkbox( + "allow uncommitted plan (demo)", + value=False, + help="사전 등록(plan.yaml 커밋) 게이트를 건너뜁니다. 데모 전용.", +) +llm_ok = fv.llm_available() +use_llm = st.sidebar.toggle( + "LLM 서술 사용", + value=False, + disabled=not llm_ok, + help=None if llm_ok else "LLM 키 환경변수(OPENAI_API_KEY 등)가 없어 비활성화됨", +) +run = st.sidebar.button("Run flow", type="primary") + +plan = fv.load_plan_dict(case) +log = None +if run: + with st.spinner("흐름 실행 중..."): + try: + log = fv.run_flow_for(case, allow_uncommitted, use_llm) + except ImportError: + st.error("core.agent를 불러올 수 없습니다. `make install` 후 다시 시도하세요.") + except Exception as exc: # show, don't crash the page + st.exception(exc) +else: + log = fv.load_flow_log(case) + +if log is None: + st.info(f"`{case.name}/flow_log.json`이 없습니다. 사이드바에서 **Run flow**를 누르세요.") + if plan.get("synthetic_data"): + st.warning("이 케이스는 합성 데이터(SYNTHETIC)입니다.") +else: + if not run: + st.caption(f"저장된 `{case.name}/flow_log.json`을 표시합니다 (재실행 안 함).") + fv.render_flow(st, case, log, plan) diff --git a/cases/_example_night_clinic/figures/raw_trends.png b/cases/_example_night_clinic/figures/raw_trends.png index 2e52676..c57ae93 100644 Binary files a/cases/_example_night_clinic/figures/raw_trends.png and b/cases/_example_night_clinic/figures/raw_trends.png differ diff --git a/cases/_example_night_clinic/flow_log.json b/cases/_example_night_clinic/flow_log.json new file mode 100644 index 0000000..d500b56 --- /dev/null +++ b/cases/_example_night_clinic/flow_log.json @@ -0,0 +1,281 @@ +{ + "schema_version": "1", + "case_dir": "cases/_example_night_clinic", + "question": "달빛어린이병원(야간·휴일 소아 경증 진료기관)이 지정된 시군구에서, 지정되지 않은 시군구 대비 소아 야간 경증 응급실 방문율이 감소했는가?\n", + "status": "warn", + "verdict": "identified", + "started_at": "2026-09-24T02:36:04+00:00", + "ended_at": "2026-09-24T02:36:05+00:00", + "steps": [ + { + "index": 1, + "step": "define_problem", + "title": "① 문제 정의", + "status": "ok", + "started_at": "2026-09-24T02:36:05+00:00", + "ended_at": "2026-09-24T02:36:05+00:00", + "duration_s": 0.01, + "message": "분석계획 검증 완료, 사전 등록(git 커밋) 확인: _example_night_clinic", + "artifacts": { + "plan": "plan.yaml", + "plan_sha256": "362e754b978adfa67b1606f06bc758843000e8fb4e738dd70b255dc0d87833f1", + "committed": true + } + }, + { + "index": 2, + "step": "collect", + "title": "② 데이터 수집", + "status": "warn", + "started_at": "2026-09-24T02:36:05+00:00", + "ended_at": "2026-09-24T02:36:05+00:00", + "duration_s": 0.0, + "message": "기존 스냅샷 재사용 (새로 받으려면 파일 삭제 후 재실행) / ⚠️ 합성 데이터 — 결과는 실제 정책 효과가 아님", + "artifacts": { + "data": "data/panel.csv", + "data_sha256": "20223bf371f64c27e897b820f9156c27ffce2588585c7d25386d082fa26d7eff", + "licenses": [ + "synthetic" + ], + "method": "기존 스냅샷 재사용 (새로 받으려면 파일 삭제 후 재실행)", + "source_meta": "data/panel.source.json" + } + }, + { + "index": 3, + "step": "structure_metrics", + "title": "③ 지표 구조화·품질 점검", + "status": "ok", + "started_at": "2026-09-24T02:36:05+00:00", + "ended_at": "2026-09-24T02:36:05+00:00", + "duration_s": 0.004, + "message": "11개 점검 모두 통과", + "artifacts": { + "rows": 640, + "units": 80 + } + }, + { + "index": 4, + "step": "estimate", + "title": "④ 효과 추정", + "status": "ok", + "started_at": "2026-09-24T02:36:05+00:00", + "ended_at": "2026-09-24T02:36:05+00:00", + "duration_s": 0.287, + "message": "event_study: -5.033 [-6.105, -3.960] → identified", + "artifacts": { + "night_ed_rate": { + "estimate": -5.0327, + "ci": [ + -6.1051, + -3.9603 + ], + "verdict": "identified" + } + } + }, + { + "index": 5, + "step": "guard", + "title": "⑤ 과잉해석 가드", + "status": "ok", + "started_at": "2026-09-24T02:36:05+00:00", + "ended_at": "2026-09-24T02:36:05+00:00", + "duration_s": 0.0, + "message": "판정 identified, 서술=template", + "artifacts": { + "triggers": [] + } + }, + { + "index": 6, + "step": "report", + "title": "⑥ 리포트·재현 기록", + "status": "ok", + "started_at": "2026-09-24T02:36:05+00:00", + "ended_at": "2026-09-24T02:36:05+00:00", + "duration_s": 0.439, + "message": "리포트 2개 그림 포함 작성", + "artifacts": { + "report": "report.md", + "manifest": "run_manifest.json", + "figures": [ + "figures/raw_trends.png", + "figures/event_study.png" + ] + } + } + ], + "quality": [ + { + "check": "필수 컬럼", + "value": 4, + "status": "ok", + "message": "" + }, + { + "check": "시간 컬럼 정수형", + "value": "int64", + "status": "ok", + "message": "" + }, + { + "check": "결측률 region_id", + "value": 0.0, + "status": "ok", + "message": "" + }, + { + "check": "결측률 year", + "value": 0.0, + "status": "ok", + "message": "" + }, + { + "check": "결측률 treated", + "value": 0.0, + "status": "ok", + "message": "" + }, + { + "check": "결측률 night_ed_rate", + "value": 0.0, + "status": "ok", + "message": "" + }, + { + "check": "사전 기간 수", + "value": 4, + "status": "ok", + "message": "" + }, + { + "check": "사후 기간 수", + "value": 4, + "status": "ok", + "message": "" + }, + { + "check": "처치 단위 수", + "value": 30, + "status": "ok", + "message": "" + }, + { + "check": "통제 단위 수", + "value": 50, + "status": "ok", + "message": "" + }, + { + "check": "단위×시점 중복", + "value": 0, + "status": "ok", + "message": "" + } + ], + "results": [ + { + "method": "event_study", + "outcome": "night_ed_rate", + "estimate": -5.032702345732548, + "se": 0.5471286269132056, + "ci_low": -6.105054749393283, + "ci_high": -3.9603499420718125, + "p_value": 3.6335376481325645e-20, + "n_obs": 640, + "n_clusters": 80, + "assumptions_checked": { + "parallel_pretrends": { + "test": "joint Wald (chi2)", + "stat": 0.6252747974975548, + "df": 3, + "p_value": 0.8906226863229644, + "passed": true + }, + "pre_periods": 4, + "n_adoption_cohorts": 1, + "placebo_time": { + "fake_treat_time": 2014, + "estimate": -0.32469326228919176, + "p_value": 0.5245409054480956, + "passed": true + } + }, + "warnings": [], + "triggers": [], + "verdict": "identified", + "extra": { + "coefs": [ + { + "rel_time": -4, + "coef": 0.44969282466849825, + "se": 0.6116580595373967, + "ci_low": -0.7677820885266659, + "ci_high": 1.6671677378636622 + }, + { + "rel_time": -3, + "coef": 0.17895423334833896, + "se": 0.6467334612792702, + "ci_low": -1.1083365206178435, + "ci_high": 1.4662449873145214 + }, + { + "rel_time": -2, + "coef": -0.020739466561547057, + "se": 0.559734439210234, + "ci_low": -1.1348629987606007, + "ci_high": 1.0933840656375067 + }, + { + "rel_time": -1, + "coef": 0.0, + "se": 0.0, + "ci_low": 0.0, + "ci_high": 0.0 + }, + { + "rel_time": 0, + "coef": -4.438542547221294, + "se": 0.6841551135912444, + "ci_low": -5.800319236899004, + "ci_high": -3.0767658575435837 + }, + { + "rel_time": 1, + "coef": -5.4012082194069535, + "se": 0.6881833855255733, + "ci_low": -6.771002983803212, + "ci_high": -4.031413455010695 + }, + { + "rel_time": 2, + "coef": -5.031560056431842, + "se": 0.6884702945405902, + "ci_low": -6.401925898937359, + "ci_high": -3.6611942139263256 + }, + { + "rel_time": 3, + "coef": -5.259498559870104, + "se": 0.6475211640626454, + "ci_low": -6.548357197007054, + "ci_high": -3.9706399227331532 + } + ] + }, + "identified": true + } + ], + "guard": { + "verdict": "identified", + "narrative_source": "template", + "passed": true, + "violations": [], + "rejected_llm_violations": [], + "narrative": "**요약**: 정책 도입 후 처치집단의 소아 야간 경증 응급실 방문율은(는) 통제집단 대비 평균 -5.03 건/천명 변화한 것으로 추정됩니다 (95% CI -6.11 ~ -3.96). 이 해석은 plan.yaml 의 식별 가정이 성립할 때에만 유효하며, 자동 점검에서 가정이 반박되지 않았다는 뜻이지 검증되었다는 뜻은 아닙니다." + }, + "report_path": "report.md" +} \ No newline at end of file diff --git a/cases/_example_night_clinic/report.md b/cases/_example_night_clinic/report.md index aba907e..e3e858f 100644 --- a/cases/_example_night_clinic/report.md +++ b/cases/_example_night_clinic/report.md @@ -7,6 +7,8 @@ **종합 판정**: ✅ 식별됨 (가정 하에서) +**요약**: 정책 도입 후 처치집단의 소아 야간 경증 응급실 방문율은(는) 통제집단 대비 평균 -5.03 건/천명 변화한 것으로 추정됩니다 (95% CI -6.11 ~ -3.96). 이 해석은 plan.yaml 의 식별 가정이 성립할 때에만 유효하며, 자동 점검에서 가정이 반박되지 않았다는 뜻이지 검증되었다는 뜻은 아닙니다. + ## 1. 설계 - 분석 단위: 시군구 (`region_id`), 기간 2012–2019 (year) @@ -20,7 +22,6 @@ | 방법 | 결과변수 | 추정치 | 95% CI | p | N | 클러스터 | 판정 | |---|---|---:|---|---:|---:|---:|---| | event_study | night_ed_rate | -5.033 | [-6.105, -3.960] | <0.001 | 640 | 80 | ✅ 식별됨 (가정 하에서) | -| did_twfe | night_ed_rate | -5.185 | [-5.956, -4.413] | <0.001 | 640 | 80 | ✅ 식별됨 (가정 하에서) | ![Raw trends](figures/raw_trends.png) diff --git a/cases/_example_night_clinic/run_manifest.json b/cases/_example_night_clinic/run_manifest.json new file mode 100644 index 0000000..90686df --- /dev/null +++ b/cases/_example_night_clinic/run_manifest.json @@ -0,0 +1,25 @@ +{ + "case_id": "_example_night_clinic", + "started_at": "2026-09-24T02:36:04+00:00", + "finished_at": "2026-09-24T02:36:05+00:00", + "git_sha": "0e70ae244069234d826817c46a9f730927ee6e99", + "git_dirty": true, + "plan_sha256": "362e754b978adfa67b1606f06bc758843000e8fb4e738dd70b255dc0d87833f1", + "data_sha256": "20223bf371f64c27e897b820f9156c27ffce2588585c7d25386d082fa26d7eff", + "data_path": "data/panel.csv", + "versions": { + "python": "3.11.15", + "pandas": "2.3.3", + "numpy": "2.4.6", + "scipy": "1.17.1", + "pyfixest": "0.60.0", + "statsmodels": "0.15.0", + "pydantic": "2.13.5", + "matplotlib": "3.11.2", + "langgraph": "1.2.12" + }, + "platform": "Linux-6.18.44-fc-v37-x86_64-with-glibc2.39", + "verdicts": { + "night_ed_rate": "identified" + } +} \ No newline at end of file diff --git a/core/agent/__init__.py b/core/agent/__init__.py index d55d427..99bcf6c 100644 --- a/core/agent/__init__.py +++ b/core/agent/__init__.py @@ -1,31 +1,26 @@ -"""에이전트 툴 인터페이스 (스텁). W4~5 에 Develop 이 LangGraph tool 로 래핑. +"""6단계 에이전트 흐름: 문제 정의 → 수집 → 지표 구조화 → 추정 → 과잉해석 가드 → 리포트. -규약: 각 툴은 JSON 직렬화 가능한 입력/출력만 사용한다. - - validate_plan(path) -> {"ok": bool, "errors": [...]} - - run_estimate(plan_path, data_path) -> EffectResult.to_dict() -LLM 은 plan.yaml 작성까지만 하고, 수치 계산은 반드시 core.estimators 를 호출한다. -""" - -from __future__ import annotations - -import pandas as pd -from pydantic import ValidationError - -from ..schema.plan import load_plan + python -m core.agent cases/_example_night_clinic --allow-uncommitted +- state.py : FlowState + flow_log.json 스키마 +- nodes.py : 6개 노드 (FlowState -> FlowState) +- graph.py : run_flow (순수 Python) / build_graph (LangGraph) +- guard.py : 과잉해석 린터 + 판정별 서술 템플릿 +- llm.py : OpenAI 호환 LLM 호출 (선택; 숫자 계산에는 절대 쓰지 않음) +- tools.py : validate_plan / run_estimate (JSON 툴) +""" -def validate_plan(path: str) -> dict: - try: - load_plan(path) - return {"ok": True, "errors": []} - except ValidationError as e: - return { - "ok": False, - "errors": [f"{'.'.join(map(str, x['loc']))}: {x['msg']}" for x in e.errors()], - } - - -def run_estimate(plan_path: str, data_path: str) -> list[dict]: - from ..pipeline import run_plan # 지연 import - - return [r.to_dict() for r in run_plan(load_plan(plan_path), pd.read_csv(data_path))] +from .graph import build_graph, langgraph_available, run_flow, write_flow_log +from .state import STEPS, FlowState +from .tools import run_estimate, validate_plan + +__all__ = [ + "run_flow", + "build_graph", + "langgraph_available", + "write_flow_log", + "FlowState", + "STEPS", + "validate_plan", + "run_estimate", +] diff --git a/core/agent/__main__.py b/core/agent/__main__.py new file mode 100644 index 0000000..779b046 --- /dev/null +++ b/core/agent/__main__.py @@ -0,0 +1,37 @@ +"""CLI: python -m core.agent [--allow-uncommitted] [--llm] [--question "..."]""" + +from __future__ import annotations + +import argparse +import sys + +from .graph import run_flow + +ICON = {"ok": "OK", "warn": "WARN", "failed": "FAIL", "needs_human": "HUMAN", "blocked": "BLOCK"} + + +def main(argv=None) -> int: + ap = argparse.ArgumentParser( + prog="python -m core.agent", description="6단계 정책효과 분석 흐름 실행" + ) + ap.add_argument("case_dir") + ap.add_argument("--question", help="plan.yaml 이 없을 때 초안 작성에 쓸 분석 질문") + ap.add_argument( + "--allow-uncommitted", action="store_true", help="사전 등록 게이트 우회 (데모 전용)" + ) + ap.add_argument( + "--llm", action="store_true", help="LLM 사용 (OPENAI_BASE_URL/OPENAI_API_KEY/LLM_MODEL)" + ) + ap.add_argument("--engine", choices=["auto", "python", "langgraph"], default="auto") + a = ap.parse_args(argv) + s = run_flow(a.case_dir, a.question, a.allow_uncommitted, a.llm, engine=a.engine) + for st in s.steps: + print( + f"{st.title:<18} {ICON.get(st.status, st.status):<5} {st.duration_s or 0:>6.2f}s {st.message}" + ) + print(f"\n상태={s.status} 판정={s.verdict} 로그={s.case_dir / 'flow_log.json'}") + return {"ok": 0, "warn": 0, "failed": 1}.get(s.status, 2) + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/core/agent/graph.py b/core/agent/graph.py new file mode 100644 index 0000000..74e9634 --- /dev/null +++ b/core/agent/graph.py @@ -0,0 +1,85 @@ +"""흐름 실행기: LangGraph StateGraph(있으면) 또는 순수 Python 순차 실행 — 의미는 동일. + + from core.agent import run_flow + state = run_flow("cases/_example_night_clinic", allow_uncommitted=True) + state.to_log() # flow_log.json 과 같은 dict + +두 실행기 모두 '마지막 단계가 ok/warn 일 때만 다음 단계로' 규칙을 따른다 +(failed / needs_human / blocked 이면 즉시 종료). +""" + +from __future__ import annotations + +import json +from pathlib import Path +from typing import TypedDict + +from .nodes import NODES +from .state import FlowState + + +class GraphState(TypedDict): + flow: FlowState # LangGraph 채널은 하나만 쓰고, 그 안에 FlowState 전체를 담는다 + + +def _init(case_dir, question, allow_uncommitted, use_llm) -> FlowState: + return FlowState( + case_dir=Path(case_dir).resolve(), + question=question, + allow_uncommitted=allow_uncommitted, + use_llm=use_llm, + ) + + +def write_flow_log(s: FlowState) -> Path: + path = s.case_dir / "flow_log.json" + if s.case_dir.exists(): + path.write_text(json.dumps(s.to_log(), ensure_ascii=False, indent=2, default=str), "utf-8") + return path + + +def langgraph_available() -> bool: + try: + import langgraph.graph # noqa: F401 + except ImportError: + return False + return True + + +def build_graph(): + """LangGraph StateGraph: ①→②→…→⑥, 각 단계 뒤 조건부 엣지로 중단 여부 결정.""" + from langgraph.graph import END, START, StateGraph + + g = StateGraph(GraphState) + for fn in NODES: + g.add_node(fn.__name__, lambda d, fn=fn: {"flow": fn(d["flow"])}) + g.add_edge(START, NODES[0].__name__) + names = [fn.__name__ for fn in NODES] + for a, b in zip(names, names[1:], strict=False): + g.add_conditional_edges(a, lambda d, b=b: b if d["flow"].can_continue else END, [b, END]) + g.add_edge(NODES[-1].__name__, END) + return g.compile() + + +def run_flow( + case_dir, + question: str | None = None, + allow_uncommitted: bool = False, + use_llm: bool = False, + engine: str = "python", + write_log: bool = True, +) -> FlowState: + """engine: "python"(기본, 의존성 없음) | "langgraph" | "auto"(설치돼 있으면 langgraph).""" + s = _init(case_dir, question, allow_uncommitted, use_llm) + if engine == "auto": + engine = "langgraph" if langgraph_available() else "python" + if engine == "langgraph": + s = build_graph().invoke({"flow": s})["flow"] + else: + for node in NODES: + s = node(s) + if not s.can_continue: + break + if write_log: + write_flow_log(s) + return s diff --git a/core/agent/guard.py b/core/agent/guard.py new file mode 100644 index 0000000..dec2656 --- /dev/null +++ b/core/agent/guard.py @@ -0,0 +1,98 @@ +"""⑤ 과잉해석 가드: 판정에 맞는 서술 템플릿 + 인과 단정 표현 린터. + +규칙 + 1) 판정(verdict)이 identified 가 아니면 인과·확정 표현(CAUSAL) 금지 + 2) 신뢰구간이 0을 포함하면 '효과 없음' 류 금지 → '판별 불가'로 써야 함 +린터는 단순 문자열/정규식 매칭이다. 부정문("입증되지 않았다")도 잡히는데, 의도된 보수성이다. +""" + +from __future__ import annotations + +import re + +from ..estimators.result import EffectResult + +CAUSAL = [ + "효과가 입증", + "입증", + "증명", + "때문에", + "덕분에", + "로 인해", + "인과적으로 확인", + "확실히", + "명백히", + r"\bcaus(e|es|ed|ing)\b", + r"\bproves?\b", + r"\bproven\b", + r"\bdue to\b", + r"\bbecause of\b", + r"\bdemonstrates?\b", + r"\bdefinitely\b", +] +NULL_EFFECT = [ + "효과 없음", + "효과가 없", + "효과는 없", + r"\bno effect\b", + r"\bhad no (significant )?effect\b", +] + +TRIGGER_KO = { + "pretrend_rejected": "사전추세 차이(평행추세 가정 기각)", + "few_clusters": "클러스터 수 부족(표준오차 과소추정 위험)", + "staggered_adoption": "시차 도입에서 TWFE 편향 위험", + "placebo_significant": "가짜 도입시점 검정 유의", + "ci_crosses_zero": "신뢰구간이 0을 포함", + "short_pre_period": "사전 기간 부족", +} + + +def _sentences(text: str) -> list[str]: + return [s.strip() for s in re.split(r"(?<=[.!?。])\s+|\n+", text) if s.strip()] + + +def lint(text: str, verdict: str, ci_crosses_zero: bool) -> list[dict]: + """금지 표현 목록을 반환. 빈 리스트면 통과.""" + rules = [] + if verdict != "identified": + rules += [(p, f"판정이 '{verdict}' 인데 인과·확정 표현 사용") for p in CAUSAL] + if ci_crosses_zero: + rules += [ + (p, "신뢰구간이 0을 포함 → '효과 없음'이 아니라 '판별 불가'로 서술") + for p in NULL_EFFECT + ] + out = [] + for sent in _sentences(text): + for pat, reason in rules: + m = re.search(pat, sent, flags=re.IGNORECASE) + if m: + out.append({"phrase": m.group(0), "sentence": sent, "reason": reason}) + return out + + +def template_narrative(r: EffectResult, unit: str | None = None, name: str | None = None) -> str: + """판정별 결정론적 서술 (LLM 미사용 시 기본값). 린터를 항상 통과하도록 작성.""" + u = f" {unit}" if unit else "" + ci = f"95% CI {r.ci_low:.2f} ~ {r.ci_high:.2f}" + why = ", ".join(TRIGGER_KO.get(t, t) for t in r.triggers) or "없음" + if r.verdict == "identified": + return ( + f"**요약**: 정책 도입 후 처치집단의 {name or r.outcome}은(는) 통제집단 대비 평균 " + f"{r.estimate:+.2f}{u} 변화한 것으로 추정됩니다 ({ci}). 이 해석은 plan.yaml 의 식별 가정이 " + "성립할 때에만 유효하며, 자동 점검에서 가정이 반박되지 않았다는 뜻이지 검증되었다는 뜻은 아닙니다." + ) + if r.verdict == "conditional": + zero = ( + " 신뢰구간이 0을 포함하므로 효과의 방향과 크기는 판별 불가입니다." + if ("ci_crosses_zero" in r.triggers) + else "" + ) + return ( + f"**요약(조건부)**: 추정치는 {r.estimate:+.2f}{u} ({ci}) 이지만 다음 위험 신호가 있어 " + f"조건부로만 보고합니다: {why}.{zero} 결론 전에 보완 분석이 필요합니다." + ) + return ( + f"**요약(식별 불가)**: 다음 사유로 이 설계에서는 정책 효과를 추정치로 보고하지 않습니다: {why}. " + f"표의 수치({r.estimate:+.2f}{u})는 진단용이며 정책 효과로 해석하면 안 됩니다." + ) diff --git a/core/agent/llm.py b/core/agent/llm.py new file mode 100644 index 0000000..891aed6 --- /dev/null +++ b/core/agent/llm.py @@ -0,0 +1,46 @@ +"""LLM 호출 (선택). OpenAI 호환 Chat Completions 엔드포인트 하나로 통일. + +환경변수 (.env.example 참고): + OPENAI_BASE_URL 기본 https://api.openai.com/v1 (Ollama: http://localhost:11434/v1, vLLM 등) + OPENAI_API_KEY 로컬 오픈 LLM 이면 비워도 됨 + LLM_MODEL 예: gpt-4o-mini, qwen2.5:7b-instruct +규칙: LLM 은 '글(계획 초안·서술)'만 쓴다. 숫자 계산은 항상 core.estimators. +""" + +from __future__ import annotations + +import os + +import requests + + +def configured() -> bool: + return bool(os.getenv("LLM_MODEL")) and bool( + os.getenv("OPENAI_API_KEY") or os.getenv("OPENAI_BASE_URL") + ) + + +def chat(system: str, user: str, temperature: float = 0.2, timeout: int = 120) -> str: + base = os.getenv("OPENAI_BASE_URL") or "https://api.openai.com/v1" + key = os.getenv("OPENAI_API_KEY", "") + r = requests.post( + f"{base.rstrip('/')}/chat/completions", + headers={"Authorization": f"Bearer {key}"} if key else {}, + json={ + "model": os.environ["LLM_MODEL"], + "temperature": temperature, + "messages": [{"role": "system", "content": system}, {"role": "user", "content": user}], + }, + timeout=timeout, + ) + r.raise_for_status() + return r.json()["choices"][0]["message"]["content"].strip() + + +def strip_fence(text: str) -> str: + """```yaml ... ``` 코드펜스 제거.""" + t = text.strip() + if t.startswith("```"): + t = t.split("\n", 1)[1] if "\n" in t else "" + t = t.rsplit("```", 1)[0] + return t.strip() diff --git a/core/agent/nodes.py b/core/agent/nodes.py new file mode 100644 index 0000000..005573e --- /dev/null +++ b/core/agent/nodes.py @@ -0,0 +1,373 @@ +"""6단계 노드. 각 노드는 순수 함수 `FlowState -> FlowState` (입력 상태를 바꾸지 않고 복사본 반환). + + ① define_problem → ② collect → ③ structure_metrics → ④ estimate → ⑤ guard → ⑥ report + +노드 본문은 `(status, message, artifacts)` 만 돌려주고, 로그·시간·예외 처리는 @step 이 맡는다. +멘티가 단계를 추가할 때도 같은 패턴을 쓰면 된다. +""" + +from __future__ import annotations + +import hashlib +import json +import platform +import subprocess +import sys +import time +from functools import wraps +from importlib import metadata +from pathlib import Path + +import pandas as pd +import yaml +from pydantic import ValidationError + +from ..adapters import REGISTRY, LocalFileAdapter, SourceMeta +from ..pipeline import run_plan +from ..report import event_study_plot, its_plot, raw_trends, render_report +from ..schema.plan import Plan, load_plan +from . import guard as G +from . import llm +from .state import STEPS, FlowState, StepLog, now + +ROOT = Path(__file__).resolve().parents[2] +TEMPLATE_PLAN = ROOT / "cases" / "_template" / "plan.yaml" +_TITLES = dict(STEPS) + + +class FlowError(Exception): + """사용자에게 그대로 보여줄 한국어 오류.""" + + +def step(fn): + name = fn.__name__ + + @wraps(fn) + def node(state: FlowState) -> FlowState: + s = state.model_copy(deep=True) + log = StepLog(index=len(s.steps) + 1, step=name, title=_TITLES[name]) + t0 = time.perf_counter() + try: + log.status, log.message, log.artifacts = fn(s) + except FlowError as e: + log.status, log.message = "failed", str(e) + except Exception as e: # noqa: BLE001 — 흐름을 멈추고 로그에 남긴다 + log.status, log.message = "failed", f"{type(e).__name__}: {e}" + log.ended_at, log.duration_s = now(), round(time.perf_counter() - t0, 3) + s.steps.append(log) + return s + + return node + + +def _sha256(path: Path) -> str: + return hashlib.sha256(path.read_bytes()).hexdigest() + + +def _git(cwd: Path, *args: str) -> subprocess.CompletedProcess: + return subprocess.run(["git", *args], cwd=cwd, capture_output=True, text=True) + + +def plan_committed(case_dir: Path) -> bool: + """plan.yaml 이 git 에 추적되고 HEAD 대비 수정이 없으면 True (사전 등록 확인).""" + tracked = _git(case_dir, "ls-files", "--error-unmatch", "plan.yaml").returncode == 0 + return tracked and _git(case_dir, "diff", "--quiet", "HEAD", "--", "plan.yaml").returncode == 0 + + +# ① ------------------------------------------------------------------------------------------ +@step +def define_problem(s: FlowState): + path = s.case_dir / "plan.yaml" + if not path.exists(): + s.case_dir.mkdir(parents=True, exist_ok=True) + template = TEMPLATE_PLAN.read_text(encoding="utf-8") + if s.use_llm and llm.configured() and s.question: + text = llm.strip_fence( + llm.chat( + "너는 정책 효과 분석 설계자다. 주어진 템플릿과 같은 필드만 사용해 plan.yaml 을 작성하라. " + "YAML 만 출력하라. 데이터를 보기 전 사전 등록 문서이므로 보류 조건(abstention)을 보수적으로 둔다.", + f"분석 질문: {s.question}\n\n템플릿:\n{template}", + ) + ) + path.write_text(text, encoding="utf-8") + try: + Plan.model_validate(yaml.safe_load(text)) + msg = "LLM 이 plan.yaml 초안을 작성했습니다. 사람이 검토·수정 후 커밋하고 다시 실행하세요." + except (ValidationError, yaml.YAMLError) as e: + msg = f"LLM 초안이 스키마 검증에 실패했습니다 — 직접 수정하세요:\n{e}" + return "needs_human", msg, {"plan": "plan.yaml", "drafted_by": "llm"} + header = f"# 분석 질문(입력): {s.question}\n" if s.question else "" + path.write_text(header + template, encoding="utf-8") + return ( + "needs_human", + "plan.yaml 이 없어 템플릿을 복사했습니다(LLM 미설정). 작성·커밋 후 다시 실행하세요.", + {"plan": "plan.yaml", "drafted_by": "template"}, + ) + + try: + s.plan = load_plan(path) + except ValidationError as e: + errs = "; ".join(f"{'.'.join(map(str, x['loc']))}: {x['msg']}" for x in e.errors()) + raise FlowError(f"plan.yaml 스키마 오류 — {errs}") from e + arts = { + "plan": "plan.yaml", + "plan_sha256": _sha256(path), + "committed": plan_committed(s.case_dir), + } + if arts["committed"]: + return "ok", f"분석계획 검증 완료, 사전 등록(git 커밋) 확인: {s.plan.case_id}", arts + if s.allow_uncommitted: + return ( + "warn", + "⚠️ plan.yaml 이 커밋되지 않았습니다(--allow-uncommitted, 데모 전용). " + "실제 분석에서는 결과를 보기 전에 계획을 커밋해야 합니다.", + arts, + ) + return ( + "blocked", + "사전 등록 게이트: plan.yaml 을 먼저 git 에 커밋하세요 " + "(데이터를 본 뒤 계획을 바꾸는 것을 막기 위함). 데모는 --allow-uncommitted.", + arts, + ) + + +# ② ------------------------------------------------------------------------------------------ +@step +def collect(s: FlowState): + ds = s.plan.data_sources[0] + path = s.case_dir / ds.path + if path.exists(): + how = "기존 스냅샷 재사용 (새로 받으려면 파일 삭제 후 재실행)" + elif (s.case_dir / "fetch.py").exists(): + r = subprocess.run( + [sys.executable, "fetch.py"], cwd=s.case_dir, capture_output=True, text=True + ) + if r.returncode != 0: + raise FlowError(f"fetch.py 실행 실패:\n{r.stderr[-800:]}") + how = "fetch.py 실행" + elif ds.adapter == "local": + meta = SourceMeta(ds.name, ds.provider, ds.license, ds.url) + ad = LocalFileAdapter(s.case_dir / ds.query["path"], meta) + ad.save(ad.fetch(), path, ds.query) + how = "LocalFileAdapter" + elif ds.adapter in REGISTRY: + ad = REGISTRY[ds.adapter]() + ad.save(ad.fetch(**ds.query), path, ds.query) + how = f"{ds.adapter} 어댑터" + else: + raise FlowError( + f"데이터가 없습니다: {ds.path} 가 없고 fetch.py 나 알려진 어댑터({list(REGISTRY)})도 없습니다." + ) + if not path.exists(): + raise FlowError(f"수집 후에도 {ds.path} 가 생성되지 않았습니다.") + s.data_path = path + licenses = sorted({d.license for d in s.plan.data_sources}) + arts = {"data": ds.path, "data_sha256": _sha256(path), "licenses": licenses, "method": how} + src = path.with_suffix(".source.json") + if src.exists(): + arts["source_meta"] = str(src.relative_to(s.case_dir)) + notes = [] + if s.plan.synthetic_data: + notes.append("⚠️ 합성 데이터 — 결과는 실제 정책 효과가 아님") + if {"KOGL-3", "KOGL-4"} & set(licenses): + notes.append("공공누리 3·4유형(변경금지) 포함 — 가공 데이터 재배포 불가") + return ("warn" if notes else "ok"), " / ".join([how, *notes]), arts + + +# ③ ------------------------------------------------------------------------------------------ +def quality_checks(p: Plan, df: pd.DataFrame) -> list[dict]: + """(점검, 값, 상태, 메시지) 목록. status: ok | warn | failed.""" + t, u, g = p.time.col, p.unit.id_col, p.treatment.group_col + outs = [o.col for o in p.outcomes] + cols = [u, t, g, p.treatment.first_treat_col, p.estimator.cluster_col, *outs] + need = [t, *outs] if p.estimator.method == "its" else [c for c in cols if c] + need = list(dict.fromkeys(need + p.estimator.covariates)) + q = [] + + def add(check, value, fail=False, warn=False, msg=""): + status = "failed" if fail else "warn" if warn else "ok" + q.append({"check": check, "value": value, "status": status, "message": msg}) + + missing = [c for c in need if c not in df.columns] + add("필수 컬럼", len(need) - len(missing), bool(missing), msg=f"없는 컬럼: {missing}") + if missing: + return q + ints = df[t].dropna() + ok = pd.api.types.is_integer_dtype(ints) or bool(ints.mod(1).eq(0).all()) + add("시간 컬럼 정수형", str(df[t].dtype), not ok, msg=f"'{t}' 는 정수(연도/0,1,2…)여야 함") + for c in need: + rate, key = float(df[c].isna().mean()), c in (u, t, g) + add(f"결측률 {c}", round(rate, 4), (key and rate > 0) or rate > 0.2, rate > 0, + msg=f"결측 {rate:.1%}" + (" — 식별 컬럼은 결측 불가" if key else "")) # fmt: skip + if (tt := p.treatment.treat_time) is not None: + periods = df[t].dropna().unique() + n_pre, n_post = int((periods < tt).sum()), int((periods >= tt).sum()) + min_pre = p.thresholds.min_pre_periods + add("사전 기간 수", n_pre, n_pre == 0, n_pre < min_pre, f"최소 {min_pre}기 권장") + add("사후 기간 수", n_post, n_post == 0, msg="사후 기간이 없습니다") + if p.estimator.method != "its": + by = df.groupby(u)[g].max() + n_t, n_c, dup = int((by == 1).sum()), int((by == 0).sum()), int(df.duplicated([u, t]).sum()) + add("처치 단위 수", n_t, n_t == 0, msg="처치 단위가 없습니다") + add("통제 단위 수", n_c, n_c == 0, msg="통제 단위가 없습니다") + add("단위×시점 중복", dup, dup > 0, msg="한 단위·시점에 행이 여러 개") + for x in q: # 통과한 항목은 메시지 비움 + x["message"] = x["message"] if x["status"] != "ok" else "" + return q + + +@step +def structure_metrics(s: FlowState): + df = pd.read_csv(s.data_path) + s.quality = q = quality_checks(s.plan, df) + if failed := [x for x in q if x["status"] == "failed"]: + detail = "; ".join(f"{x['check']}: {x['message']}" for x in failed) + raise FlowError(f"데이터 품질 점검 실패 — {detail}") + warns = [x["check"] for x in q if x["status"] == "warn"] + msg = f"{len(q)}개 점검" + (f", 주의: {warns}" if warns else " 모두 통과") + return ( + ("warn" if warns else "ok"), + msg, + {"rows": len(df), "units": int(df[s.plan.unit.id_col].nunique())}, + ) + + +# ④ ------------------------------------------------------------------------------------------ +@step +def estimate(s: FlowState): + s.results = run_plan(s.plan, pd.read_csv(s.data_path)) # 숫자는 여기서만 계산 (LLM 관여 없음) + arts = { + r.outcome: { + "estimate": round(r.estimate, 4), + "ci": [round(r.ci_low, 4), round(r.ci_high, 4)], + "verdict": r.verdict, + } + for r in s.results + } + r = s.results[0] + msg = f"{r.method}: {r.estimate:+.3f} [{r.ci_low:.3f}, {r.ci_high:.3f}] → {r.verdict}" + return ("warn" if any(x.warnings for x in s.results) else "ok"), msg, arts + + +# ⑤ ------------------------------------------------------------------------------------------ +@step +def guard(s: FlowState): + r = next(x for x in s.results if x.outcome == s.plan.primary_outcome.col) + zero = r.ci_low <= 0 <= r.ci_high + o = s.plan.primary_outcome + text, source, rejected = G.template_narrative(r, o.unit, o.name), "template", [] + if s.use_llm and llm.configured(): + facts = {k: v for k, v in r.to_dict().items() if k != "extra"} + draft = llm.chat( + "너는 정책 효과 분석 결과를 3문장 이내 한국어로 요약한다. 주어진 수치만 쓰고 새 숫자를 만들지 마라. " + "판정이 identified 가 아니면 인과·확정 표현(입증, 때문에, 덕분에 등)을 쓰지 마라. " + "신뢰구간이 0을 포함하면 '효과 없음' 대신 '판별 불가'라고 써라.", + json.dumps(facts, ensure_ascii=False, default=str), + ) + rejected = G.lint(draft, r.verdict, zero) + if not rejected: + text, source = draft, "llm" + violations = G.lint(text, r.verdict, zero) + s.guard = { + "verdict": r.verdict, + "narrative_source": source, + "passed": not violations, + "violations": violations, + "rejected_llm_violations": rejected, + "narrative": text, + } + if violations: + raise FlowError(f"서술 린트 실패: {violations}") + msg = f"판정 {r.verdict}, 서술={source}" + if rejected: + msg += f" (LLM 서술에서 과잉해석 {len(rejected)}건 → 템플릿으로 대체)" + return ( + ("ok" if r.verdict == "identified" and not rejected else "warn"), + msg, + {"triggers": r.triggers}, + ) + + +# ⑥ ------------------------------------------------------------------------------------------ +def _checks(plan, r) -> dict: + out, ac = {}, r.assumptions_checked + for a in plan.assumptions: + if a.check == "pretrend_test" and "parallel_pretrends" in ac: + x = ac["parallel_pretrends"] + out[a.name] = ( + f"사전계수 결합 Wald p={x['p_value']:.3f} → {'통과' if x['passed'] else '기각'}" + ) + elif a.check == "placebo_time" and "placebo_time" in ac: + x = ac["placebo_time"] + out[a.name] = ( + f"가짜 도입({x['fake_treat_time']}) 추정 {x['estimate']:.3f}, " + f"p={x['p_value']:.3f} → {'통과' if x['passed'] else '실패'}" + ) + return out + + +def _versions() -> dict: + pk = "pandas numpy scipy pyfixest statsmodels pydantic matplotlib langgraph".split() + out = {"python": platform.python_version()} + for p in pk: + try: + out[p] = metadata.version(p) + except metadata.PackageNotFoundError: + pass + return out + + +@step +def report(s: FlowState): + p, df = s.plan, pd.read_csv(s.data_path) + r = next(x for x in s.results if x.outcome == p.primary_outcome.col) + figdir = s.case_dir / "figures" + figdir.mkdir(exist_ok=True) + figs = {} + if p.estimator.method == "its": + its_plot( + df, + r.outcome, + p.time.col, + p.treatment.treat_time, + r.extra["fitted"], + r.extra["counterfactual"], + figdir / "its.png", + ) + figs["ITS"] = "figures/its.png" + else: + tt = p.treatment.treat_time or df[p.treatment.first_treat_col].min() + raw_trends(df, r.outcome, p.time.col, p.treatment.group_col, tt, figdir / "raw_trends.png") + figs["Raw trends"] = "figures/raw_trends.png" + if "coefs" in r.extra: + event_study_plot(r.extra["coefs"], figdir / "event_study.png") + figs["Event study"] = "figures/event_study.png" + md = render_report(p, s.results, figs, _checks(p, r), narrative=s.guard.get("narrative")) + s.report_path = s.case_dir / "report.md" + s.report_path.write_text(md, encoding="utf-8") + + sha = _git(s.case_dir, "rev-parse", "HEAD") + dirty = _git(s.case_dir, "status", "--porcelain", "--", ".") + manifest = { + "case_id": p.case_id, + "started_at": s.started_at, + "finished_at": now(), + "git_sha": sha.stdout.strip() if sha.returncode == 0 else None, + "git_dirty": bool(dirty.stdout.strip()) if dirty.returncode == 0 else None, + "plan_sha256": _sha256(s.case_dir / "plan.yaml"), + "data_sha256": _sha256(s.data_path), + "data_path": str(s.data_path.relative_to(s.case_dir)), + "versions": _versions(), + "platform": platform.platform(), + "verdicts": {x.outcome: x.verdict for x in s.results}, + } + (s.case_dir / "run_manifest.json").write_text( + json.dumps(manifest, ensure_ascii=False, indent=2), "utf-8" + ) + arts = {"report": "report.md", "manifest": "run_manifest.json", "figures": list(figs.values())} + return "ok", f"리포트 {len(figs)}개 그림 포함 작성", arts + + +NODES = [define_problem, collect, structure_metrics, estimate, guard, report] + +__all__ = ["NODES", "FlowError", "plan_committed"] diff --git a/core/agent/state.py b/core/agent/state.py new file mode 100644 index 0000000..311d5d3 --- /dev/null +++ b/core/agent/state.py @@ -0,0 +1,135 @@ +"""에이전트 흐름 상태(FlowState)와 flow_log.json 스키마. + +flow_log.json (schema_version "1") — Streamlit Flow 페이지·CI 가 읽는 계약: +{ + "schema_version": "1", + "case_dir": "cases/_example_night_clinic", + "question": str | null, + "status": "ok" | "warn" | "failed" | "needs_human" | "blocked" | "running", + "verdict": "identified" | "conditional" | "not_identified" | null, # 가장 보수적인 판정 + "started_at": ISO8601, "ended_at": ISO8601 | null, + "steps": [ # 실행된 단계만 순서대로 (중단 시 이후 단계 없음) + {"index": 1, "step": "define_problem", "title": "① 문제 정의", + "status": "ok"|"warn"|"failed"|"needs_human"|"blocked", + "started_at": ISO8601, "ended_at": ISO8601, "duration_s": float, + "message": str (한국어), "artifacts": {name: 경로(케이스 기준) 또는 값}} + ], + "quality": [{"check": str, "value": any, "status": "ok"|"warn"|"failed", "message": str}], + "results": [EffectResult.to_dict()], # 대용량 배열(extra.fitted 등)은 제외 + "guard": {"verdict": str, "narrative_source": "template"|"llm", "passed": bool, + "violations": [{"phrase": str, "sentence": str, "reason": str}], "narrative": str}, + "report_path": str | null +} +run_flow() 는 FlowState 를 반환하고, FlowState.to_log() 가 위 dict 를 만든다. +""" + +from __future__ import annotations + +from datetime import UTC, datetime +from pathlib import Path +from typing import Any, Literal + +from pydantic import BaseModel, ConfigDict, Field + +from ..estimators.result import EffectResult +from ..schema.plan import Plan + +Status = Literal["ok", "warn", "failed", "needs_human", "blocked", "running"] +CONTINUE = {"ok", "warn"} # 이 상태일 때만 다음 단계로 진행 +STEPS = [ + ("define_problem", "① 문제 정의"), + ("collect", "② 데이터 수집"), + ("structure_metrics", "③ 지표 구조화·품질 점검"), + ("estimate", "④ 효과 추정"), + ("guard", "⑤ 과잉해석 가드"), + ("report", "⑥ 리포트·재현 기록"), +] +_ORDER = ["ok", "warn", "needs_human", "blocked", "failed"] + + +def now() -> str: + return datetime.now(UTC).isoformat(timespec="seconds") + + +class StepLog(BaseModel): + index: int + step: str + title: str + status: Status = "running" + started_at: str = Field(default_factory=now) + ended_at: str | None = None + duration_s: float | None = None + message: str = "" + artifacts: dict[str, Any] = {} + + +class FlowState(BaseModel): + model_config = ConfigDict(arbitrary_types_allowed=True) + + case_dir: Path + question: str | None = None + allow_uncommitted: bool = False + use_llm: bool = False + started_at: str = Field(default_factory=now) + plan: Plan | None = None + data_path: Path | None = None + quality: list[dict[str, Any]] = [] + results: list[EffectResult] = [] + guard: dict[str, Any] = {} + report_path: Path | None = None + steps: list[StepLog] = [] + + @property + def last_status(self) -> str: + return self.steps[-1].status if self.steps else "running" + + @property + def can_continue(self) -> bool: + return not self.steps or self.last_status in CONTINUE + + @property + def status(self) -> str: + if not self.steps: + return "running" + worst = max((s.status for s in self.steps), key=_ORDER.index) + done = len(self.steps) == len(STEPS) or worst not in CONTINUE + return worst if done else "running" + + @property + def verdict(self) -> str | None: + v = [r.verdict for r in self.results] + order = ["identified", "conditional", "not_identified"] + return max(v, key=order.index) if v else None + + def to_log(self) -> dict: + def _res(r: EffectResult) -> dict: + d = r.to_dict() + d["extra"] = { + k: v for k, v in d["extra"].items() if k not in {"fitted", "counterfactual"} + } + return d + + def _rel(p: Path | None) -> str | None: + # 공개 레포에 실행자의 로컬 경로가 남지 않도록 케이스 폴더 기준 상대경로만 기록 + if p is None: + return None + p = Path(p) + try: + return str(p.resolve().relative_to(Path(self.case_dir).resolve())) + except ValueError: + return p.name + + return { + "schema_version": "1", + "case_dir": f"cases/{Path(self.case_dir).name}", + "question": self.question or (self.plan.question if self.plan else None), + "status": self.status, + "verdict": self.verdict, + "started_at": self.started_at, + "ended_at": self.steps[-1].ended_at if self.steps else None, + "steps": [s.model_dump() for s in self.steps], + "quality": self.quality, + "results": [_res(r) for r in self.results], + "guard": self.guard, + "report_path": _rel(self.report_path), + } diff --git a/core/agent/tools.py b/core/agent/tools.py new file mode 100644 index 0000000..f68d301 --- /dev/null +++ b/core/agent/tools.py @@ -0,0 +1,25 @@ +"""에이전트 툴 (JSON 입출력). LangGraph/MCP tool 로 그대로 래핑할 수 있는 얇은 함수들.""" + +from __future__ import annotations + +import pandas as pd +from pydantic import ValidationError + +from ..schema.plan import load_plan + + +def validate_plan(path: str) -> dict: + try: + load_plan(path) + return {"ok": True, "errors": []} + except ValidationError as e: + return { + "ok": False, + "errors": [f"{'.'.join(map(str, x['loc']))}: {x['msg']}" for x in e.errors()], + } + + +def run_estimate(plan_path: str, data_path: str) -> list[dict]: + from ..pipeline import run_plan # 지연 import + + return [r.to_dict() for r in run_plan(load_plan(plan_path), pd.read_csv(data_path))] diff --git a/core/report/markdown.py b/core/report/markdown.py index ff6139b..b27df63 100644 --- a/core/report/markdown.py +++ b/core/report/markdown.py @@ -32,7 +32,11 @@ def results_table(results: list[EffectResult]) -> str: def render_report( - plan, results: list[EffectResult], figures: dict[str, str], checks: dict | None = None + plan, + results: list[EffectResult], + figures: dict[str, str], + checks: dict | None = None, + narrative: str | None = None, ) -> str: worst = max( results, key=lambda r: ["identified", "conditional", "not_identified"].index(r.verdict) @@ -49,6 +53,7 @@ def render_report( "", f"**종합 판정**: {VERDICT_KO[worst.verdict]}", "", + *([narrative, ""] if narrative else []), "## 1. 설계", "", f"- 분석 단위: {plan.unit.name} (`{plan.unit.id_col}`), 기간 {plan.time.start}–{plan.time.end} ({plan.time.freq})", diff --git a/core/schema/plan.py b/core/schema/plan.py index 4e664c8..058d3a3 100644 --- a/core/schema/plan.py +++ b/core/schema/plan.py @@ -120,6 +120,8 @@ class DataSource(_Base): url: str | None = None license: License adapter: str | None = Field(None, description="core.adapters 의 어댑터 이름") + path: str = Field("data/panel.csv", description="케이스 폴더 기준 분석용 스냅샷 경로") + query: dict = Field({}, description="어댑터 fetch() 인자 (fetch.py 가 없을 때 사용)") notes: str = "" diff --git a/docs/ops/group-guide.md b/docs/ops/group-guide.md new file mode 100644 index 0000000..8e0b999 --- /dev/null +++ b/docs/ops/group-guide.md @@ -0,0 +1,57 @@ +# 조별 운영 가이드 + +조(약 5명)마다 공공데이터로 정책 효과 질문을 **하나** 정하고, 공용 저장소의 `cases/<조-주제>/` 폴더 하나에 결과를 쌓습니다. +전체 흐름은 6단계 에이전트 Flow(README의 Flow 절)와 같고, 조원 역할도 이 단계에 맞춰 나눕니다. + +- 완성 예시: [`cases/_example_night_clinic/`](../../cases/_example_night_clinic/) — 합성 데이터로 6단계를 끝까지 돌린 샘플 +- 주제 후보: [`docs/strategy/topic-guide.md`](../strategy/topic-guide.md) (T1~T8) +- 분석계획 작성법: [`docs/strategy/plan-guide.md`](../strategy/plan-guide.md) + +## 1. 역할 (5인 기준, 겸임 가능) + +| 역할 | 맡는 단계 | 주 산출물 | +|---|---|---| +| **문제 정의** (조장 겸임 권장) | ① 문제 정의 | `plan.yaml` — 질문·처치·대조·결과지표·가정·보류 규칙 | +| **데이터 수집** | ② 수집 | `fetch.py`, `data/README.md` (출처·라이선스·수집일) | +| **지표·품질** | ③ 지표 구조화 | 정제 패널(`data/panel.csv`), 품질 점검 통과 | +| **추정** | ④ 효과 추정 | `estimate.py`, 사전추세·placebo 결과 | +| **리포트·에이전트** | ⑤ 과잉해석 방지, ⑥ 리포트 | `report.md`, 그림, `flow_log.json` | + +4명 이하면 **수집 + 지표·품질**, **추정 + 리포트**를 한 사람이 맡습니다. 역할은 주마다 바꿔도 됩니다. + +## 2. 주차별 마일스톤 + +| 모임 | 주차 | 조가 끝낼 것 | 인증(커밋/PR) | +|---|---|---|---| +| 9.27(일) | W2 | 조 편성, 주제 후보 2개 선택, 데이터 접근 확인(API 키 발급) | `cases/<조>/` 폴더 + 이슈 "케이스 제안" | +| 10.2(금) 오프라인 | — | 주제 1개 확정, **`plan.yaml` 사전 등록 PR** | plan.yaml PR 병합 | +| 10.4(일) | W3 | `fetch.py` 동작, ③ 품질 점검 통과 | 수집·정제 PR | +| 10.11(일) | W4 | ④ 1차 추정 + 사전추세 그림 | 추정 PR | +| 10.18(일) | W5 | 반박 검정 추가, 에이전트 Flow로 ①~⑥ 한 번에 실행 | `flow_log.json` 커밋 | +| 10.25(일) | W6 | ⑤ 가드 통과(과잉 인과 표현 제거), 판정 문구 확정 | 리포트 PR | +| 11.1(일) | W7 | 다른 조가 README만 보고 재현 성공 → 공개 | 교차 재현 결과 이슈 | + +Level 1(정리된 데이터 + 문서)에서 멈춰도 산출물입니다. 판정이 "식별 불가"로 나와도 근거가 명확하면 좋은 결과입니다. + +## 3. 브랜치·PR 규칙 + +- 브랜치: `group/<짧은-설명>` (예: `group2/fetch-airkorea`). `main`에 직접 push 하지 않습니다. +- 수정 범위: **자기 조 `cases/<조>/` 폴더만.** `core/` 수정이 필요하면 이슈를 먼저 엽니다. +- **사전 등록**: `plan.yaml`을 먼저 PR로 병합한 뒤 데이터를 봅니다. Flow ①단계가 커밋되지 않은 plan.yaml을 막습니다. +- PR 하나에 한 단계. 리뷰 1명 승인 + CI 통과 후 병합(Squash). +- 올리면 안 되는 것: API 키(`.env`), 원자료 대용량 파일(`data/raw/`), 개인정보가 담긴 데이터, 재배포가 금지된 데이터(라이선스 확인). + +## 4. 실행 방법 + +```bash +make install # 최초 1회 +cp -r cases/_template cases/group1-topic # 조 폴더 만들기 +make flow CASE=cases/group1-topic # ①~⑥ 실행 (plan.yaml 커밋 후) +make app # Streamlit → Flow 페이지에서 단계별 결과 확인 +``` + +## 5. 활동 인증 + +- 매 모임 전까지 조 브랜치에 커밋 1회 이상 + PR 링크를 디스코드 조 채널에 공유합니다. +- 멘토는 `python scripts/weekly_activity.py`로 조별 커밋·변경 파일을 확인합니다 ([monitoring.md](monitoring.md)). +- 조모임은 전체 모임과 별도로 자유롭게 잡고, 결정 사항은 PR 설명이나 이슈에 남깁니다. diff --git a/tests/core/test_agent.py b/tests/core/test_agent.py new file mode 100644 index 0000000..e767794 --- /dev/null +++ b/tests/core/test_agent.py @@ -0,0 +1,146 @@ +import json +import shutil +import subprocess +from pathlib import Path + +import pytest + +from core.agent import STEPS, FlowState, build_graph, langgraph_available, run_flow +from core.agent.guard import lint, template_narrative +from core.agent.nodes import define_problem +from core.estimators.result import EffectResult + +ROOT = Path(__file__).resolve().parents[2] +EXAMPLE = ROOT / "cases" / "_example_night_clinic" + + +@pytest.fixture +def case(tmp_path): + """예제 케이스 복사본 (저장소 파일을 건드리지 않도록).""" + dst = tmp_path / "case" + shutil.copytree(EXAMPLE, dst, ignore=shutil.ignore_patterns("figures", "*.json", "report.md")) + shutil.copy(EXAMPLE / "data" / "panel.source.json", dst / "data") + return dst + + +def test_full_flow_offline_no_llm(case): + s = run_flow(case, allow_uncommitted=True, engine="python") + assert [x.step for x in s.steps] == [n for n, _ in STEPS] + assert s.status == "warn" # 미커밋 plan + 합성 데이터 경고 + assert s.verdict == "identified" and s.guard["passed"] + r = s.results[0] + assert r.ci_low <= -5.0 <= r.ci_high # 합성 데이터 참효과 + for f in ("report.md", "run_manifest.json", "flow_log.json", "figures/event_study.png"): + assert (case / f).exists(), f + log = json.loads((case / "flow_log.json").read_text(encoding="utf-8")) + assert log["schema_version"] == "1" and log["verdict"] == "identified" + assert { + "index", + "step", + "title", + "status", + "started_at", + "ended_at", + "message", + "artifacts", + } <= set(log["steps"][0]) + assert len(log["steps"][1]["artifacts"]["data_sha256"]) == 64 + assert "요약" in (case / "report.md").read_text(encoding="utf-8") + + +def _r(verdict, lo, hi, triggers=()): + return EffectResult( + "did_twfe", + "y", + (lo + hi) / 2, + 1.0, + lo, + hi, + 0.5, + 100, + 30, + triggers=list(triggers), + verdict=verdict, + ) + + +def test_guard_flags_overclaim(): + bad = "정책 때문에 방문율이 줄었고 효과가 입증되었다. This proves the policy works." + v = lint(bad, "conditional", ci_crosses_zero=False) + assert {x["phrase"].lower() for x in v} >= {"때문에", "효과가 입증", "proves"} + assert lint(bad, "identified", ci_crosses_zero=False) == [] # 식별된 경우엔 허용 + null = lint("정책은 효과 없음으로 나타났다.", "conditional", ci_crosses_zero=True) + assert any("판별 불가" in x["reason"] for x in null) + + +@pytest.mark.parametrize( + "r", + [ + _r("identified", -6, -4), + _r("conditional", -2, 1, ["ci_crosses_zero"]), + _r("not_identified", -6, -4, ["pretrend_rejected"]), + ], +) +def test_template_narrative_always_passes_lint(r): + assert lint(template_narrative(r), r.verdict, r.ci_low <= 0 <= r.ci_high) == [] + + +def _git(cwd, *a): + subprocess.run(["git", *a], cwd=cwd, check=True, capture_output=True) + + +def test_gate_blocks_uncommitted_plan(case): + _git(case, "init", "-q") + _git( + case, + "-c", + "user.name=t", + "-c", + "user.email=t@example.invalid", + "commit", + "-q", + "--allow-empty", + "-m", + "init", + ) + s = run_flow(case, write_log=False) + assert [x.status for x in s.steps] == ["blocked"] and s.status == "blocked" + + _git(case, "add", "plan.yaml") + _git( + case, + "-c", + "user.name=t", + "-c", + "user.email=t@example.invalid", + "commit", + "-q", + "-m", + "plan", + ) + assert define_problem(FlowState(case_dir=case)).steps[-1].status == "ok" + + (case / "plan.yaml").write_text((case / "plan.yaml").read_text("utf-8") + "\n# edit\n", "utf-8") + assert define_problem(FlowState(case_dir=case)).steps[-1].status == "blocked" # 커밋 후 수정 + + +def test_missing_plan_falls_back_to_template(tmp_path): + s = run_flow(tmp_path / "new_case", question="X 정책이 Y를 줄였는가?", write_log=False) + assert s.status == "needs_human" and len(s.steps) == 1 + assert (tmp_path / "new_case" / "plan.yaml").exists() + + +def test_quality_check_fails_fast(case): + (case / "data" / "panel.csv").write_text("region_id,year\nR1,2015\n", encoding="utf-8") + s = run_flow(case, allow_uncommitted=True, write_log=False) + assert s.steps[-1].step == "structure_metrics" and s.status == "failed" + assert "필수 컬럼" in s.steps[-1].message + + +@pytest.mark.skipif(not langgraph_available(), reason="langgraph 미설치") +def test_langgraph_path_matches_python(case): + assert build_graph() is not None + g = run_flow(case, allow_uncommitted=True, engine="langgraph", write_log=False) + p = run_flow(case, allow_uncommitted=True, engine="python", write_log=False) + assert [x.status for x in g.steps] == [x.status for x in p.steps] + assert g.results[0].estimate == pytest.approx(p.results[0].estimate) diff --git a/tests/core/test_plan_and_case.py b/tests/core/test_plan_and_case.py index c2411d5..7f7b42e 100644 --- a/tests/core/test_plan_and_case.py +++ b/tests/core/test_plan_and_case.py @@ -1,4 +1,6 @@ import copy +import os +import shutil import subprocess import sys from pathlib import Path @@ -44,7 +46,11 @@ def test_invalid_plans_rejected(mutate): def test_example_case_runs_end_to_end(tmp_path): + # 추적 중인 예제 산출물을 건드리지 않도록 임시 복사본에서 실행 + case = tmp_path / "cases" / CASE.name + shutil.copytree(CASE, case) + env = {**os.environ, "PYTHONPATH": str(ROOT)} for script in ("fetch.py", "estimate.py"): - subprocess.run([sys.executable, str(CASE / script)], check=True, cwd=ROOT) - report = (CASE / "report.md").read_text(encoding="utf-8") - assert "합성(synthetic)" in report and (CASE / "figures" / "event_study.png").exists() + subprocess.run([sys.executable, str(case / script)], check=True, cwd=case, env=env) + report = (case / "report.md").read_text(encoding="utf-8") + assert "합성(synthetic)" in report and (case / "figures" / "event_study.png").exists() diff --git a/tests/fixtures/flow_case/flow_log.json b/tests/fixtures/flow_case/flow_log.json new file mode 100644 index 0000000..fa8bf15 --- /dev/null +++ b/tests/fixtures/flow_case/flow_log.json @@ -0,0 +1,225 @@ +{ + "schema_version": "1", + "case_dir": "tests/fixtures/flow_case", + "question": "달빛어린이병원(야간·휴일 소아 경증 진료기관)이 지정된 시군구에서, 지정되지 않은 시군구 대비 소아 야간 경증 응급실 방문율이 감소했는가?\n", + "status": "warn", + "verdict": "identified", + "started_at": "2026-09-24T02:30:44+00:00", + "ended_at": "2026-09-24T02:30:45+00:00", + "steps": [ + { + "index": 1, + "step": "define_problem", + "title": "① 문제 정의", + "status": "ok", + "started_at": "2026-09-24T02:30:44+00:00", + "ended_at": "2026-09-24T02:30:44+00:00", + "duration_s": 0.01, + "message": "분석계획 검증 완료, 사전 등록(git 커밋) 확인: _example_night_clinic", + "artifacts": { + "plan": "plan.yaml", + "plan_sha256": "362e754b978adfa67b1606f06bc758843000e8fb4e738dd70b255dc0d87833f1", + "committed": true + } + }, + { + "index": 2, + "step": "collect", + "title": "② 데이터 수집", + "status": "warn", + "started_at": "2026-09-24T02:30:44+00:00", + "ended_at": "2026-09-24T02:30:44+00:00", + "duration_s": 0.0, + "message": "기존 스냅샷 재사용 (새로 받으려면 파일 삭제 후 재실행) / ⚠️ 합성 데이터 — 결과는 실제 정책 효과가 아님", + "artifacts": { + "data": "data/panel.csv", + "data_sha256": "20223bf371f64c27e897b820f9156c27ffce2588585c7d25386d082fa26d7eff", + "licenses": [ + "synthetic" + ], + "method": "기존 스냅샷 재사용 (새로 받으려면 파일 삭제 후 재실행)", + "source_meta": "data/panel.source.json" + } + }, + { + "index": 3, + "step": "structure_metrics", + "title": "③ 지표 구조화·품질 점검", + "status": "ok", + "started_at": "2026-09-24T02:30:44+00:00", + "ended_at": "2026-09-24T02:30:44+00:00", + "duration_s": 0.004, + "message": "11개 점검 모두 통과", + "artifacts": { + "rows": 640, + "units": 80 + } + }, + { + "index": 4, + "step": "estimate", + "title": "④ 효과 추정", + "status": "ok", + "started_at": "2026-09-24T02:30:44+00:00", + "ended_at": "2026-09-24T02:30:44+00:00", + "duration_s": 0.282, + "message": "event_study: -5.033 [-6.105, -3.960] → identified", + "artifacts": { + "night_ed_rate": { + "estimate": -5.0327, + "ci": [ + -6.1051, + -3.9603 + ], + "verdict": "identified" + } + } + }, + { + "index": 5, + "step": "guard", + "title": "⑤ 과잉해석 가드", + "status": "warn", + "started_at": "2026-09-24T02:30:44+00:00", + "ended_at": "2026-09-24T02:30:44+00:00", + "duration_s": 0.0, + "message": "판정 identified, 서술=template", + "artifacts": { + "triggers": [] + } + }, + { + "index": 6, + "step": "report", + "title": "⑥ 리포트·재현 기록", + "status": "ok", + "started_at": "2026-09-24T02:30:44+00:00", + "ended_at": "2026-09-24T02:30:45+00:00", + "duration_s": 0.414, + "message": "리포트 2개 그림 포함 작성", + "artifacts": { + "report": "report.md", + "manifest": "run_manifest.json", + "figures": [] + } + } + ], + "quality": [ + { + "check": "필수 컬럼", + "value": 4, + "status": "ok", + "message": "" + }, + { + "check": "시간 컬럼 정수형", + "value": "int64", + "status": "ok", + "message": "" + }, + { + "check": "결측률 region_id", + "value": 0.0, + "status": "ok", + "message": "" + }, + { + "check": "결측률 year", + "value": 0.0, + "status": "ok", + "message": "" + }, + { + "check": "결측률 treated", + "value": 0.0, + "status": "ok", + "message": "" + }, + { + "check": "결측률 night_ed_rate", + "value": 0.0, + "status": "ok", + "message": "" + }, + { + "check": "사전 기간 수", + "value": 4, + "status": "ok", + "message": "" + }, + { + "check": "사후 기간 수", + "value": 4, + "status": "ok", + "message": "" + }, + { + "check": "처치 단위 수", + "value": 30, + "status": "ok", + "message": "" + }, + { + "check": "통제 단위 수", + "value": 50, + "status": "ok", + "message": "" + }, + { + "check": "단위×시점 중복", + "value": 0, + "status": "ok", + "message": "" + } + ], + "results": [ + { + "method": "event_study", + "outcome": "night_ed_rate", + "estimate": -5.032702345732548, + "se": 0.547128626913206, + "ci_low": -6.105054749393284, + "ci_high": -3.960349942071812, + "p_value": 3.6335376481327216e-20, + "n_obs": 640, + "n_clusters": 80, + "assumptions_checked": { + "parallel_pretrends": { + "test": "joint Wald (chi2)", + "stat": 0.6252747974975572, + "df": 3, + "p_value": 0.8906226863229638, + "passed": true + }, + "pre_periods": 4, + "n_adoption_cohorts": 1, + "placebo_time": { + "fake_treat_time": 2014, + "estimate": -0.32469326228919176, + "p_value": 0.5245409054480956, + "passed": true + } + }, + "warnings": [], + "triggers": [], + "verdict": "identified", + "extra": {}, + "identified": true + } + ], + "guard": { + "verdict": "identified", + "narrative_source": "template", + "passed": false, + "violations": [ + { + "phrase": "입증", + "sentence": "정책이 효과를 입증했다.", + "reason": "인과 단정 표현" + } + ], + "rejected_llm_violations": [], + "narrative": "**요약**: 정책 도입 후 처치집단의 소아 야간 경증 응급실 방문율은(는) 통제집단 대비 평균 -5.03 건/천명 변화한 것으로 추정됩니다 (95% CI -6.11 ~ -3.96). 이 해석은 plan.yaml 의 식별 가정이 성립할 때에만 유효하며, 자동 점검에서 가정이 반박되지 않았다는 뜻이지 검증되었다는 뜻은 아닙니다." + }, + "report_path": "report.md" +} \ No newline at end of file diff --git a/tests/fixtures/flow_case/plan.yaml b/tests/fixtures/flow_case/plan.yaml new file mode 100644 index 0000000..4706429 --- /dev/null +++ b/tests/fixtures/flow_case/plan.yaml @@ -0,0 +1,10 @@ +# Minimal fixture for app tests (not validated against core schema). +case_id: flow_case +title: Flow fixture +question: 야간 소아 진료 확대가 경증 응급실 방문을 줄였는가? +synthetic_data: true +treatment: {definition: 정책 도입 지역} +control: {definition: 미도입 지역} +estimator: {method: event_study} +data_sources: + - {name: synthetic panel, provider: generated, license: synthetic} diff --git a/tests/fixtures/flow_case/report.md b/tests/fixtures/flow_case/report.md new file mode 100644 index 0000000..11b30e8 --- /dev/null +++ b/tests/fixtures/flow_case/report.md @@ -0,0 +1,3 @@ +# Flow fixture report + +FIXTURE-REPORT-BODY diff --git a/tests/test_app_flow.py b/tests/test_app_flow.py new file mode 100644 index 0000000..8bd984e --- /dev/null +++ b/tests/test_app_flow.py @@ -0,0 +1,73 @@ +"""Flow page renders an existing flow_log.json without running the agent (AppTest, no network).""" + +from __future__ import annotations + +import json +import shutil +import sys +from pathlib import Path + +import pytest + +ROOT = Path(__file__).resolve().parents[1] +FIXTURE = ROOT / "tests" / "fixtures" / "flow_case" +PAGE = ROOT / "app" / "pages" / "1_Flow.py" + +sys.path.insert(0, str(ROOT / "app")) +import flow_view as fv # noqa: E402 + +AppTest = pytest.importorskip("streamlit.testing.v1").AppTest + + +@pytest.fixture +def cases_root(tmp_path, monkeypatch): + shutil.copytree(FIXTURE, tmp_path / "flow_case") + monkeypatch.setenv("PEA_CASES_DIR", str(tmp_path)) + for k in fv.LLM_KEY_ENVS: + monkeypatch.delenv(k, raising=False) + return tmp_path + + +def _text(at) -> str: + parts = [ + e.value + for kind in ("markdown", "error", "info", "caption", "success") + for e in getattr(at, kind) + ] + return "\n".join(str(p) for p in parts) + + +def test_normalize_fills_six_steps(): + log = fv.normalize_flow_log({"steps": [{"step": "collect", "status": "ok"}]}) + assert [s["step"] for s in log["steps"]] == [n for n, _ in fv.STEPS] + assert [s["status"] for s in log["steps"]] == ["skipped", "ok"] + ["skipped"] * 4 + + +def test_steps_match_core_contract(): + state = pytest.importorskip("core.agent.state") + assert [n for n, _ in state.STEPS] == [n for n, _ in fv.STEPS] + + +def test_fixture_matches_schema_version(): + raw = json.loads((FIXTURE / "flow_log.json").read_text(encoding="utf-8")) + assert raw["schema_version"] == "1" and len(raw["steps"]) == 6 + + +def test_flow_page_renders_existing_log(cases_root): + at = AppTest.from_file(str(PAGE), default_timeout=60).run() + assert not at.exception, at.exception + text = _text(at) + assert "SYNTHETIC" in text + assert "식별됨" in text # verdict label, not raw "identified" + assert "과잉해석 표현 1건" in text # guard finding highlighted + assert "FIXTURE-REPORT-BODY" in text # report preview + assert len(at.expander) >= 6 + assert at.sidebar.toggle[0].disabled # no LLM key -> toggle disabled + assert at.sidebar.checkbox[0].label == "allow uncommitted plan (demo)" + + +def test_flow_page_without_log_prompts_run(cases_root): + (cases_root / "flow_case" / "flow_log.json").unlink() + at = AppTest.from_file(str(PAGE), default_timeout=60).run() + assert not at.exception + assert any("Run flow" in i.value for i in at.info)