from __future__ import annotations from online_mining_common import ( case_matches, compact_text, emit_error, emit_success, extract_case, load_json_payload, search_docs, source_config, string_values, write_optional_jsonl, ) def build_index_filters(source: str, filters: dict) -> list[dict]: """尽量把可索引条件下推到 ES,其余条件在客户端二次过滤。""" config = source_config(source) result: list[dict] = [] if source == "main": for key in ("domain", "func"): values = string_values(filters.get(key)) if values: result.append({"terms": {key: values}}) contains = string_values(filters.get("query_contains")) if contains: # query 是 keyword 字段,wildcard 可用但不要放太宽,调用方要限制日期和 scan_size。 for item in contains: result.append({"wildcard": {"query": {"value": f"*{item}*"}}}) if filters.get("request_id"): result.append({"terms": {str(config["id_field"]): string_values(filters.get("request_id"))}}) return result def main() -> int: try: payload = load_json_payload() source = str(payload.get("source") or "pre_processing") date = payload.get("date") date_from = payload.get("date_from") date_to = payload.get("date_to") lookback_days = payload.get("lookback_days") lookback_days = int(lookback_days) if lookback_days not in (None, "") else None size = int(payload.get("size") or 50) scan_size = int(payload.get("scan_size") or max(size * 5, size)) filters = payload.get("filters") or {} if not isinstance(filters, dict): raise ValueError("filters must be an object") docs = search_docs( source=source, date=date, size=scan_size, query_filters=build_index_filters(source, filters), date_from=date_from, date_to=date_to, lookback_days=lookback_days, ) cases = [] for doc in docs: case = extract_case(source, doc) if not case_matches(case, filters): continue cases.append(case) if len(cases) >= size: break output_path = write_optional_jsonl(str(payload.get("output_path") or ""), cases) emit_success( { "source": source, "date": date or (f"{date_from}..{date_to}" if date_from and date_to else f"past-{lookback_days}d" if lookback_days else "past-48h"), "scanned": len(docs), "matched": len(cases), "output_path": output_path, "cases": [ { key: compact_text(value, 1200) for key, value in case.items() if key != "raw" and value not in (None, "", [], {}) } for case in cases ], } ) return 0 except Exception as exc: # noqa: BLE001 emit_error(exc) return 1 if __name__ == "__main__": raise SystemExit(main())