diff --git a/scripts/data_platform/corpus_download.py b/scripts/data_platform/corpus_download.py index f6d4e8bb..45a11d78 100644 --- a/scripts/data_platform/corpus_download.py +++ b/scripts/data_platform/corpus_download.py @@ -89,7 +89,10 @@ UA = {"User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.3 ID_KEYS = {"ann_meta": ["announcement_id"], "news_meta": ["art_code", "stock_code"], "news_fulltext": ["art_code"], - "flash_meta": ["flash_id"]} + "flash_meta": ["flash_id"], + # P4-3 事件快照两域(spec §4.8 F2): 输入切片按内容寻址, 判读按事件 id + "events_in": ["content_hash"], + "events_llm": ["event_id"]} FLASH_EM_MAX_PAGES = 20 # 快照型时效数据: 50条/页×20≈2-3天窗, 不深回补 FLASH_LOOKBACK_DAYS = 1 # 东财快讯翻页拉到昨日 00:00 边界即停 @@ -736,6 +739,7 @@ class Ctx: self.deadline = dt.datetime.combine( dt.date.today(), dt.time(int(hh), int(mm))) self.failed = 0 + self.parse_fail = 0 # 解析/契约违规计数(P4-3 漏斗) self.rate_limited = False self.hard_cool = False self.wallclock = False diff --git a/scripts/data_platform/corpus_funnel.py b/scripts/data_platform/corpus_funnel.py new file mode 100644 index 00000000..9e99759d --- /dev/null +++ b/scripts/data_platform/corpus_funnel.py @@ -0,0 +1,375 @@ +# -*- coding: utf-8 -*- +"""corpus_funnel.py — 漏斗分诊+事件快照夜批(spec §4.8 决议 G/C7+F1-F5, P4-3). + +两层: 第一层确定性词表(funnel_lexicon, 零成本处理大头) → 词表 miss 且含 +信号词 → 第二层 LLM(只碰漏下来的, sanguo_api/llm client 串行慢爬). 产物= +事件快照两域(append-only dt 分区=采集日, 永不重算): + events_in LLM 候选输入切片(防编数: 喂模型的确切文本先落盘, 可回放对拍) + events_llm 事件判读(词表+LLM 统一域, model 列区分来源) + +基建全复用 corpus_download(Ctx 预算断路/_absorb_new 账本去重/marker); +自有 flock(state/funnel.lock, corpus 班错峰后跑互不锁). 增量=state/ +funnel_watermark.json 各源域独立水位只前向; 历史回填=--start 按预算分批; +LLM 候选持久队列 state/funnel_pending.jsonl(预算停摆残余次夜续). + +退出码(与 corpus_download 同契约): 0=完成 1=致命(unit 失败/LLM 配置缺失 +且有候选) 2=限流让路 3=预算/墙钟让路(checkpoint, 次夜续). +用法: python corpus_funnel.py --until 23:00 --llm-cap 200 + python corpus_funnel.py --no-llm # 只跑第一层 + python corpus_funnel.py --start 2026-09-06 # 历史回填(按预算分批) +""" +from __future__ import annotations + +import argparse +import asyncio +import datetime as dt +import fcntl +import hashlib +import json +import logging +import os +import sys +from pathlib import Path + +_HERE = os.path.dirname(os.path.abspath(__file__)) +_REPO = os.path.dirname(os.path.dirname(_HERE)) +for _p in (_HERE, _REPO): + if _p not in sys.path: + sys.path.insert(0, _p) + +import pandas as pd # noqa: E402 + +from scripts.data_platform import corpus_download as cd # noqa: E402 +from scripts.data_platform import funnel_lexicon as lex # noqa: E402 +from sanguo_data import partition_reader as pr # noqa: E402 + +log = logging.getLogger("funnel") + +PENDING_NAME = "funnel_pending.jsonl" +WATERMARK_NAME = "funnel_watermark.json" +SRC_DOMAINS = ("news_meta", "ann_meta") +LLM_CONSEC_FAIL_BREAK = 3 # 连续失败熔断(防烧夜, 已处理件有 marker 不重) + +EVENTS_IN_COLUMNS = ["content_hash", "src_domain", "src_id", "stock_code", + "event_date", "title", "text_snippet", "selected_at"] +EVENTS_LLM_COLUMNS = ["event_id", "content_hash", "src_domain", "src_id", + "stock_code", "event_date", "event_type", "direction", + "confidence", "model", "extracted_at"] + +PROMPT_SYSTEM = ( + "你是A股公告/新闻事件抽取器。输入是单一股票单一日期的一段标题与摘要。" + "只依据输入文本判断, 禁止编造输入中没有的信息。" + '输出严格 JSON: {"events": [{"event_type": "...", ' + '"direction": "pos|neg|neutral", "confidence": 0.0到1.0}]}。' + "event_type 只能取以下之一: " + + ", ".join(f"{k}({v})" for k, v in lex.EVENT_TYPES.items()) + + "。无把握或无事件时输出空 events 数组。stock_code 与 event_date 由系统" + "从输入行带入, 你不要输出它们。") + + +# ---------- 快照行构造 ---------- + +def _content_hash(domain, src_id, code, title): + s = f"{domain}|{src_id}|{code}|{title}" + return hashlib.blake2b(s.encode("utf-8"), digest_size=8).hexdigest() + + +def _event_row(content_hash, src_domain, src_id, code, edate, hit, model): + eid = hashlib.blake2b( + f"{content_hash}|{hit['event_type']}|{code}|{hit['direction']}" + .encode("utf-8"), digest_size=8).hexdigest() + now = dt.datetime.now().strftime("%Y-%m-%d %H:%M:%S") + return {"event_id": eid, "content_hash": content_hash, + "src_domain": src_domain, "src_id": str(src_id), + "stock_code": str(code), "event_date": edate, + "event_type": hit["event_type"], "direction": hit["direction"], + "confidence": float(hit["confidence"]), "model": model, + "extracted_at": now} + + +def _extract_day(domain, day, df): + """一天 corpus 行 → (词表事件行, LLM 候选输入行)。""" + events, cands = [], [] + for row in df.to_dict("records"): + if domain == "news_meta": + src_id, code = row.get("art_code"), row.get("stock_code") + when, title = row.get("show_time"), row.get("title") + text = row.get("summary") + else: + src_id, code = row.get("announcement_id"), row.get("sec_code") + when, title = row.get("ann_time"), row.get("title") + text = row.get("content") or row.get("short_title") + if not src_id or not code or not title: + continue + edate = str(when or "")[:10] or day + ch = _content_hash(domain, src_id, code, title) + for hit in lex.match_events(title, text): + events.append(_event_row(ch, domain, src_id, code, edate, hit, + model="lexicon")) + if lex.is_llm_candidate(title, text): + cands.append({"content_hash": ch, "src_domain": domain, + "src_id": str(src_id), "stock_code": str(code), + "event_date": edate, "title": title, + "text_snippet": f"{title}\n{text or ''}"[:500], + "selected_at": dt.date.today().isoformat()}) + return events, cands + + +# ---------- 水位与 pending 队列 ---------- + +def _load_watermark(): + p = cd.CORPUS_ROOT / "state" / WATERMARK_NAME + if p.exists(): + try: + return json.loads(p.read_text(encoding="utf-8")) + except ValueError: + log.warning("%s 损坏, 重置", WATERMARK_NAME) + return {} + + +def _save_watermark(wm): + p = cd.CORPUS_ROOT / "state" / WATERMARK_NAME + p.parent.mkdir(parents=True, exist_ok=True) + tmp = p.with_suffix(".tmp") + tmp.write_text(json.dumps(wm), encoding="utf-8") + os.replace(tmp, p) + + +def _pending_path(): + return cd.CORPUS_ROOT / "state" / PENDING_NAME + + +def _load_pending(): + p = _pending_path() + if not p.exists(): + return [] + out = [] + for ln in p.read_text(encoding="utf-8").splitlines(): + ln = ln.strip() + if not ln: + continue + try: + row = json.loads(ln) + except ValueError: + continue + if not cd.is_done("funnel", "llm", row.get("content_hash")): + out.append(row) # 消费 marker 幂等: done 件不重排 + return out + + +def _save_pending(rows): + p = _pending_path() + p.parent.mkdir(parents=True, exist_ok=True) + tmp = p.with_suffix(".tmp") + tmp.write_text( + "".join(json.dumps(r, ensure_ascii=False) + "\n" for r in rows), + encoding="utf-8") + os.replace(tmp, p) + + +# ---------- 单实例锁(flock, corpus 班错峰互不干扰) ---------- + +def _acquire_lock(): + lock = cd.CORPUS_ROOT / "state" / "funnel.lock" + lock.parent.mkdir(parents=True, exist_ok=True) + fd = os.open(str(lock), os.O_CREAT | os.O_RDWR, 0o666) + try: + fcntl.flock(fd, fcntl.LOCK_EX | fcntl.LOCK_NB) + except BlockingIOError: + os.close(fd) + return None + return fd + + +# ---------- 卡片需求面(F5: 记录不驱动) ---------- + +def _read_card_demand(pipeline_db): + if not pipeline_db or not os.path.exists(pipeline_db): + return [] + try: + from sanguo_portfolio import pipeline_store + items = pipeline_store.list_hypotheses(pipeline_db) + except Exception as e: # noqa: BLE001 读侧可选件 + log.warning("卡片需求面读取失败(非致命): %s", e) + return [] + return [{"hyp_id": i.get("id"), "dataNeeds": i.get("dataNeeds")} + for i in items if i.get("state") == "data_check"] + + +# ---------- 第一层: 词表走水位 ---------- + +def _walk_domain(ctx, domain, args, ledgers, recents, fresh_cands, stats): + root = cd.CORPUS_ROOT + wm = _load_watermark() + todo = [d for d in pr.list_dates(str(root), domain) + if (wm.get(domain) is None or d > wm[domain]) + and (args.start is None or d >= args.start)] + for d in sorted(todo): + ctx.stop_now() # WallClock/BudgetStop + df = pr.read_partitions(str(root), domain, start=d, end=d) + if not df.empty: + events, cands = _extract_day(domain, d, df) + if events: + cd._absorb_new("events_llm", events, ledgers, recents) + stats["lex_events"] += len(events) + fresh_cands.extend(cands[: args.llm_cap]) + stats["candidates_fresh"] += len(cands[: args.llm_cap]) + ctx.unit_done(domain, d) + stats["days"] += 1 + wm[domain] = d + _save_watermark(wm) + + +# ---------- 第二层: LLM lane(串行慢爬, 预算断路) ---------- + +async def _llm_lane(ctx, client, model_name, cands, ledgers, recents): + """逐候选串行抽取。返回写出事件数。 + + - 已 done(content_hash marker)件跳过(崩溃恢复幂等); + - 预算触顶抛 BudgetStop 由调用方收拾(残余回写 pending); + - 非法判读项确定性丢弃计 ctx.parse_fail, 不整批废(F3); + - 消费即落 events_in(防编数: 喂模型文本先落盘)。 + """ + written = 0 + consec_fail = 0 + for row in cands: + if cd.is_done("funnel", "llm", row["content_hash"]): + continue + ctx.stop_now() + try: + resp, usage = await client.chat_json( + [{"role": "system", "content": PROMPT_SYSTEM}, + {"role": "user", "content": _prompt_user(row)}], + want_usage=True) + except Exception as e: # noqa: BLE001 LLMError 族 + log.warning("llm fail %s: %s", row["content_hash"], e) + ctx.failed += 1 + consec_fail += 1 + if consec_fail >= LLM_CONSEC_FAIL_BREAK: + log.error("llm 连续 %d 败, 熔断余量", consec_fail) + break + continue + consec_fail = 0 + ctx.add_tokens(usage.get("prompt_tokens"), usage.get("completion_tokens")) + events = [] + for ev in (resp.get("events") or []) if isinstance(resp, dict) else []: + etype = ev.get("event_type") + direction = ev.get("direction") + try: + conf = float(ev.get("confidence")) + except (TypeError, ValueError): + conf = -1.0 + if (etype not in lex.EVENT_TYPES + or direction not in ("pos", "neg", "neutral") + or not 0.0 <= conf <= 1.0): + ctx.parse_fail += 1 + continue + events.append(_event_row( + row["content_hash"], row["src_domain"], row["src_id"], + row["stock_code"], row["event_date"], + {"event_type": etype, "direction": direction, + "confidence": conf}, model_name)) + if events: + cd._absorb_new("events_llm", events, ledgers, recents) + written += len(events) + cd._absorb_new("events_in", [dict(row)], ledgers, recents) + ctx.unit_done("llm", row["content_hash"]) + return written + + +def _prompt_user(row): + return (f"stock_code={row['stock_code']}\nevent_date={row['event_date']}\n" + f"标题: {row['title']}\n摘要: {row['text_snippet']}") + + +def _dump_raw_fail(content_hash, resp_text): + d = cd.CORPUS_ROOT / "state" / "funnel_raw" + d.mkdir(parents=True, exist_ok=True) + (d / f"{content_hash}.txt").write_text(str(resp_text), encoding="utf-8") + + +# ---------- 入口 ---------- + +def main(argv=None): + ap = argparse.ArgumentParser(description=__doc__) + ap.add_argument("--until", help="墙钟自停 HH:MM(rc=3 checkpoint)") + ap.add_argument("--limit", type=int, help="LLM 调用次数上限(冒烟不落 marker)") + ap.add_argument("--token-budget", type=int, help="夜批 token 上限(rc=3)") + ap.add_argument("--llm-cap", type=int, default=200, + help="每夜新鲜候选上限(默认 200)") + ap.add_argument("--start", help="回填起点 YYYY-MM-DD(默认=水位续跑)") + ap.add_argument("--no-llm", action="store_true", help="只跑第一层") + ap.add_argument("--pipeline-db", default=os.environ.get("SANGUO_PIPELINE_DB"), + help="卡片库路径(data_check 态需求面, 缺省 env/缺文件=跳过)") + args = ap.parse_args(argv) + + logging.basicConfig(level=logging.INFO, + format="%(asctime)s %(levelname)s %(name)s %(message)s") + fd = _acquire_lock() + if fd is None: + print("[funnel] 锁让路退出(另一实例在跑)") + return 3 + + ctx = cd.Ctx("funnel", until=args.until, limit=args.limit, + token_budget=args.token_budget) + ledgers = {d: cd.IdLedger(d) for d in ("events_in", "events_llm")} + recents = {d: cd.RecentIndex(d) for d in ("events_in", "events_llm")} + ctx.set_stores(ledgers) + cd._reconcile_ledgers(ledgers, recents) + + stats = {"days": 0, "lex_events": 0, "candidates_fresh": 0, + "llm_calls": 0, "llm_events": 0, "parse_fail": 0, + "tokens": 0, "pending_before": 0, "pending_after": 0, + "card_demand": []} + config_error = False + + pending = _load_pending() + stats["pending_before"] = len(pending) + fresh: list[dict] = [] + stats["card_demand"] = _read_card_demand(args.pipeline_db) + + stopped = None + try: + for domain in SRC_DOMAINS: + _walk_domain(ctx, domain, args, ledgers, recents, fresh, stats) + if not args.no_llm and (pending or fresh): + try: + from sanguo_api.llm import resolve_config, LLMClient + cfg = resolve_config() + except Exception as e: # noqa: BLE001 配置缺失 + log.error("LLM 配置缺失(env 三件见 runbook P4-1 节): %s", e) + config_error = True + else: + client = LLMClient(cfg) + ctx.reset_stage() # limit=LLM 调用数, 段内预算 + try: + n = asyncio.run(_llm_lane(ctx, client, cfg.model, + pending + fresh, ledgers, recents)) + stats["llm_events"] = n + except cd.BudgetStop: + stopped = "budget" + ctx.commit() + except cd.WallClockStop: + stopped = "wallclock" + ctx.commit() + except cd.BudgetStop: + stopped = "budget" + ctx.commit() + + # 残余回写: 消费 marker 过滤后重排(done 件自然消失, 幂等收敛) + remainder = [r for r in (pending + fresh) + if not cd.is_done("funnel", "llm", r["content_hash"])] + _save_pending(remainder) + stats["pending_after"] = len(remainder) + stats["tokens"] = ctx.tokens_used + stats["llm_calls"] = ctx.stage_units + + rc = 1 if config_error else ctx.rc() + print("[funnel] 完成 统计: " + json.dumps(stats, ensure_ascii=False) + + (f" stopped={stopped}" if stopped else "")) + os.close(fd) + return rc + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/tests/data_platform/test_corpus_funnel.py b/tests/data_platform/test_corpus_funnel.py new file mode 100644 index 00000000..973e4d81 --- /dev/null +++ b/tests/data_platform/test_corpus_funnel.py @@ -0,0 +1,239 @@ +# -*- coding: utf-8 -*- +"""TDD for corpus_funnel.py — 两层漏斗+事件快照(P4-3, spec §4.8 F1-F5). + +全 mock 零网络(LLM 层注入 FakeClient); 日期全相对 today; CORPUS_ROOT 重定向 +tmp_path(patch cd 模块属性, cd._day_dir/_absorb_new 均动态查全局)。 +""" +import asyncio +import datetime as dt +import json + +import pandas as pd +import pytest + +from scripts.data_platform import corpus_download as cd +from scripts.data_platform import corpus_funnel as cf +from scripts.data_platform import monitor_checks +from sanguo_data import partition_reader as pr + + +@pytest.fixture +def root(tmp_path, monkeypatch): + monkeypatch.setattr(cd, "CORPUS_ROOT", tmp_path) + return tmp_path + + +def _day(offset=0): + return (dt.date.today() - dt.timedelta(days=offset)).isoformat() + + +def _seed_day(root, domain, day, rows): + d = root / domain / f"dt={day}" + d.mkdir(parents=True, exist_ok=True) + pd.DataFrame(rows).to_parquet(d / "part-0.parquet", index=False) + + +def _news_rows(day, specs): + """specs=[(art, code, title, summary)] → news_meta 行.""" + return [{"art_code": a, "stock_code": c, "show_time": f"{day} 09:00:00", + "title": t, "summary": s, "first_seen_date": day} + for a, c, t, s in specs] + + +def _read_domain(root, domain): + return pr.read_partitions(str(root), domain) + + +def _pending(root): + p = root / "state" / cf.PENDING_NAME + if not p.exists(): + return [] + return [json.loads(ln) for ln in p.read_text(encoding="utf-8").splitlines() + if ln.strip()] + + +# ---------- 第一层: 词表+水位+pending ---------- + +def test_first_layer_walk_and_watermark(root): + day = _day(1) + _seed_day(root, "news_meta", day, _news_rows(day, [ + ("A1", "600519", "控股股东拟增持公司股份公告", None), + ("A2", "000001", "控股股东不减持公告", None), + ("A3", "300750", "日常经营动态", "运营平稳"), + ])) + rc = cf.main(["--no-llm"]) + assert rc == 0 + ev = _read_domain(root, "events_llm").to_dict("records") + assert len(ev) == 1 and ev[0]["event_type"] == "holder_increase" + assert ev[0]["model"] == "lexicon" and ev[0]["confidence"] == 1.0 + wm = cf._load_watermark() + assert wm["news_meta"] == day + assert len(_pending(root)) == 1 and _pending(root)[0]["src_id"] == "A2" + + +def test_ann_meta_domain_walked(root): + day = _day(1) + _seed_day(root, "ann_meta", day, [{ + "announcement_id": "B1", "sec_code": "000002", "ann_time": day, + "title": "关于回购公司股份的进展公告", "content": None, + "short_title": None}]) + assert cf.main(["--no-llm"]) == 0 + ev = _read_domain(root, "events_llm").to_dict("records") + assert any(r["event_type"] == "buyback" and r["src_domain"] == "ann_meta" + for r in ev) + + +def test_idempotent_rerun_absorbed(root): + day = _day(1) + _seed_day(root, "news_meta", day, _news_rows(day, [ + ("A1", "600519", "控股股东拟增持公司股份公告", None)])) + cf.main(["--no-llm"]) + n1 = len(_read_domain(root, "events_llm")) + cf.main(["--no-llm"]) # 水位已过, 全零重扫 + assert len(_read_domain(root, "events_llm")) == n1 + # 水位重置重跑(模拟崩溃恢复) → 账本吸收, 行数不变 + cf._save_watermark({}) + cf.main(["--no-llm"]) + assert len(_read_domain(root, "events_llm")) == n1 + + +def test_old_dates_behind_watermark_skipped(root): + d0, d1 = _day(2), _day(1) + _seed_day(root, "news_meta", d0, _news_rows(d0, [ + ("A0", "600519", "控股股东拟增持公告", None)])) + _seed_day(root, "news_meta", d1, _news_rows(d1, [ + ("A1", "000001", "股东减持计划公告", None)])) + cf._save_watermark({"news_meta": d0}) + cf.main(["--no-llm"]) + ev = _read_domain(root, "events_llm").to_dict("records") + assert len(ev) == 1 and ev[0]["src_id"] == "A1" + + +def test_no_signal_rows_not_candidated(root): + day = _day(1) + _seed_day(root, "news_meta", day, _news_rows(day, [ + ("A3", "300750", "日常经营动态", "运营平稳")])) + cf.main(["--no-llm"]) + assert _read_domain(root, "events_llm").empty + assert _pending(root) == [] + + +def test_stats_line_monitor_parseable(root, capsys): + day = _day(1) + _seed_day(root, "news_meta", day, _news_rows(day, [ + ("A1", "600519", "控股股东拟增持公告", None)])) + cf.main(["--no-llm"]) + out = capsys.readouterr().out + stats = monitor_checks.parse_stat_lines(out).get("funnel") + assert stats is not None and stats["lex_events"] == 1 + + +def test_start_flag_bounds_backfill(root): + d0, d1 = _day(5), _day(1) + _seed_day(root, "news_meta", d0, _news_rows(d0, [ + ("A0", "600519", "控股股东拟增持公告", None)])) + _seed_day(root, "news_meta", d1, _news_rows(d1, [ + ("A1", "000001", "股东减持计划公告", None)])) + cf.main(["--no-llm", "--start", d1]) + ev = _read_domain(root, "events_llm").to_dict("records") + assert len(ev) == 1 and ev[0]["src_id"] == "A1" + + +# ---------- 第二层: LLM lane ---------- + +class FakeClient: + def __init__(self, resp, usage=None): + self.calls = [] + self._resp = resp + self._usage = usage or {"prompt_tokens": 10, "completion_tokens": 5} + + async def chat_json(self, messages, *, temperature=0.2, max_tokens=2000, + want_usage=False): + self.calls.append(messages) + return (self._resp, dict(self._usage)) if want_usage else self._resp + + +def _mk_ctx(**kw): + return cd.Ctx("funnel", **kw) + + +def _ledgers(): + led = {d: cd.IdLedger(d) for d in ("events_in", "events_llm")} + rec = {d: cd.RecentIndex(d) for d in ("events_in", "events_llm")} + return led, rec + + +def test_llm_lane_valid_event_written(root): + day = _day(1) + _seed_day(root, "news_meta", day, _news_rows(day, [ + ("A2", "000001", "控股股东不减持公告", None)])) + cf.main(["--no-llm"]) # 入队 1 候选 + cands = _pending(root) + ctx = _mk_ctx() + led, rec = _ledgers() + fake = FakeClient({"events": [{"event_type": "holder_decrease", + "direction": "neg", "confidence": 0.8}]}) + n = asyncio.run(cf._llm_lane(ctx, fake, "glm-test", cands, led, rec)) + ctx.commit() + assert n == 1 and len(fake.calls) == 1 + ins = _read_domain(root, "events_in").to_dict("records") + evs = _read_domain(root, "events_llm").to_dict("records") + assert any(r["model"] == "glm-test" and r["event_type"] == "holder_decrease" + for r in evs) + assert len(ins) == 1 and ins[0]["content_hash"] == evs[0]["content_hash"] + assert ctx.tokens_used == 15 + assert cd.is_done("funnel", "llm", cands[0]["content_hash"]) + + +def test_llm_lane_invalid_items_discarded(root): + ctx = _mk_ctx() + led, rec = _ledgers() + cands = [{"content_hash": "h1", "src_domain": "news_meta", "src_id": "X", + "stock_code": "000001", "event_date": _day(1), "title": "t", + "text_snippet": "s", "selected_at": _day(1)}] + fake = FakeClient({"events": [ + {"event_type": "not_a_type", "direction": "pos", "confidence": 0.5}, + {"event_type": "buyback", "direction": "maybe", "confidence": 0.5}, + {"event_type": "buyback", "direction": "pos", "confidence": 1.5}, + {"event_type": "buyback", "direction": "pos", "confidence": 0.7}]}) + n = asyncio.run(cf._llm_lane(ctx, fake, "m", cands, led, rec)) + assert n == 1 and ctx.parse_fail == 3 + + +def test_llm_lane_budget_stop_leaves_rest_pending(root): + ctx = _mk_ctx(token_budget=1) + led, rec = _ledgers() + cands = [{"content_hash": f"h{i}", "src_domain": "news_meta", + "src_id": f"X{i}", "stock_code": "000001", + "event_date": _day(1), "title": "t", "text_snippet": "s", + "selected_at": _day(1)} for i in range(3)] + fake = FakeClient({"events": []}) + with pytest.raises(cd.BudgetStop): + asyncio.run(cf._llm_lane(ctx, fake, "m", cands, led, rec)) + assert ctx.tokens_used > 0 and ctx.budget_stopped + + +def test_llm_prompt_carries_row_text(root): + ctx = _mk_ctx() + led, rec = _ledgers() + cand = [{"content_hash": "h1", "src_domain": "news_meta", "src_id": "X", + "stock_code": "600519", "event_date": "2026-10-02", + "title": "标题甲", "text_snippet": "摘要乙", + "selected_at": _day(0)}] + fake = FakeClient({"events": []}) + asyncio.run(cf._llm_lane(ctx, fake, "m", cand, led, rec)) + sysmsg, usrmsg = fake.calls[0][0]["content"], fake.calls[0][1]["content"] + assert "holder_increase" in sysmsg and "事件抽取器" in sysmsg + assert "600519" in usrmsg and "标题甲" in usrmsg and "2026-10-02" in usrmsg + + +def test_llm_lane_skips_done_markers(root): + ctx = _mk_ctx() + led, rec = _ledgers() + cd.mark_done("funnel", "llm", "h1") + cands = [{"content_hash": "h1", "src_domain": "news_meta", "src_id": "X", + "stock_code": "000001", "event_date": _day(1), "title": "t", + "text_snippet": "s", "selected_at": _day(1)}] + fake = FakeClient({"events": []}) + n = asyncio.run(cf._llm_lane(ctx, fake, "m", cands, led, rec)) + assert n == 0 and fake.calls == []