| """Decision Index engine over the Core ML window models: |
| `--engine dmodel_mac.engine:CoreMLDecisionEngine --option model_dir=... --option name=...`. |
| |
| Per question: render windows (dmodel_mac.render), run each window through the Core ML bucket |
| that fits it, pool marker logits per option (log-mean-exp over occurrences), divide by the |
| calibrated temperature and softmax. Nothing is truncated and no option is dropped. |
| """ |
|
|
| from __future__ import annotations |
|
|
| import json |
| import math |
| import time |
| from pathlib import Path |
| from typing import Any |
|
|
| import numpy as np |
|
|
| from dmodel_mac.render import MAX_OPTIONS, Renderer, WindowConfig, option_list |
| from dmodel_mac.render import Unsupported as RenderUnsupported |
|
|
| try: |
| from decision_index.engines import Engine, Unsupported |
| except ImportError: |
| class Engine: |
| def __init__(self, **options: Any) -> None: |
| self.options = options |
|
|
| def warmup(self): |
| pass |
|
|
| class Unsupported(ValueError): |
| pass |
|
|
| UNITS = {"all": "ALL", "ane": "CPU_AND_NE", "gpu": "CPU_AND_GPU", "cpu": "CPU_ONLY"} |
|
|
|
|
| def to_answer(question: dict[str, Any], probs: np.ndarray) -> dict[str, Any]: |
| keys, descriptions = option_list(question) |
| values = [float(p) for p in probs] |
| if len(values) != len(keys) or any(not math.isfinite(v) or v < 0 for v in values): |
| raise RuntimeError("invalid distribution") |
| total = sum(values) |
| values = [v / total for v in values] |
| if question["type"] == "noul": |
| return {"type": "noul", "noul": values[1]} |
| best = max(range(len(values)), key=values.__getitem__) |
| dist = dict(zip(keys, values)) |
| if question["type"] == "choice": |
| return {"type": "choice", "choice": keys[best], "probabilities": dist} |
| return {"type": "score", "probabilities": dist, "legend": dict(zip(keys, descriptions)), |
| "score": sum(i * p for i, p in enumerate(values)), "level": keys[best]} |
|
|
|
|
| class WindowRunner: |
| """Loads one Core ML model per bucket and scores windows.""" |
|
|
| def __init__(self, model_dir: str, name: str, slots: int, buckets: tuple[int, ...], units: dict[int, str]): |
| import coremltools as ct |
| self.slots, self.models = slots, {} |
| for b in buckets: |
| unit = getattr(ct.ComputeUnit, UNITS[units[b]]) |
| self.models[b] = ct.models.MLModel(str(Path(model_dir) / f"{name}_L{b}_K{slots}.mlpackage"), compute_units=unit) |
| self.buckets = buckets |
|
|
| def window_logits(self, window) -> np.ndarray: |
| """Logits of the window's markers, in marker order (chunks of `slots` markers per call).""" |
| b = next(x for x in self.buckets if window.length <= x) |
| ids = np.zeros((1, b), dtype=np.int32) |
| ids[0, :] = 50283 |
| ids[0, :window.length] = window.ids |
| mask = np.zeros((1, b), dtype=np.int32) |
| mask[0, :window.length] = 1 |
| out = [] |
| for s in range(0, len(window.positions), self.slots): |
| mm = np.zeros((1, self.slots, b), dtype=np.float16) |
| chunk = window.positions[s:s + self.slots] |
| mm[0, np.arange(len(chunk)), chunk] = 1 |
| res = self.models[b].predict({"input_ids": ids, "attention_mask": mask, "marker_map": mm}) |
| out.append(np.asarray(res["logits"], dtype=np.float32).reshape(-1)[:len(chunk)]) |
| return np.concatenate(out) |
|
|
|
|
| def pool_lme(windows, logits_per_window, n_opt: int) -> np.ndarray: |
| occ: list[list[float]] = [[] for _ in range(n_opt)] |
| for w, z in zip(windows, logits_per_window): |
| for o, v in zip(w.options, z): |
| occ[o].append(float(v)) |
| pooled = np.empty(n_opt, dtype=np.float64) |
| for o, vals in enumerate(occ): |
| a = np.asarray(vals, dtype=np.float64) |
| m = a.max() |
| pooled[o] = m + math.log(np.exp(a - m).mean()) |
| return pooled |
|
|
|
|
| class CoreMLDecisionEngine(Engine): |
| name = "dmodel-mac" |
| latency = ("in-process wall time of the whole request: rendering, tokenization, every Core ML window " |
| "prediction, pooling and softmax; excludes model loading") |
|
|
| def __init__(self, **options: Any) -> None: |
| super().__init__(**options) |
| cfg = json.loads(Path(options["config"]).read_text()) if "config" in options else {} |
| cfg.update({k: v for k, v in options.items() if k != "config"}) |
| self.temperature = float(cfg.get("temperature", 1.0)) |
| self.slots = int(cfg.get("slots", 64)) |
| buckets = tuple(int(x) for x in str(cfg.get("buckets", "128,256,512")).split(",")) |
| default_units = {b: ("ane" if b <= 128 else "all") for b in buckets} |
| for b in buckets: |
| if f"units_{b}" in cfg: |
| default_units[b] = cfg[f"units_{b}"] |
| self.renderer = Renderer(cfg.get("tokenizer", "bases/modernbert-base/tokenizer.json"), |
| WindowConfig(max_len=buckets[-1], buckets=buckets)) |
| self.runner = WindowRunner(cfg["model_dir"], cfg["name"], self.slots, buckets, default_units) |
| self.provenance = {"model_dir": cfg["model_dir"], "name": cfg["name"], "temperature": self.temperature, |
| "buckets": list(buckets), "units": {str(k): v for k, v in default_units.items()}, |
| "slots": self.slots, "window_config": self.renderer.cfg.__dict__} |
|
|
| def logits(self, state: Any, question: dict[str, Any]) -> np.ndarray: |
| keys, _ = option_list(question) |
| if len(keys) > MAX_OPTIONS: |
| raise Unsupported(f"{len(keys)} options exceeds the declared limit of {MAX_OPTIONS}") |
| try: |
| windows, n = self.renderer.windows(state, question) |
| except RenderUnsupported as e: |
| raise Unsupported(str(e)) from e |
| return pool_lme(windows, [self.runner.window_logits(w) for w in windows], n) |
|
|
| def __call__(self, state: Any, questions: dict[str, Any]): |
| t0 = time.perf_counter() |
| answers, windows = {}, 0 |
| for key, q in questions.items(): |
| z = self.logits(state, q) / self.temperature |
| p = np.exp(z - z.max()) |
| p /= p.sum() |
| answers[key] = to_answer(q, p) |
| return {"model": self.name, "answers": answers}, {"wall_ms": (time.perf_counter() - t0) * 1000, "windows": windows} |
|
|