File size: 3,708 Bytes
d10ad42
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
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
"""One decision request -> probability distributions (same scoring as the kev benchmark engine).

Every question of a request is scored against the same state in token-budgeted, length-sorted
batches, one forward pass per batch. Nothing is truncated: a request that does not fit (too long,
too many options, or out of GPU memory) raises CapacityError, which the server maps to HTTP 413.
"""

import json
from collections.abc import Mapping, Sequence
from typing import Any, cast

import torch

from maincode_jev_serve.model import MAX_OPTIONS, DecisionModel, options
from maincode_jev_serve.types import DecisionInput, ImageInput


class CapacityError(ValueError):
    pass


def length_estimate(row: Mapping[str, Any]) -> int:
    """Cheap prompt-length estimate used only to order and pack batches (same as kev.train)."""
    measured = row.get("source", {}).get("input_tokens")
    if isinstance(measured, int) and not isinstance(measured, bool):
        return measured + 16
    return len(json.dumps([row["state"], row["question"]], ensure_ascii=False)) // 3 + 192 + 512 * len(row.get("images", []))


def microbatches(rows: Sequence[dict[str, Any]], batch_size: int, token_budget: int) -> list[list[dict[str, Any]]]:
    result: list[list[dict[str, Any]]] = []
    pending: list[dict[str, Any]] = []
    longest = 0
    for row in rows:
        length = length_estimate(row)
        if pending and (len(pending) == batch_size or max(longest, length) * (len(pending) + 1) > token_budget):
            result.append(pending)
            pending, longest = [], 0
        pending.append(row)
        longest = max(longest, length)
    if pending:
        result.append(pending)
    return result


def decide(model: DecisionModel, state: object, questions: Mapping[str, Mapping[str, Any]], *, temperature: float,
           max_tokens: int, token_budget: int, batch_size: int, images: Sequence[ImageInput] = ()) -> tuple[dict[str, list[float]], int]:
    """Normalized probabilities per question key (in option order) and the input tokens used."""
    rows: list[dict[str, Any]] = []
    for key, question in questions.items():
        count = len(options(cast(Any, question))[0])
        if count > MAX_OPTIONS:
            raise CapacityError(f"at most {MAX_OPTIONS} options per choice question are supported ({count} given)")
        row: dict[str, Any] = {"state": state, "question": dict(question), "id": key, "source": {}}
        if images:
            row["images"] = list(images)
        rows.append(row)
    distributions: dict[str, list[float]] = {}
    input_tokens = 0
    for batch in microbatches(sorted(rows, key=length_estimate), batch_size, token_budget):
        try:
            prepared = model.prepare(cast(list[DecisionInput], batch), max_length=max_tokens)
        except ValueError as error:
            if "token limit" in str(error):
                raise CapacityError(f"request exceeds the maximum context length of {max_tokens} tokens") from error
            raise
        input_tokens += prepared.input_tokens
        try:
            with torch.inference_mode():
                probabilities = (model(prepared) / temperature).softmax(-1).float().cpu().tolist()
        except torch.OutOfMemoryError as error:
            del prepared
            torch.cuda.empty_cache()
            raise CapacityError("request exceeds the maximum context length that fits in GPU memory") from error
        for item, values, count in zip(batch, probabilities, prepared.counts, strict=True):
            total = sum(values[:count])
            distributions[item["id"]] = [value / total for value in values[:count]]
    return {key: distributions[key] for key in questions}, input_tokens