File size: 6,325 Bytes
26b3403
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
"""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:  # usable without the kit installed
    class Engine:  # type: ignore[no-redef]
        def __init__(self, **options: Any) -> None:
            self.options = options

        def warmup(self):
            pass

    class Unsupported(ValueError):  # type: ignore[no-redef]
        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}