"""NCA-ADS 原创练习。仅用合成数据；CPU 默认，GPU 需兼容的 RAPIDS 环境。
运行：python ads_labs.py --backend cpu  或 --backend gpu
输出写入当前目录下 ads_lab_output，不修改已有数据或安装环境。
"""
import argparse
import importlib.metadata
import json
import platform
import sys
import time
from pathlib import Path

import joblib
import numpy as np
import pandas as pd
from sklearn.metrics import confusion_matrix, f1_score, mean_absolute_error
from sklearn.model_selection import GroupShuffleSplit


def environment(backend="cpu"):
    info = {"backend": backend, "python": sys.version, "executable": sys.executable,
            "architecture": platform.machine(), "packages": {}}
    for name in ("numpy", "pandas", "scikit-learn", "joblib"):
        info["packages"][name] = importlib.metadata.version(name)
    if backend == "gpu":
        import cupy as cp
        import cudf
        import cuml
        if cp.cuda.runtime.getDeviceCount() < 1:
            raise RuntimeError("GPU 不可见；请核对当前环境及容器设备配置。")
        info["packages"].update(cupy=cp.__version__, cudf=cudf.__version__, cuml=cuml.__version__)
        info["gpu_count"] = cp.cuda.runtime.getDeviceCount()
    print(json.dumps(info, ensure_ascii=False, indent=2))
    return info


def to_host(values):
    if hasattr(values, "to_numpy"):
        return values.to_numpy()
    if hasattr(values, "get"):
        return values.get()
    return np.asarray(values)


def etl(backend="cpu"):
    if backend == "gpu":
        import cudf as frame
    else:
        frame = pd
    records = frame.DataFrame({"subject_id": [1, 1, 2, 3], "power": [2., 4., None, 8.]})
    people = frame.DataFrame({"subject_id": [1, 2, 3], "site": ["A", "A", "B"]})
    assert len(people) == int(people.subject_id.nunique()), "右表键必须唯一"
    joined = records.merge(people, on="subject_id", how="left")
    assert len(joined) == 4
    summary = joined.groupby("site").power.mean().sort_index()
    print(summary)
    np.testing.assert_allclose(to_host(summary), [3., 8.])
    # 故意制造右表重复，体会连接膨胀；不覆盖原表
    duplicated = frame.concat([people, people.iloc[:1]], ignore_index=True)
    expanded = records.merge(duplicated, on="subject_id", how="left")
    assert len(expanded) == 6
    print("原连接 4 行；信息表中 subject_id=1 重复后为", len(expanded), "行")
    return summary


def classification(backend="cpu"):
    rng = np.random.default_rng(42)
    groups = np.repeat(np.arange(200), 10)
    X = rng.normal(size=(len(groups), 4)).astype("float32")
    subject_effect = rng.normal(scale=.3, size=200)[groups]
    y = (X[:, 0] - .5 * X[:, 1] + subject_effect + rng.normal(scale=.8, size=len(groups)) > .4).astype("int32")
    X[rng.random(X.shape) < .04] = np.nan
    split = GroupShuffleSplit(n_splits=1, test_size=.25, random_state=42)
    train_idx, valid_idx = next(split.split(X, y, groups))
    assert set(groups[train_idx]).isdisjoint(set(groups[valid_idx]))
    # 只用训练数据拟合填补参数
    medians = np.nanmedian(X[train_idx], axis=0)
    Xtr = np.where(np.isnan(X[train_idx]), medians, X[train_idx])
    Xva = np.where(np.isnan(X[valid_idx]), medians, X[valid_idx])
    if backend == "gpu":
        import cupy as cp
        from cuml.preprocessing import StandardScaler
        from cuml.linear_model import LogisticRegression
        Xtr, Xva, ytr = cp.asarray(Xtr), cp.asarray(Xva), cp.asarray(y[train_idx])
        model = LogisticRegression(max_iter=1000, output_type="numpy")
    else:
        from sklearn.preprocessing import StandardScaler
        from sklearn.linear_model import LogisticRegression
        ytr = y[train_idx]
        model = LogisticRegression(max_iter=1000, random_state=42)
    scaler = StandardScaler()
    Xtr = scaler.fit_transform(Xtr)
    Xva = scaler.transform(Xva)
    model.fit(Xtr, ytr)
    pred = to_host(model.predict(Xva)).reshape(-1)
    metrics = {"valid_f1": float(f1_score(y[valid_idx], pred)),
               "confusion_matrix_rows_true_cols_pred": confusion_matrix(y[valid_idx], pred).tolist(),
               "train_subjects": len(set(groups[train_idx])),
               "valid_subjects": len(set(groups[valid_idx])), "subject_overlap": 0}
    print(json.dumps(metrics, indent=2))
    print("这是合成数据实际运行的结果，不是临床性能或官方考试分数。")
    return {"model": model, "scaler": scaler, "medians": medians,
            "feature_names": ["feature_0", "feature_1", "feature_2", "feature_3"],
            "backend": backend}, X[valid_idx][:10].copy(), metrics


