Laya / AX637_BACKEND /python /infer.py
yongqiang
Add AX620E and AX637 backends
14a0ae5
Raw History Blame Contribute Delete
15.2 kB
#!/usr/bin/env python3
"""Run a packaged Laya AXModel with the native AX637 ax_run_model utility."""
import argparse
import os
import re
import subprocess
import tempfile
import json
import sys
import time
from pathlib import Path
from typing import Any, Dict, Iterable, List, Tuple
os.environ.setdefault("TRANSFORMERS_NO_ADVISORY_WARNINGS", "1")
import numpy as np
from transformers import AutoTokenizer
QTYPES = {"choice": 0, "score": 1, "noul": 2}
QTYPE_NAMES = {value: key for key, value in QTYPES.items()}
def softmax(values: np.ndarray) -> np.ndarray:
values = values.astype(np.float64)
values -= values.max()
result = np.exp(values)
return result / result.sum()
def confidence_from_probs(probabilities: np.ndarray, count: int) -> float:
if count < 2:
return 1.0
values = probabilities[:count]
entropy = -(values * np.log(np.clip(values, 1e-12, 1.0))).sum()
return float(np.clip(1.0 - entropy / np.log(count), 0.0, 1.0))
def temp_bucket(qtype: int, count: int) -> str:
size = "2" if count <= 2 else "3-5" if count <= 5 else "6-10" if count <= 10 else "11+"
return f"{QTYPE_NAMES[int(qtype)]}:{size}"
def render(value: Any) -> str:
if isinstance(value, str):
return value
return json.dumps(value, ensure_ascii=False, separators=(", ", ": "), default=str)
def internal_question(question: Dict[str, Any]) -> Dict[str, Any]:
qtype = question["type"]
criteria = question.get("criteria")
if qtype == "choice" and isinstance(criteria, list):
criteria = {item: None for item in criteria}
instructions = question["instructions"]
if not isinstance(instructions, str):
instructions = json.dumps(instructions)
return {"t": qtype, "ins": instructions, "crit": criteria}
def render_options(question: Dict[str, Any]) -> List[str]:
qtype = question["t"]
criteria = question.get("crit")
if qtype == "choice":
if not isinstance(criteria, dict):
raise ValueError("choice.criteria must be an object or list")
return [
key if value is None or value == "" else f"{key}: {render(value)}"
for key, value in criteria.items()
]
if qtype == "score":
if not isinstance(criteria, list):
raise ValueError("score.criteria must be a list")
return [f"level {index}: {render(value)}" for index, value in enumerate(criteria)]
criteria = criteria or {}
false_value = criteria.get("false")
true_value = criteria.get("true")
return [
"false: " + (render(false_value) if false_value not in (None, "") else "no, the statement does not hold"),
"true: " + (render(true_value) if true_value not in (None, "") else "yes, the statement holds"),
]
def build_sequence(
tokenizer,
state: Any,
question: Dict[str, Any],
max_len: int,
head_max_len: int,
) -> Tuple[List[int], List[int]]:
mask_token = tokenizer.mask_token
options = render_options(question)
instructions = str(question["ins"]).replace(mask_token, " ")
head_ids = tokenizer(
f"{question['t']} question: {instructions}",
add_special_tokens=False,
)["input_ids"]
option_ids = [
[tokenizer.mask_token_id]
+ tokenizer(
" " + option.replace(mask_token, " "),
add_special_tokens=False,
)["input_ids"][:48]
for option in options
]
budget = head_max_len - sum(len(option) for option in option_ids)
if budget < 16:
per_option = max(4, (head_max_len - 16) // max(1, len(option_ids)))
option_ids = [option[:per_option] for option in option_ids]
budget = head_max_len - sum(len(option) for option in option_ids)
head_ids = head_ids[: max(8, budget)]
ids = [tokenizer.cls_token_id] + head_ids + [tokenizer.sep_token_id]
markers = []
for option in option_ids:
markers.append(len(ids))
ids.extend(option)
ids.append(tokenizer.sep_token_id)
room = max(0, max_len - len(ids) - 1)
state_text = state if isinstance(state, str) else json.dumps(state, ensure_ascii=False)
state_ids = tokenizer(
state_text.replace(mask_token, " "),
add_special_tokens=False,
)["input_ids"][:room]
ids = ids + state_ids + [tokenizer.sep_token_id]
return ids[:max_len], [marker for marker in markers if marker < max_len]
def encode_question(
tokenizer,
state: Any,
question: Dict[str, Any],
seq_len: int,
num_options: int,
head_max_len: int,
) -> Dict[str, np.ndarray]:
internal = internal_question(question)
options = render_options(internal)
if not 2 <= len(options) <= num_options:
raise ValueError(
f"question has {len(options)} options; this graph supports 2..{num_options}"
)
ids, markers = build_sequence(tokenizer, state, internal, seq_len, head_max_len)
if len(markers) != len(options):
raise ValueError("an option marker was truncated; shorten the question")
pad = seq_len - len(ids)
return {
"input_ids": np.asarray(
[ids + [tokenizer.pad_token_id] * pad], dtype=np.int32
),
"attention_mask": np.asarray(
[[1] * len(ids) + [0] * pad], dtype=np.int32
),
"marker_pos": np.asarray(
[markers + [0] * (num_options - len(markers))], dtype=np.int32
),
"marker_mask": np.asarray(
[[1] * len(markers) + [0] * (num_options - len(markers))],
dtype=np.int32,
),
"qtype": np.asarray([QTYPES[internal["t"]]], dtype=np.int32),
}
def validate_request(request: Any) -> Dict[str, Any]:
if not isinstance(request, dict):
raise ValueError("request must be a JSON object")
if "state" not in request:
raise ValueError("request is missing required field: state")
questions = request.get("questions")
if not isinstance(questions, dict) or not questions:
raise ValueError("request.questions must be a non-empty object")
for question_id, question in questions.items():
if not isinstance(question, dict):
raise ValueError(f"question {question_id!r} must be an object")
if question.get("type") not in QTYPES:
raise ValueError(
f"question {question_id!r} has unsupported type {question.get('type')!r}"
)
if "instructions" not in question:
raise ValueError(f"question {question_id!r} is missing instructions")
return request
def format_answer(
question: Dict[str, Any],
probabilities: np.ndarray,
confidence: float,
act_probability: float,
latency_ms: float,
) -> Dict[str, Any]:
common = {
"type": question["type"],
"confidence": float(confidence),
"action": {"act_probability": float(act_probability)},
"python_latency_ms": float(latency_ms),
}
if question["type"] == "choice":
labels = list(question["criteria"])
return {
"type": "choice",
"choice": labels[int(probabilities.argmax())],
"probabilities": {
label: float(value) for label, value in zip(labels, probabilities)
},
**{key: common[key] for key in ("confidence", "action", "python_latency_ms")},
}
if question["type"] == "score":
score = float(np.dot(np.arange(len(probabilities)), probabilities))
return {
"type": "score",
"probabilities": {
str(index): float(value)
for index, value in enumerate(probabilities)
},
"legend": {
str(index): value
for index, value in enumerate(question["criteria"])
},
"score": score,
**{key: common[key] for key in ("confidence", "action", "python_latency_ms")},
}
return {
"type": "noul",
"noul": float(probabilities[1]),
**{key: common[key] for key in ("confidence", "action", "python_latency_ms")},
}
class LayaAX637:
def __init__(self, model_dir: Path, runner: Path):
self.model_dir = model_dir.resolve()
config_path = self.model_dir / "config.json"
if not config_path.is_file():
raise FileNotFoundError(config_path)
self.config = json.loads(config_path.read_text(encoding="utf-8"))
self.seq_len = int(self.config.get("sequence_length", 256))
self.num_options = int(self.config.get("num_options", 4))
self.head_max_len = int(self.config.get("head_max_len", 128))
self.temperatures = self.config.get("temperature", [1.0, 1.0, 1.0])
self.temperatures_by_options = self.config.get("temperature_by_options", {})
tokenizer_dir = self.model_dir / self.config.get("tokenizer_dir", "tokenizer")
self.tokenizer = AutoTokenizer.from_pretrained(tokenizer_dir)
self.model_path = self.model_dir / self.config.get("filename_axmodel", "model.axmodel")
if not self.model_path.is_file():
raise FileNotFoundError(self.model_path)
self.runner = runner.resolve()
if not self.runner.is_file():
raise FileNotFoundError(self.runner)
def _run_native(
self, feeds: List[Dict[str, np.ndarray]]
) -> Tuple[List[Tuple[np.ndarray, np.ndarray]], float, float]:
with tempfile.TemporaryDirectory(prefix="laya_ax637_") as temporary:
work_dir = Path(temporary)
input_root = work_dir / "input"
output_root = work_dir / "output"
output_root.mkdir(parents=True)
sample_names = []
for index, feed in enumerate(feeds):
sample_name = f"sample_{index:05d}"
sample_names.append(sample_name)
sample_dir = input_root / sample_name
sample_dir.mkdir(parents=True)
for name, value in feed.items():
np.ascontiguousarray(value).tofile(sample_dir / f"{name}.bin")
list_path = work_dir / "list.txt"
list_path.write_text("\n".join(sample_names) + "\n", encoding="utf-8")
command = [
str(self.runner),
"--model",
str(self.model_path),
"--input-folder",
str(input_root),
"--output-folder",
str(output_root),
"--list",
str(list_path),
"--repeat",
"1",
]
started = time.perf_counter()
environment = os.environ.copy()
board_library_dir = "/opt/lib"
current_library_path = environment.get("LD_LIBRARY_PATH", "")
environment["LD_LIBRARY_PATH"] = (
board_library_dir
if not current_library_path
else board_library_dir + ":" + current_library_path
)
completed = subprocess.run(
command,
stdout=subprocess.PIPE,
stderr=subprocess.STDOUT,
text=True,
env=environment,
)
wall_ms = (time.perf_counter() - started) * 1000.0
if completed.returncode != 0:
raise RuntimeError(
f"ax_run_model failed ({completed.returncode}):\n{completed.stdout}"
)
match = re.search(r"avg\s*=\s*([0-9.]+)\s*ms", completed.stdout)
average_npu_ms = float(match.group(1)) if match else float("nan")
outputs = []
for sample_name in sample_names:
output_dir = output_root / sample_name
logits = np.fromfile(output_dir / "logits.bin", dtype=np.float32)
act_logits = np.fromfile(
output_dir / "act_logits.bin", dtype=np.float32
)
if logits.size != 4 or act_logits.size != 2:
raise RuntimeError(
f"Unexpected output sizes for {sample_name}: "
f"logits={logits.size}, act_logits={act_logits.size}"
)
outputs.append(
(logits.reshape(1, 4), act_logits.reshape(1, 2))
)
return outputs, average_npu_ms, wall_ms
def predict(self, request: Dict[str, Any]) -> Dict[str, Any]:
request = validate_request(request)
items = list(request["questions"].items())
feeds = [
encode_question(
self.tokenizer,
request["state"],
question,
self.seq_len,
self.num_options,
self.head_max_len,
)
for _, question in items
]
outputs, average_npu_ms, native_wall_ms = self._run_native(feeds)
answers = {}
for (name, question), inputs, (logits, act_logits) in zip(
items, feeds, outputs
):
count = int(inputs["marker_mask"].sum())
qtype = QTYPES[question["type"]]
scale = self.temperatures_by_options.get(
temp_bucket(qtype, count), self.temperatures[qtype]
)
probabilities = softmax(
logits[0, :count] / max(1e-3, float(scale))
)
act_probability = softmax(act_logits[0])[0]
confidence = (
max(float(probabilities[1]), 1.0 - float(probabilities[1]))
if question["type"] == "noul"
else confidence_from_probs(probabilities, count)
)
answers[name] = format_answer(
question,
probabilities,
confidence,
float(act_probability),
average_npu_ms,
)
answers[name]["npu_latency_ms"] = answers[name].pop(
"python_latency_ms"
)
return {
"model": self.config.get("model_name", self.model_dir.name),
"backend": "ax_run_model",
"answers": answers,
"average_npu_latency_ms": average_npu_ms,
"total_npu_latency_ms": average_npu_ms * len(items),
"native_process_wall_ms": native_wall_ms,
}
def parse_args() -> argparse.Namespace:
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument(
"model_dir",
type=Path,
help="Packaged checkpoint directory, such as multilingual.",
)
parser.add_argument(
"--input",
type=Path,
required=True,
help="Request JSON file.",
)
parser.add_argument(
"--runner",
type=Path,
default=Path("/opt/bin/ax_run_model"),
help="Path to the board-native ax_run_model executable.",
)
return parser.parse_args()
def main() -> None:
args = parse_args()
request = json.loads(args.input.read_text(encoding="utf-8"))
runner = LayaAX637(args.model_dir, args.runner)
result = runner.predict(request)
print(json.dumps(result, indent=2, ensure_ascii=False))
if __name__ == "__main__":
main()