"""특징 프로그램(feature program) DSL → 검증 → 시점 제약을 부가한 SQL → DuckDB 실행. 특허 구성 (c)(d): 의미 모델의 조인 경로로 탐색 범위를 한정하고, 모든 프로그램에 '관측 기준 시점(seed_time) 이전 레코드만' 이라는 시점 제약을 컴파일러가 강제로 부가한다. LLM 은 시점 제약을 쓰지 않는다 — 쓸 수 없다. """ from __future__ import annotations import re from dataclasses import dataclass import duckdb import pandas as pd from semantic_model import SemanticModel AGGS = {"count": "count", "sum": "sum", "mean": "avg", "min": "min", "max": "max", "std": "stddev_samp", "nunique": "count(distinct", "last": "arg_max"} WINDOWS = (30, 90, 365, 1095, 3650, None) # days; None = 전체 이력 # expr 안에 허용되는 비-컬럼 토큰 (DuckDB 스칼라 표현식 부분집합). 이 밖의 식별자는 전부 거부. SQL_WORDS = { "case", "when", "then", "else", "end", "is", "null", "not", "and", "or", "in", "between", "like", "true", "false", "as", "cast", "integer", "int", "double", "float", "bigint", "date", "timestamp", "coalesce", "nullif", "abs", "round", "floor", "ceil", "ln", "log", "sqrt", "greatest", "least", "date_diff", "datediff", "date_part", "extract", "year", "month", "day", "interval", "days", "seed_time", "epoch", "sign", "power", } IDENT = re.compile(r"[A-Za-z_][A-Za-z0-9_.]*") STRING = re.compile(r"'[^']*'") @dataclass class FeatureProgram: name: str path: list[str] # entity table 에서 시작하는 조인 경로 agg: str # AGGS 키, 또는 "none" (엔티티 자체 속성) expr: str # 리프/경로 테이블 컬럼에 대한 스칼라식 window_days: int | None # 시점 제약 창; None = 전체 이력 rationale: str = "" @classmethod def from_dict(cls, d: dict) -> "FeatureProgram": return cls(name=str(d["name"]), path=list(d["path"]), agg=str(d.get("agg", "none")).lower(), expr=str(d["expr"]).strip(), window_days=d.get("window_days"), rationale=str(d.get("rationale", ""))) def to_dict(self) -> dict: return {"name": self.name, "path": self.path, "agg": self.agg, "expr": self.expr, "window_days": self.window_days, "rationale": self.rationale} def validate(fp: FeatureProgram, sm: SemanticModel, entity_table: str) -> str | None: """None = 통과. 문자열 = 거부 사유. (선언 안 된 테이블·컬럼·경로는 참조 자체를 거부)""" if not re.fullmatch(r"[A-Za-z_][A-Za-z0-9_]{0,63}", fp.name): return f"bad feature name {fp.name!r}" if fp.path[:1] != [entity_table]: return f"path must start at entity table {entity_table}" if (e := sm.path_error(fp.path)): return e if fp.agg == "none": if len(fp.path) != 1: return "agg=none is only for entity-table attributes (path length 1)" elif fp.agg not in AGGS: return f"unknown agg {fp.agg!r}; use one of {sorted(AGGS)}" elif sm.time_table(fp.path) is None: return "aggregation path has no timestamped table; as-of constraint impossible" if fp.window_days is not None and fp.window_days not in WINDOWS: return f"window_days must be one of {WINDOWS}" if ";" in fp.expr or "--" in fp.expr or "/*" in fp.expr: return "statement separators/comments are not allowed in expr" allowed = sm.columns_on(fp.path) dropped = sm.sensitive("drop") # 문자열 리터럴 제거 후 식별자 검사 for tok in IDENT.findall(STRING.sub("''", fp.expr)): if tok.lower() in SQL_WORDS: continue if tok not in allowed: return f"undeclared identifier {tok!r} in expr (allowed: columns of {fp.path})" if tok in dropped or any(f"{t}.{tok}" in dropped for t in fp.path): return f"column {tok!r} is classified PII; not usable" return None def compile_sql(fp: FeatureProgram, sm: SemanticModel, label_view: str = "L") -> str: """L(row_id, entity_id, seed_time) 와 조인하여 row_id 별 특징값 1개를 내는 SQL.""" ent = fp.path[0] pk = sm.tables[ent]["pkey"] joins = [f"JOIN {ent} ON {ent}.{pk} = {label_view}.entity_id"] for a, b in zip(fp.path, fp.path[1:]): joins.append(f"JOIN {b} ON {sm.join_condition(a, b)}") expr = fp.expr.replace("seed_time", f"{label_view}.seed_time") if fp.agg == "none": return f"SELECT {label_view}.row_id, ({expr}) AS value FROM {label_view} {' '.join(joins)}" # 시점 제약: 경로 위의 시간열 보유 테이블 *전부* 에 부가한다 (하나만 걸면 하류 테이블로 미래가 샌다) timed = [(t, f"{t}.{sm.tables[t]['time_col']}") for t in fp.path if sm.tables[t]["time_col"]] tcol = timed[0][1] # last/arg_max 의 기준 시간열 where = [] for _, c in timed: where.append(f"{c} < {label_view}.seed_time") if fp.window_days is not None: where.append(f"{c} >= {label_view}.seed_time - INTERVAL {int(fp.window_days)} DAY") if fp.agg == "last": agg_sql = f"arg_max(({expr}), {tcol})" elif fp.agg == "nunique": agg_sql = f"count(distinct ({expr}))" else: agg_sql = f"{AGGS[fp.agg]}(({expr}))" return (f"SELECT {label_view}.row_id, {agg_sql} AS value FROM {label_view} {' '.join(joins)} " f"WHERE {' AND '.join(where)} GROUP BY {label_view}.row_id") class Executor: """제1 컴퓨팅 환경: DB 옆에서 특징 프로그램을 실행해 특징 행렬을 만든다.""" def __init__(self, db_tables: dict[str, pd.DataFrame]): self.con = duckdb.connect() for name, df in db_tables.items(): self.con.register(name, df) def run(self, programs: list[FeatureProgram], sm: SemanticModel, labels: pd.DataFrame, entity_col: str, time_col: str) -> tuple[pd.DataFrame, dict[str, str]]: """labels: 예측 요청 표 (entity_col, time_col[, target]). 반환: 특징 행렬(row 순서 유지), 실패 사유.""" L = pd.DataFrame({"row_id": range(len(labels)), "entity_id": labels[entity_col].values, "seed_time": pd.to_datetime(labels[time_col]).values}) self.con.register("L", L) out = pd.DataFrame(index=L.row_id) errors = {} for fp in programs: try: r = self.con.execute(compile_sql(fp, sm)).df().set_index("row_id")["value"] out[fp.name] = pd.to_numeric(r.reindex(out.index), errors="coerce").astype("float64") except Exception as e: # 실행 실패는 거부 사유로 기록, 행렬에서 제외 errors[fp.name] = f"{type(e).__name__}: {str(e).splitlines()[0][:160]}" return out.reset_index(drop=True), errors