def benchmark(backend="cpu", rows=200_000):
    rng = np.random.default_rng(42)
    host = pd.DataFrame({"group": rng.integers(0, 100, rows), "value": rng.normal(size=rows)})
    def sync():
        if backend == "gpu":
            import cupy as cp
            cp.cuda.runtime.deviceSynchronize()
    if backend == "gpu":
        import cudf
        sync()
        t0 = time.perf_counter()
        data = cudf.from_pandas(host)
        sync()
        conversion = time.perf_counter()-t0
    else:
        data = host
        conversion = 0.
    # 预热，不计入稳态
    data.groupby("group").value.mean()
    sync()
    times = []
    for _ in range(3):
        sync()
        t0 = time.perf_counter()
        out = data.groupby("group").value.mean()
        sync()
        times.append(time.perf_counter()-t0)
    expected = host.groupby("group").value.mean().sort_index().to_numpy()
    np.testing.assert_allclose(to_host(out.sort_index()), expected, rtol=1e-6, atol=1e-7)
    result = {"backend": backend, "rows": rows, "conversion_seconds": conversion,
              "aggregation_seconds": times, "median_aggregation_seconds": float(np.median(times)),
              "scope": "预热后分组均值；排除数据生成、输入转换与最终结果传回"}
    print(json.dumps(result, ensure_ascii=False, indent=2))
    print("比较硬件需分别运行两种 backend；这里不预设 GPU 会更快。")
    return result


def time_series():
    # 独立 CPU 小例子，用于看清时间含义；不是 GPU 性能测试
    rng = np.random.default_rng(42)
    df = pd.DataFrame({"time": pd.date_range("2026-01-01", periods=100, freq="h"),
                       "value": np.sin(np.arange(100)/8) + rng.normal(scale=.1, size=100)})
    df = df.sort_values("time")
    df["lag1"] = df.value.shift(1)
    df["past_mean3"] = df.value.shift(1).rolling(3).mean()
    np.testing.assert_allclose(df.loc[3, "past_mean3"], df.loc[:2, "value"].mean())
    ready = df.dropna().reset_index(drop=True)
    cut = int(len(ready)*.7)
    train, valid = ready.iloc[:cut], ready.iloc[cut:]
    assert train.time.max() < valid.time.min()
    # 一步向前评估：每次允许使用刚观测到的上一真实值；不是多步盲预测
    mae = mean_absolute_error(valid.value, valid.lag1)
    print(df.head(6).to_string(index=False))
    print("一步持久性预测 MAE:", float(mae))
    return float(mae)


def save_reload(bundle, example, metrics, info, output="ads_lab_output"):
    out = Path(output)
    out.mkdir(parents=True, exist_ok=True)
    def predict(b):
        filled = np.where(np.isnan(example), b["medians"], example)
        if b["backend"] == "gpu":
            import cupy as cp
            filled = cp.asarray(filled)
        return to_host(b["model"].predict(b["scaler"].transform(filled))).reshape(-1)
    before = predict(bundle)
    joblib.dump(bundle, out/"trusted_local_model.joblib")
    restored = joblib.load(out/"trusted_local_model.joblib")  # 只加载自己生成的可信文件
    np.testing.assert_array_equal(before, predict(restored))
    report = {"seed": 42, "split": "GroupShuffleSplit by synthetic subject_id",
              "feature_names": bundle["feature_names"], "environment": info, "metrics": metrics,
              "reload_predictions_equal": True, "data": "synthetic only"}
    (out/"run_report.json").write_text(json.dumps(report, ensure_ascii=False, indent=2), encoding="utf-8")
    print("已保存并验证：", out.resolve())
    return report


def run_all(backend="cpu"):
    info = environment(backend)
    etl(backend)
    bundle, example, metrics = classification(backend)
    benchmark(backend)
    time_series()
    return save_reload(bundle, example, metrics, info)


if __name__ == "__main__":
    parser = argparse.ArgumentParser(description=__doc__)
    parser.add_argument("--backend", choices=["cpu", "gpu"], default="cpu")
    args = parser.parse_args()
    run_all(args.backend)
