#!/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()