diff --git a/Cargo.lock b/Cargo.lock index 2bb95d2f..b993c0a9 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1123,6 +1123,24 @@ dependencies = [ "libloading", ] +[[package]] +name = "omni-decider-native" +version = "0.1.0" +dependencies = [ + "anyhow", + "axum", + "half", + "memmap2", + "omni-qwen3-5-native", + "omni-runtime", + "safetensors 0.8.0", + "serde", + "serde_json", + "sha2", + "tokenizers 0.22.2", + "tokio", +] + [[package]] name = "omni-jev" version = "0.1.0" diff --git a/Cargo.toml b/Cargo.toml index 242d51ca..57bf6308 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,3 +1,3 @@ [workspace] -members = ["src/frontend", "src/runtime", "src/models/clm", "src/models/cua_s1/native", "src/models/qwen3_5/native", "src/models/open_jev/native", "src/models/laya", "src/backends/cuda"] +members = ["src/frontend", "src/runtime", "src/models/clm", "src/models/cua_s1/native", "src/models/qwen3_5/native", "src/models/open_jev/native", "src/models/laya", "src/backends/cuda", "src/models/decider/native"] resolver = "3" diff --git a/recipe/decider/download_weights.py b/recipe/decider/download_weights.py new file mode 100644 index 00000000..efe74c0a --- /dev/null +++ b/recipe/decider/download_weights.py @@ -0,0 +1,58 @@ +#!/usr/bin/env python3 +"""Download and SHA-256 verify the immutable Decider-2B v11 artifacts.""" +import argparse +import hashlib +import os +from pathlib import Path +import tempfile +import urllib.request + +REVISION = "533964dae8be954c5b5e19fa4948e48408094c1e" +FILES = { + "config.json": (1790, "6cb8daca9fb653c61485ff7452fc068bacd5c27cbee659ecd24b47186b0d1b52"), + "decider_config.json": (1240, "6e4891f2754a1c18a10f8dadb0c04e439e7f79fab0333d56641491bd4a05e722"), + "tokenizer.json": (19989325, "06b9509352d2af50381ab2247e083b80d32d5c0aba91c272ca9ff729b6a0e523"), + "tokenizer_config.json": (1127, "171ecbe7ddae98d11840698f7df2b8d5b4722139db0f0620d3bbf429bd656250"), + "model.safetensors": (3763692048, "acaef2228b134dcdc20cad4ee79219482c927ec819aa3687b9b8a575c338817f"), +} + + +def verify(path, expected): + size, digest = expected + if path.stat().st_size != size: + raise ValueError(f"{path}: size mismatch") + actual = hashlib.sha256() + with path.open("rb") as source: + for block in iter(lambda: source.read(1024 * 1024), b""): + actual.update(block) + if actual.hexdigest() != digest: + raise ValueError(f"{path}: checksum mismatch") + + +def main(): + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("directory", type=Path) + parser.add_argument("--endpoint", default="https://huggingface.co") + args = parser.parse_args() + args.directory.mkdir(parents=True, exist_ok=True) + for name, expected in FILES.items(): + destination = args.directory / name + if destination.exists(): + verify(destination, expected) + else: + fd, temporary = tempfile.mkstemp(prefix=name + ".", suffix=".part", dir=args.directory) + temporary = Path(temporary) + try: + url = f"{args.endpoint.rstrip('/')}/Mapika/decider-2b/resolve/{REVISION}/{name}" + with os.fdopen(fd, "wb") as output, urllib.request.urlopen(url, timeout=120) as response: + for block in iter(lambda: response.read(1024 * 1024), b""): + output.write(block) + verify(temporary, expected) + temporary.replace(destination) + finally: + temporary.unlink(missing_ok=True) + print(f"verified {name}", flush=True) + + +if __name__ == "__main__": + main() diff --git a/src/models/decider/native/Cargo.toml b/src/models/decider/native/Cargo.toml new file mode 100644 index 00000000..465ad45b --- /dev/null +++ b/src/models/decider/native/Cargo.toml @@ -0,0 +1,47 @@ +[package] +name = "omni-decider-native" +version = "0.1.0" +edition = "2024" +publish = false +description = "Native Rust/CUDA text worker for Decider-2B v11" +[dependencies] +anyhow = "1.0.100" +serde = "1.0.229" +serde_json = { version = "1.0.149", features = ["float_roundtrip", "preserve_order", "raw_value"] } +tokenizers = { version = "=0.22.2", default-features = false, features = ["onig"] } +sha2 = "0.10" +half = "2.7.1" +memmap2 = "0.9.9" +safetensors = "0.8.0" +omni-qwen3-5-native = { path = "../../qwen3_5/native" } +omni-runtime = { path = "../../../runtime" } +axum = "0.8.8" +tokio = { version = "1.49.0", features = ["macros", "net", "rt-multi-thread", "sync", "signal"] } + +[[bin]] +name = "omni-decider" +path = "src/main.rs" + +[[test]] +name = "gpu" +path = "../../../../tests/decider/gpu.rs" + +[[test]] +name = "contract" +path = "../../../../tests/decider/contract.rs" + +[[test]] +name = "config" +path = "../../../../tests/decider/config.rs" + +[[test]] +name = "batching" +path = "../../../../tests/decider/batching.rs" + +[[test]] +name = "options" +path = "../../../../tests/decider/options.rs" + +[[example]] +name = "decider-run" +path = "../../../../tests/decider/runner.rs" diff --git a/src/models/decider/native/LICENSE.decider b/src/models/decider/native/LICENSE.decider new file mode 100644 index 00000000..916a53c7 --- /dev/null +++ b/src/models/decider/native/LICENSE.decider @@ -0,0 +1,190 @@ + Apache License + Version 2.0, January 2004 + http://www.apache.org/licenses/ + + TERMS AND CONDITIONS FOR USE, REPRODUCTION, AND DISTRIBUTION + + 1. Definitions. + + "License" shall mean the terms and conditions for use, reproduction, + and distribution as defined by Sections 1 through 9 of this document. + + "Licensor" shall mean the copyright owner or entity authorized by + the copyright owner that is granting the License. + + "Legal Entity" shall mean the union of the acting entity and all + other entities that control, are controlled by, or are under common + control with that entity. For the purposes of this definition, + "control" means (i) the power, direct or indirect, to cause the + direction or management of such entity, whether by contract or + otherwise, or (ii) ownership of fifty percent (50%) or more of the + outstanding shares, or (iii) beneficial ownership of such entity. + + "You" (or "Your") shall mean an individual or Legal Entity + exercising permissions granted by this License. + + "Source" form shall mean the preferred form for making modifications, + including but not limited to software source code, documentation + source, and configuration files. + + "Object" form shall mean any form resulting from mechanical + transformation or translation of a Source form, including but + not limited to compiled object code, generated documentation, + and conversions to other media types. + + "Work" shall mean the work of authorship, whether in Source or + Object form, made available under the License, as indicated by a + copyright notice that is included in or attached to the work + (an example is provided in the Appendix below). + + "Derivative Works" shall mean any work, whether in Source or Object + form, that is based on (or derived from) the Work and for which the + editorial revisions, annotations, elaborations, or other modifications + represent, as a whole, an original work of authorship. For the purposes + of this License, Derivative Works shall not include works that remain + separable from, or merely link (or bind by name) to the interfaces of, + the Work and Derivative Works thereof. + + "Contribution" shall mean any work of authorship, including + the original version of the Work and any modifications or additions + to that Work or Derivative Works thereof, that is intentionally + submitted to Licensor for inclusion in the Work by the copyright owner + or by an individual or Legal Entity authorized to submit on behalf of + the copyright owner. For the purposes of this definition, "submitted" + means any form of electronic, verbal, or written communication sent + to the Licensor or its representatives, including but not limited to + communication on electronic mailing lists, source code control systems, + and issue tracking systems that are managed by, or on behalf of, the + Licensor for the purpose of discussing and improving the Work, but + excluding communication that is conspicuously marked or otherwise + designated in writing by the copyright owner as "Not a Contribution." + + "Contributor" shall mean Licensor and any individual or Legal Entity + on behalf of whom a Contribution has been received by Licensor and + subsequently incorporated within the Work. + + 2. Grant of Copyright License. Subject to the terms and conditions of + this License, each Contributor hereby grants to You a perpetual, + worldwide, non-exclusive, no-charge, royalty-free, irrevocable + copyright license to reproduce, prepare Derivative Works of, + publicly display, publicly perform, sublicense, and distribute the + Work and such Derivative Works in Source or Object form. + + 3. Grant of Patent License. Subject to the terms and conditions of + this License, each Contributor hereby grants to You a perpetual, + worldwide, non-exclusive, no-charge, royalty-free, irrevocable + (except as stated in this section) patent license to make, have made, + use, offer to sell, sell, import, and otherwise transfer the Work, + where such license applies only to those patent claims licensable + by such Contributor that are necessarily infringed by their + Contribution(s) alone or by combination of their Contribution(s) + with the Work to which such Contribution(s) was submitted. If You + institute patent litigation against any entity (including a + cross-claim or counterclaim in a lawsuit) alleging that the Work + or a Contribution incorporated within the Work constitutes direct + or contributory patent infringement, then any patent licenses + granted to You under this License for that Work shall terminate + as of the date such litigation is filed. + + 4. Redistribution. You may reproduce and distribute copies of the + Work or Derivative Works thereof in any medium, with or without + modifications, and in Source or Object form, provided that You + meet the following conditions: + + (a) You must give any other recipients of the Work or + Derivative Works a copy of this License; and + + (b) You must cause any modified files to carry prominent notices + stating that You changed the files; and + + (c) You must retain, in the Source form of any Derivative Works + that You distribute, all copyright, patent, trademark, and + attribution notices from the Source form of the Work, + excluding those notices that do not pertain to any part of + the Derivative Works; and + + (d) If the Work includes a "NOTICE" text file as part of its + distribution, then any Derivative Works that You distribute must + include a readable copy of the attribution notices contained + within such NOTICE file, excluding those notices that do not + pertain to any part of the Derivative Works, in at least one + of the following places: within a NOTICE text file distributed + as part of the Derivative Works; within the Source form or + documentation, if provided along with the Derivative Works; or, + within a display generated by the Derivative Works, if and + wherever such third-party notices normally appear. The contents + of the NOTICE file are for informational purposes only and + do not modify the License. You may add Your own attribution + notices within Derivative Works that You distribute, alongside + or as an addendum to the NOTICE text from the Work, provided + that such additional attribution notices cannot be construed + as modifying the License. + + You may add Your own copyright statement to Your modifications and + may provide additional or different license terms and conditions + for use, reproduction, or distribution of Your modifications, or + for any such Derivative Works as a whole, provided Your use, + reproduction, and distribution of the Work otherwise complies with + the conditions stated in this License. + + 5. Submission of Contributions. Unless You explicitly state otherwise, + any Contribution intentionally submitted for inclusion in the Work + by You to the Licensor shall be under the terms and conditions of + this License, without any additional terms or conditions. + Notwithstanding the above, nothing herein shall supersede or modify + the terms of any separate license agreement you may have executed + with Licensor regarding such Contributions. + + 6. Trademarks. This License does not grant permission to use the trade + names, trademarks, service marks, or product names of the Licensor, + except as required for reasonable and customary use in describing the + origin of the Work and reproducing the content of the NOTICE file. + + 7. Disclaimer of Warranty. Unless required by applicable law or + agreed to in writing, Licensor provides the Work (and each + Contributor provides its Contributions) on an "AS IS" BASIS, + WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or + implied, including, without limitation, any warranties or conditions + of TITLE, NON-INFRINGEMENT, MERCHANTABILITY, or FITNESS FOR A + PARTICULAR PURPOSE. You are solely responsible for determining the + appropriateness of using or redistributing the Work and assume any + risks associated with Your exercise of permissions under this License. + + 8. Limitation of Liability. In no event and under no legal theory, + whether in tort (including negligence), contract, or otherwise, + unless required by applicable law (such as deliberate and grossly + negligent acts) or agreed to in writing, shall any Contributor be + liable to You for damages, including any direct, indirect, special, + incidental, or consequential damages of any character arising as a + result of this License or out of the use or inability to use the + Work (including but not limited to damages for loss of goodwill, + work stoppage, computer failure or malfunction, or any and all + other commercial damages or losses), even if such Contributor + has been advised of the possibility of such damages. + + 9. Accepting Warranty or Additional Liability. While redistributing + the Work or Derivative Works thereof, You may choose to offer, + and charge a fee for, acceptance of support, warranty, or other + liability obligations and/or rights consistent with this License. + However, in accepting such obligations, You may act only on Your + own behalf and on Your sole responsibility, not on behalf of any + other Contributor, and only if You agree to indemnify, defend, and + hold each Contributor harmless for any liability incurred by, or + claims asserted against, such Contributor by reason of your + accepting any such warranty or additional liability. + + END OF TERMS AND CONDITIONS + + Copyright 2026 Mark Marosi + + Licensed under the Apache License, Version 2.0 (the "License"); + you may not use this file except in compliance with the License. + You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + + Unless required by applicable law or agreed to in writing, software + distributed under the License is distributed on an "AS IS" BASIS, + WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + See the License for the specific language governing permissions and + limitations under the License. diff --git a/src/models/decider/native/THIRD_PARTY_NOTICES.md b/src/models/decider/native/THIRD_PARTY_NOTICES.md new file mode 100644 index 00000000..1d7b7fa6 --- /dev/null +++ b/src/models/decider/native/THIRD_PARTY_NOTICES.md @@ -0,0 +1,16 @@ +# Decider native contract attribution + +The request rendering, row planning, token construction and answer formulas in +`src/contract.rs` and `src/processing.rs` adapt Mapika/decider's +`decider/systemone.py`, `decider/prompt.py`, `decider/prompt_fast.py` and +`decider/temperature.py`, revision +`50d0be0d7cb43d2066965ce5fa7f3fe4e489a60f` (decider-ai 1.8.1). +Copyright the Mapika/decider contributors. Licensed under Apache-2.0; +see [the retained license](LICENSE.decider). + +The adaptation uses Rust CPU buffers, the frozen Decider-2B v11 tokenizer and +configuration, explicit native input/admission restrictions, independent rows +and isolated Score levels. Native execution reuses the repository Qwen3.5 CUDA +backbone and the reference BF16 selected tied-embedding projection contract. The model-private JSON renderer follows this repository's shared JSON +helper's Python notation conventions, with stricter integer decoding to prevent +arbitrary_precision feature unification from changing state values. diff --git a/src/models/decider/native/src/batching.rs b/src/models/decider/native/src/batching.rs new file mode 100644 index 00000000..8204f639 --- /dev/null +++ b/src/models/decider/native/src/batching.rs @@ -0,0 +1,82 @@ +//! Bounded packing of independent complete rows within a single admitted request. +use anyhow::{Context, Result, ensure}; +use std::ops::Range; + +#[derive(Clone, Copy, Debug)] +pub struct BatchLimits { + max_rows: usize, + max_tokens: usize, +} +impl Default for BatchLimits { + fn default() -> Self { + Self { + max_rows: 1, + max_tokens: 4096, + } + } +} +impl BatchLimits { + pub fn new(max_rows: usize, max_tokens: usize) -> Result { + ensure!( + (1..=4).contains(&max_rows), + "batch rows must be within 1..=4" + ); + ensure!( + (1..=4096).contains(&max_tokens), + "batch tokens must be within 1..=4096" + ); + Ok(Self { + max_rows, + max_tokens, + }) + } + pub fn from_values(rows: Option<&str>, tokens: Option<&str>) -> Result { + let rows = rows + .map_or(Ok(1), str::parse) + .context("DECIDER_BATCH_MAX_ROWS must be an integer")?; + let tokens = tokens + .map_or(Ok(4096), str::parse) + .context("DECIDER_BATCH_MAX_TOKENS must be an integer")?; + Self::new(rows, tokens) + } + pub fn from_env() -> Result { + let get = |name| match std::env::var(name) { + Ok(value) => Ok(Some(value)), + Err(std::env::VarError::NotPresent) => Ok(None), + Err(e) => Err(e), + }; + let rows = get("DECIDER_BATCH_MAX_ROWS")?; + let tokens = get("DECIDER_BATCH_MAX_TOKENS")?; + Self::from_values(rows.as_deref(), tokens.as_deref()) + } + pub fn max_rows(&self) -> usize { + self.max_rows + } + pub fn max_tokens(&self) -> usize { + self.max_tokens + } + /// Retain original row order; a row above the packing token budget runs alone. + /// Total request/row admission remains the processor/executor's separate limit. + pub fn ranges(&self, lengths: &[usize]) -> Vec> { + let mut ranges = Vec::new(); + let mut start = 0; + let mut tokens = 0usize; + for (index, &length) in lengths.iter().enumerate() { + if index > start + && (index - start == self.max_rows + || tokens + .checked_add(length) + .is_none_or(|total| total > self.max_tokens)) + { + ranges.push(start..index); + start = index; + tokens = 0; + } + tokens = tokens.saturating_add(length); + } + if start < lengths.len() { + ranges.push(start..lengths.len()); + } + ranges + } +} diff --git a/src/models/decider/native/src/checkpoint.rs b/src/models/decider/native/src/checkpoint.rs new file mode 100644 index 00000000..05051eac --- /dev/null +++ b/src/models/decider/native/src/checkpoint.rs @@ -0,0 +1,149 @@ +//! Verify the released immutable artifacts before device initialization. +use anyhow::{Context, Result, ensure}; +use safetensors::{Dtype, SafeTensors}; +use sha2::{Digest, Sha256}; +use std::{fs::File, io::Read, path::Path}; + +pub(crate) const HIDDEN: usize = 2048; +pub(crate) const VOCAB: usize = 248320; +pub(crate) const LABELS: usize = 255; +pub(crate) const PADDED_LABELS: usize = 256; +const FILES: &[(&str, u64, &str)] = &[ + ( + "config.json", + 1790, + "6cb8daca9fb653c61485ff7452fc068bacd5c27cbee659ecd24b47186b0d1b52", + ), + ( + "decider_config.json", + 1240, + "6e4891f2754a1c18a10f8dadb0c04e439e7f79fab0333d56641491bd4a05e722", + ), + ( + "tokenizer.json", + 19989325, + "06b9509352d2af50381ab2247e083b80d32d5c0aba91c272ca9ff729b6a0e523", + ), + ( + "model.safetensors", + 3763692048, + "acaef2228b134dcdc20cad4ee79219482c927ec819aa3687b9b8a575c338817f", + ), +]; + +pub(crate) struct Checkpoint { + pub head: Vec, +} +impl Checkpoint { + pub fn load(dir: &Path, labels: &[u32]) -> Result { + ensure!(labels.len() == LABELS, "expected all 255 label IDs"); + ensure!( + !dir.join("model.safetensors.index.json").exists(), + "sharded checkpoint index would override the verified single-file release" + ); + for &(name, size, digest) in FILES { + verify_file(&dir.join(name), size, digest)?; + } + // In addition to provenance, check every dimension consumed by the shared kernels. + let cfg = omni_qwen3_5_native::model::Config::load(dir)?; + ensure!( + ( + cfg.hidden, + cfg.intermediate, + cfg.heads, + cfg.kv_heads, + cfg.head_dim, + cfg.lin_k_heads, + cfg.lin_v_heads, + cfg.lin_k_dim, + cfg.lin_v_dim + ) == (HIDDEN, 6144, 8, 2, 256, 16, 16, 128, 128) + && cfg.full_attention.len() == 24 + && cfg + .full_attention + .iter() + .enumerate() + .all(|(i, &full)| full == (i % 4 == 3)), + "expected released Decider-2B backbone layout" + ); + let file = File::open(dir.join("model.safetensors"))?; + // SAFETY: checkpoint artifacts must remain immutable while the worker runs. + let mmap = unsafe { memmap2::Mmap::map(&file)? }; + let tensors = SafeTensors::deserialize(&mmap)?; + ensure!( + tensors + .names() + .iter() + .all(|name| tensors.tensor(name).is_ok_and(|v| v.dtype() == Dtype::BF16)), + "expected BF16 tensors" + ); + Ok(Self { + head: selected_rows(&tensors, labels, HIDDEN, VOCAB)?, + }) + } +} + +fn verify_file(path: &Path, size: u64, expected: &str) -> Result<()> { + let mut file = + File::open(path).with_context(|| format!("missing pinned artifact {}", path.display()))?; + ensure!( + file.metadata()?.len() == size, + "{} size mismatch", + path.display() + ); + let mut hash = Sha256::new(); + let mut buffer = vec![0; 1024 * 1024]; + loop { + let n = file.read(&mut buffer)?; + if n == 0 { + break; + } + hash.update(&buffer[..n]); + } + ensure!( + format!("{:x}", hash.finalize()) == expected, + "{} checksum mismatch", + path.display() + ); + Ok(()) +} + +fn selected_rows( + tensors: &SafeTensors<'_>, + ids: &[u32], + hidden: usize, + vocab: usize, +) -> Result> { + ensure!( + !ids.is_empty() && ids.len() <= LABELS && hidden > 0, + "invalid head dimensions" + ); + ensure!( + ids.iter().all(|&id| (id as usize) < vocab) + && ids.iter().collect::>().len() == ids.len(), + "invalid label IDs" + ); + let view = tensors.tensor("model.language_model.embed_tokens.weight")?; + ensure!( + view.dtype() == Dtype::BF16 && view.shape() == [vocab, hidden], + "tied embedding shape/dtype mismatch" + ); + let mut rows = vec![0; ids.len().next_multiple_of(8) * hidden * 2]; + for (i, &id) in ids.iter().enumerate() { + let start = id as usize * hidden * 2; + let row = &view.data()[start..start + hidden * 2]; + ensure!( + row.as_chunks::<2>() + .0 + .iter() + .all(|&bytes| half::bf16::from_le_bytes(bytes).is_finite()), + "nonfinite selected weights" + ); + rows[i * hidden * 2..(i + 1) * hidden * 2].copy_from_slice(row); + } + Ok(rows) +} + +#[cfg(test)] +#[path = "../../../../../tests/decider/checkpoint.rs"] +mod tests; diff --git a/src/models/decider/native/src/config.rs b/src/models/decider/native/src/config.rs new file mode 100644 index 00000000..beb9a07f --- /dev/null +++ b/src/models/decider/native/src/config.rs @@ -0,0 +1,95 @@ +use crate::{contract::Kind, json}; +use anyhow::{Context, Result, ensure}; +use serde_json::Value; +use std::path::Path; + +/// Validated calibration for the released plain independent/isolated readout. +#[derive(Clone, Debug)] +pub struct Config { + temperatures: [f32; 3], +} +impl Config { + pub fn load(dir: impl AsRef) -> Result { + let dir = dir.as_ref(); + let model = json::parse(&std::fs::read(dir.join("config.json"))?)?; + json::object(&model)?; + ensure!( + model["model_type"] == "qwen3_5_text" + && model["hidden_size"] == 2048 + && model["num_hidden_layers"] == 24 + && model["tie_word_embeddings"] == true + && model["vocab_size"] == 248320 + && model["dtype"] == "bfloat16", + "expected Decider-2B v11 text configuration" + ); + Self::from_value(&json::parse(&std::fs::read( + dir.join("decider_config.json"), + )?)?) + } + /// Positive calibration overrides are supported, with missing types using fallback. + /// Nonreleased prompt/readout modes and option-dependent temperatures are rejected. + pub fn from_value(value: &Value) -> Result { + let c = json::object(value)?; + ensure!( + c.get("version").and_then(Value::as_str) == Some("2b-v11"), + "expected 2b-v11 configuration" + ); + ensure!( + c.get("layout").and_then(Value::as_str) == Some("plain"), + "only plain layout supported" + ); + for (key, want) in [ + ("chat_template", false), + ("schema_first", false), + ("schema_first_trained", false), + ("neutralize_none", false), + ("isolated_levels", true), + ] { + json::default_mode(c, key, want)?; + } + ensure!( + c.get("max_options").and_then(Value::as_u64) == Some(255) + && c.get("max_state_tokens").and_then(Value::as_u64) == Some(32768), + "unsupported option/state limits" + ); + for key in [ + "temperature_by_options", + "temperature_schema_first", + "temperature_schema_first_by_type", + ] { + ensure!(!c.contains_key(key), "unsupported calibration {key}"); + } + for key in ["neutralize_none", "isolated_levels"] { + json::required(c, key)?; + } + let raw = json::required(c, "temperature")?; + let scalar = if let Some(text) = raw.as_str() { + serde_json::json!(json::python_float(text).context("temperature must be numeric")?) + } else { + raw.clone() + }; + let fallback = positive(&scalar)?; + let mut temperatures = [fallback; 3]; + if let Some(m) = c.get("temperature_by_type").filter(|v| !v.is_null()) { + for (key, value) in json::object(m)? { + let kind = Kind::parse(key)?; + ensure!(key != "bool", "unknown calibration type bool"); + temperatures[kind.index()] = positive(value)?; + } + } + Ok(Self { temperatures }) + } + pub fn temperature(&self, kind: Kind) -> f32 { + self.temperatures[kind.index()] + } +} +fn positive(value: &Value) -> Result { + // The serving by-type path materializes its temperature list as FP32. + let x = value.as_f64().context("temperature must be numeric")?; + let rounded = x as f32; + ensure!( + x.is_finite() && x > 0.0 && rounded.is_finite() && rounded > 0.0, + "temperature must be finite and positive in FP32" + ); + Ok(rounded) +} diff --git a/src/models/decider/native/src/contract.rs b/src/models/decider/native/src/contract.rs new file mode 100644 index 00000000..96ad5022 --- /dev/null +++ b/src/models/decider/native/src/contract.rs @@ -0,0 +1,283 @@ +//! Adapted from Mapika/decider systemone.py at 50d0be0; Apache-2.0. +use crate::json; +use anyhow::{Context, Result, bail, ensure}; +use serde_json::{Map, Value}; + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum Kind { + Choice, + Noul, + Score, +} +impl Kind { + pub(crate) fn parse(t: &str) -> Result { + match t { + "choice" => Ok(Self::Choice), + "noul" | "bool" => Ok(Self::Noul), + "score" => Ok(Self::Score), + _ => bail!("unknown question type {t:?}"), + } + } + pub(crate) fn index(self) -> usize { + match self { + Self::Choice => 0, + Self::Noul => 1, + Self::Score => 2, + } + } +} +#[derive(Clone, Debug)] +pub(crate) struct Question { + pub id: String, + pub kind: Kind, + pub text: String, + pub options: Vec, + pub names: Vec, + pub legend: Vec, +} +fn empty(v: &Value) -> bool { + v.is_null() || v.as_str() == Some("") +} +fn annotate(v: &Value) -> Value { + match v { + Value::Array(a) => Value::Array( + a.iter() + .enumerate() + .map(|(i, x)| { + let x = annotate(x); + if a.len() < 8 { + x + } else { + let mut m = Map::new(); + m.insert("_index".into(), Value::from(i)); + if let Value::Object(fields) = x { + for (k, v) in fields { + m.insert(k, v); + } + } else { + m.insert("value".into(), x); + } + Value::Object(m) + } + }) + .collect(), + ), + Value::Object(m) => { + Value::Object(m.iter().map(|(k, v)| (k.clone(), annotate(v))).collect()) + } + _ => v.clone(), + } +} +pub(crate) fn render_state(v: &Value) -> String { + if v.is_string() { + json::text(v) + } else { + json::dumps(&annotate(v)) + } +} +pub(crate) fn render_question(id: &str, spec: &Value) -> Result { + let s = json::object(spec).context("question definition must be an object")?; + json::default_mode(s, "isolated", true)?; + let kind = Kind::parse( + s.get("type") + .map(|v| v.as_str().context("type must be text")) + .transpose()? + .unwrap_or("choice"), + )?; + let none = Value::Null; + let blank = Value::String(String::new()); + let criteria = s + .get("criteria") + .or_else(|| s.get("options")) + .unwrap_or(&none); + let raw = s + .get("instructions") + .or_else(|| s.get("question")) + .unwrap_or(&blank); + let text = if kind == Kind::Noul && empty(raw) { + let described = criteria.as_object().is_some_and(|m| { + ["true", "false"] + .iter() + .any(|k| m.get(*k).is_some_and(|v| !empty(v))) + }); + ensure!( + described, + "noul without instructions requires true or false description" + ); + "Which answer fits the context?".into() + } else { + json::text(raw) + }; + ensure!(!text.is_empty(), "question without instructions"); + let mut names = Vec::new(); + let mut options = Vec::new(); + let mut legend = Vec::new(); + match kind { + Kind::Choice => { + let fields = if let Some(a) = criteria.as_array() { + let mut fields = Map::new(); + for v in a { + let name = v.as_str().context( + "Choice list alias requires string names; use a map for other JSON values", + )?; + fields.insert(name.into(), Value::Null); + } + fields + } else { + json::object(criteria) + .context("Choice requires map or string list")? + .clone() + }; + ensure!( + (2..=255).contains(&fields.len()), + "Choice requires 2..255 options" + ); + for (name, v) in fields { + options.push(if empty(&v) { + name.clone() + } else { + format!("{name}: {}", json::text(&v)) + }); + names.push(name); + } + } + Kind::Score => { + let levels = if let Some(a) = criteria.as_array() { + a.clone() + } else { + let mut ordered = json::object(criteria)? + .iter() + .map(|(k, v)| { + let number = + json::python_float(k).context("Score legend keys must be numeric")?; + ensure!(number.is_finite(), "Score legend keys must be finite"); + Ok((number, v.clone())) + }) + .collect::>>()?; + ordered.sort_by(|a, b| a.0.partial_cmp(&b.0).expect("finite key")); + ordered.into_iter().map(|(_, v)| v).collect() + }; + ensure!( + (2..=10).contains(&levels.len()), + "Score requires 2..10 levels" + ); + legend = levels.iter().map(json::text).collect(); + for (i, v) in legend.iter().enumerate() { + names.push(i.to_string()); + options.push(format!("{i}: {v}")); + } + } + Kind::Noul => { + ensure!( + criteria.is_null() || criteria.is_object(), + "Noul criteria must be a map" + ); + for (key, label) in [("false", "no"), ("true", "yes")] { + let v = criteria.get(key).unwrap_or(&none); + options.push(if empty(v) { + label.into() + } else { + format!("{label}: {}", json::text(v)) + }); + names.push(key.into()); + } + } + } + Ok(Question { + id: id.into(), + kind, + text, + options, + names, + legend, + }) +} +pub(crate) fn isolated_text(q: &Question, level: usize) -> String { + use crate::text::{decimal, whitespace}; + let original = &q.legend[level]; + let trimmed = original.trim_start_matches(whitespace); + let unsigned = trimmed.strip_prefix('-').unwrap_or(trimmed); + let digits = unsigned + .char_indices() + .find(|(_, c)| decimal(*c).is_none()) + .map_or(unsigned.len(), |(i, _)| i); + let stripped = if digits > 0 { + unsigned[digits..] + .trim_start_matches(whitespace) + .strip_prefix(':') + .map(|v| v.trim_start_matches(whitespace)) + } else { + None + }; + format!( + "{}\nProposed answer: {}\nDoes the proposed answer fit?", + q.text, + stripped.unwrap_or(original) + ) +} + +fn normalized(p: &[f64]) -> Vec { + let sum: f64 = p.iter().sum(); + if sum == 0.0 { + vec![1.0 / p.len() as f64; p.len()] + } else { + p.iter().map(|v| v / sum).collect() + } +} +fn mode(p: &[f64]) -> usize { + (1..p.len()).fold(0, |m, i| if p[i] > p[m] { i } else { m }) +} +pub(crate) fn answer(q: &Question, probabilities: &[f64]) -> Value { + use crate::math::round; + use serde_json::json; + let sum: f64 = probabilities.iter().take(q.options.len()).sum(); + let divisor = if sum == 0.0 { 1.0 } else { sum }; + let p: Vec = probabilities + .iter() + .take(q.options.len()) + .map(|v| v / divisor) + .collect(); + if q.kind == Kind::Noul { + return json!({"type":"noul","noul":round(p[1],4)}); + } + let j = mode(&p); + let n = p.len(); + let entropy: f64 = p.iter().filter(|v| **v > 0.0).map(|v| -v * v.ln()).sum(); + let certainty = (1.0 - entropy / (n as f64).ln()).max(0.0); + // Confidence's zero-mass uniform fallback never replaces emitted probabilities. + let norm = normalized(&p); + let confidence = if q.kind == Kind::Choice { + ((n as f64 * norm[mode(&norm)] - 1.0) / (n - 1) as f64).clamp(0.0, 1.0) + } else { + let modal = mode(&norm); + let spread: f64 = norm + .iter() + .enumerate() + .map(|(i, v)| v * i.abs_diff(modal) as f64) + .sum(); + let center = (n - 1) as f64 / 2.0; + let uniform = (0..n).map(|i| (i as f64 - center).abs()).sum::() / n as f64; + (1.0 - spread / uniform).clamp(0.0, 1.0) + }; + let probabilities: Map = q + .names + .iter() + .cloned() + .zip(p.iter().map(|v| json!(round(*v, 4)))) + .collect(); + let mut out = json!({"type":if q.kind==Kind::Choice {"choice"} else {"score"},"confidence":round(confidence,4),"x_p_max":round(p[j],4),"certainty":round(certainty,4),"probabilities":probabilities}); + if q.kind == Kind::Choice { + out["choice"] = json!(q.names[j]); + } else { + let score: f64 = p.iter().enumerate().map(|(i, v)| i as f64 * v).sum(); + out["score"] = json!(round(score, 2)); + out["legend"] = Value::Object( + q.legend + .iter() + .enumerate() + .map(|(i, v)| (i.to_string(), json!(v))) + .collect(), + ); + } + out +} diff --git a/src/models/decider/native/src/engine.rs b/src/models/decider/native/src/engine.rs new file mode 100644 index 00000000..d7e39be8 --- /dev/null +++ b/src/models/decider/native/src/engine.rs @@ -0,0 +1,89 @@ +//! Assemble verified artifacts, model-owned processing, and serial eager execution. +use crate::{ + Config, Kind, Limits, Processor, batching::BatchLimits, checkpoint::Checkpoint, + executor::Executor, +}; +use anyhow::Result; +use omni_runtime::SerialScheduler; +use serde_json::{Value, json}; +use std::{path::Path, sync::Arc}; + +pub struct Engine { + pub processor: Arc, + pub executor: Executor, + pub scheduler: SerialScheduler, + metadata: Value, +} +impl Engine { + pub async fn load(dir: &Path, library: &Path) -> Result { + Self::load_with_options( + dir, + library, + BatchLimits::from_env()?, + crate::options::graph_env()?, + ) + .await + } + pub async fn load_with_batch( + dir: &Path, + library: &Path, + batching: BatchLimits, + ) -> Result { + Self::load_with_options(dir, library, batching, false).await + } + pub async fn load_with_options( + dir: &Path, + library: &Path, + batching: BatchLimits, + graph: bool, + ) -> Result { + let directory = dir.to_owned(); + let (processor,checkpoint,mut metadata) = tokio::task::spawn_blocking(move || -> Result<_> { + let processor = Processor::load(&directory,Limits::default())?; + let labels: Vec = processor.labels().iter().map(|label| label.id).collect(); + let checkpoint = Checkpoint::load(&directory,&labels)?; + let config = Config::load(&directory)?; + let metadata = json!({"model":crate::MODEL_ID,"checkpoint_revision":crate::CHECKPOINT_REVISION,"reference_revision":crate::RUNTIME_REVISION,"execution":"eager","dtype":"bfloat16","temperatures":{"choice":config.temperature(Kind::Choice),"score":config.temperature(Kind::Score),"noul":config.temperature(Kind::Noul)}}); + Ok((processor,checkpoint,metadata)) + }).await??; + let labels = processor.labels().iter().map(|label| label.id).collect(); + let executor = Executor::load(dir, library, checkpoint, labels, batching, graph).await?; + metadata["batch_limits"] = + json!({"rows":batching.max_rows(),"tokens":batching.max_tokens()}); + let engine = Self { + processor: Arc::new(processor), + executor, + scheduler: SerialScheduler::default(), + metadata, + }; + // Loading returns only after a real complete model decision, before any socket binds. + engine.predict(br#"{"state":"The worker is initialized.","questions":{"ready":{"type":"choice","instructions":"Choose the next action.","criteria":{"continue":"Continue","stop":"Stop"}}}}"#).await?; + Ok(engine) + } + pub async fn predict(&self, raw: &[u8]) -> Result { + let prepared = self.processor.prepare(raw)?; + let logits = self + .executor + .execute(&self.scheduler, prepared.rows) + .await?; + prepared.context.finish(logits) + } + pub fn health(&self) -> Value { + let mut metadata = self.metadata.clone(); + let stats = self.executor.graph_stats(); + metadata["execution"] = json!(if !self.executor.is_ready() { + "unavailable" + } else if stats.enabled { + "graph" + } else { + "eager" + }); + metadata["graph"] = json!(stats); + metadata["status"] = json!(if self.executor.is_ready() { + "ready" + } else { + "unavailable" + }); + metadata + } +} diff --git a/src/models/decider/native/src/executor.rs b/src/models/decider/native/src/executor.rs new file mode 100644 index 00000000..c5b8aa15 --- /dev/null +++ b/src/models/decider/native/src/executor.rs @@ -0,0 +1,311 @@ +//! Eager complete-row prefills and a selected tied-embedding CUDA readout. +use crate::{ + Limits, RowInput, + batching::BatchLimits, + checkpoint::{Checkpoint, HIDDEN, LABELS, PADDED_LABELS, VOCAB}, +}; +use anyhow::{Context, Result, ensure}; +use half::bf16; +use omni_qwen3_5_native::{ + cuda::{self, DeviceBuffer, Stream}, + model::{GraphStats, Model}, +}; +use omni_runtime::SerialScheduler; +use std::{ + ffi::c_void, + path::Path, + sync::{ + Arc, Mutex, + atomic::{AtomicBool, Ordering}, + }, +}; + +struct ProjectionStream(Stream); +impl Drop for ProjectionStream { + fn drop(&mut self) { + // SAFETY: this owner destroys its created stream after Head synchronizes it. + unsafe { + (cuda::api().cs1_stream_destroy)(self.0); + } + } +} +struct Gemm(*mut c_void); +// SAFETY: only the serialized executor uses this handle; each call selects device 0. +unsafe impl Send for Gemm {} +impl Drop for Gemm { + fn drop(&mut self) { + // SAFETY: handle created by cs1_gemm_create, owned exclusively here. + unsafe { + (cuda::api().cs1_gemm_destroy)(self.0); + } + } +} +struct Head { + stream: ProjectionStream, + gemm: Gemm, + weights: DeviceBuffer, + input: DeviceBuffer, + output: DeviceBuffer, + capacity: usize, +} +impl Drop for Head { + fn drop(&mut self) { + let _ = cuda::set_device(0); + let _ = cuda::synchronize(self.stream.0); + } +} +impl Head { + fn new(weights: &[u8], capacity: usize) -> Result { + ensure!((1..=4).contains(&capacity), "invalid head batch capacity"); + ensure!( + weights.len() == PADDED_LABELS * HIDDEN * 2, + "invalid selected head buffer" + ); + let stream = ProjectionStream(cuda::new_stream()?); + // SAFETY: null-checked immediately, then retained in the owning guard. + let gemm = Gemm(unsafe { (cuda::api().cs1_gemm_create)(32 << 20) }); + ensure!(!gemm.0.is_null(), "Decider cuBLASLt setup failed"); + let head = Self { + stream, + gemm, + weights: DeviceBuffer::new(weights.len())?, + input: DeviceBuffer::new(capacity * HIDDEN * 2)?, + output: DeviceBuffer::new(capacity * PADDED_LABELS * 2)?, + capacity, + }; + // SAFETY: the device allocation has exactly weights.len() bytes. + unsafe { + cuda::upload(head.weights.at(0), weights, head.stream.0)?; + } + Ok(head) + } + fn project(&self, hidden: &[f32], count: usize) -> Result> { + Ok(self.project_batch(&[hidden], &[count])?.pop().unwrap()) + } + fn project_batch(&self, hidden: &[&[f32]], counts: &[usize]) -> Result>> { + let rows = hidden.len(); + ensure!( + rows > 0 + && rows <= self.capacity + && rows == counts.len() + && hidden + .iter() + .all(|h| h.len() == HIDDEN && h.iter().all(|x| x.is_finite())) + && counts.iter().all(|n| (2..=LABELS).contains(n)), + "invalid hidden batch/readout" + ); + let input: Vec = hidden + .iter() + .flat_map(|h| h.iter()) + .flat_map(|&x| bf16::from_f32(x).to_le_bytes()) + .collect(); + // SAFETY: [rows,HIDDEN] * [256,HIDDEN]^T -> [rows,256]; rows <= capacity. + // The final zero weight row is padding and never participates in normalization. + unsafe { + cuda::upload(self.input.at(0), &input, self.stream.0)?; + cuda::check( + (cuda::api().cs1_gemm)( + self.gemm.0, + self.input.at(0), + self.weights.at(0), + self.output.at(0), + rows as i32, + PADDED_LABELS as i32, + HIDDEN as i32, + PADDED_LABELS as i32, + self.stream.0, + ), + "Decider label projection", + )?; + } + let mut output = vec![0; rows * PADDED_LABELS * 2]; + // SAFETY: every returned row is inside the allocation; download synchronizes GEMM. + unsafe { + cuda::download(&mut output, self.output.at(0), self.stream.0)?; + } + output + .as_chunks::<{ PADDED_LABELS * 2 }>() + .0 + .iter() + .zip(counts) + .map(|(row, &count)| { + let logits: Vec = row.as_chunks::<2>().0[..count] + .iter() + .map(|&b| bf16::from_le_bytes(b).to_f32()) + .collect(); + ensure!( + logits.iter().all(|x| x.is_finite()), + "nonfinite candidate logits" + ); + Ok(logits) + }) + .collect() + } +} +struct Loaded { + model: Model, + head: Head, + batching: BatchLimits, +} +impl Loaded { + fn execute(&mut self, rows: &[RowInput]) -> Result>> { + cuda::set_device(0)?; + let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { + let lengths: Vec = rows.iter().map(|row| row.ids.len()).collect(); + let mut logits = Vec::with_capacity(rows.len()); + for range in self.batching.ranges(&lengths) { + let batch = &rows[range]; + if batch.len() == 1 { + let hidden = self.model.forward(&batch[0].ids)?; + logits.push(self.head.project(&hidden, batch[0].candidate_ids.len())?); + } else { + let inputs: Vec<&[u32]> = batch.iter().map(|row| row.ids.as_slice()).collect(); + let hidden = self.model.forward_batch(&inputs)?; + ensure!( + hidden.len() == batch.len(), + "backbone batch output count mismatch" + ); + let hidden: Vec<&[f32]> = hidden.iter().map(Vec::as_slice).collect(); + let counts: Vec = + batch.iter().map(|row| row.candidate_ids.len()).collect(); + logits.extend(self.head.project_batch(&hidden, &counts)?); + } + } + Ok(logits) + })) + .map_err(|_| anyhow::anyhow!("Decider execution panicked")) + .and_then(|result| result); + // Complete both streams on success and failure before admission is released. + let model_sync = self.model.synchronize(); + let head_sync = cuda::synchronize(self.head.stream.0); + model_sync?; + head_sync?; + result + } +} + +pub struct Executor { + loaded: Arc>>, + labels: Vec, + ready: Arc, + graph_stats: Arc>, +} +impl Executor { + pub(crate) async fn load( + dir: &Path, + library: &Path, + checkpoint: Checkpoint, + labels: Vec, + batching: BatchLimits, + graph: bool, + ) -> Result { + let (dir, library) = (dir.to_owned(), library.to_owned()); + let loaded = tokio::task::spawn_blocking(move || -> Result { + let model = Model::load_with_graph(&dir, &library, graph)?; + let head = Head::new(&checkpoint.head, batching.max_rows())?; + Ok(Loaded { + model, + head, + batching, + }) + }) + .await??; + let graph_stats = Arc::new(Mutex::new(loaded.model.graph_stats())); + Ok(Self { + graph_stats, + loaded: Arc::new(Mutex::new(Some(loaded))), + labels, + ready: Arc::new(AtomicBool::new(true)), + }) + } + pub fn graph_stats(&self) -> GraphStats { + *self + .graph_stats + .lock() + .unwrap_or_else(|error| error.into_inner()) + } + pub fn is_ready(&self) -> bool { + self.ready.load(Ordering::Acquire) + } + + /// One admitted unit holds every row and its GPU head; no partial answers escape. + /// Cancellation after dispatch retains resources and admission until synchronization. + pub async fn execute( + &self, + scheduler: &SerialScheduler, + rows: Vec, + ) -> Result>> { + validate_rows(&rows, &self.labels)?; + ensure!(self.is_ready(), "executor unavailable"); + if rows.is_empty() { + return Ok(Vec::new()); + } + let (loaded, ready, stats) = ( + self.loaded.clone(), + self.ready.clone(), + self.graph_stats.clone(), + ); + scheduler + .run(move || { + let mut guard = match loaded.lock() { + Ok(guard) => guard, + Err(_) => { + let mut snapshot = stats.lock().unwrap_or_else(|error| error.into_inner()); + snapshot.enabled = false; + snapshot.cached_shapes = 0; + ready.store(false, Ordering::Release); + return Err(anyhow::anyhow!("poisoned executor")); + } + }; + let loaded = guard.as_mut().context("executor unavailable")?; + let result = loaded.execute(&rows); + let mut snapshot = loaded.model.graph_stats(); + if result.is_err() { + snapshot.enabled = false; + snapshot.cached_shapes = 0; + } + *stats.lock().unwrap_or_else(|error| error.into_inner()) = snapshot; + if result.is_err() { + ready.store(false, Ordering::Release); + // Loaded::execute synchronized both streams, including failure paths. + // Retire failed device state before releasing the scheduler permit. + guard.take(); + } + result + }) + .await + } +} + +fn validate_rows(rows: &[RowInput], labels: &[u32]) -> Result<()> { + let limits = Limits::default(); + ensure!(rows.len() <= limits.max_rows, "too many executor rows"); + let mut total = 0usize; + for row in rows { + let n = row.ids.len(); + ensure!( + n > 0 && n <= limits.max_row_tokens && row.readout_position == n - 1, + "invalid complete row/readout position" + ); + ensure!( + row.ids.iter().all(|&id| (id as usize) < VOCAB), + "token outside vocabulary" + ); + let count = row.candidate_ids.len(); + ensure!( + (2..=LABELS).contains(&count) + && labels.get(..count) == Some(row.candidate_ids.as_slice()), + "candidate IDs differ from tied head ordering" + ); + total = total.checked_add(n).context("executor token overflow")?; + ensure!( + total <= limits.max_request_tokens, + "executor processed-token budget exceeded" + ); + } + Ok(()) +} + +#[cfg(test)] +#[path = "../../../../../tests/decider/executor.rs"] +mod tests; diff --git a/src/models/decider/native/src/json.rs b/src/models/decider/native/src/json.rs new file mode 100644 index 00000000..f73819cc --- /dev/null +++ b/src/models/decider/native/src/json.rs @@ -0,0 +1,176 @@ +//! Model-private strict JSON decoding avoids arbitrary_precision feature-unification traps. +use anyhow::{Context, Result, bail, ensure}; +use serde::de::{self, MapAccess, SeqAccess, Visitor}; +use serde_json::{Map, Number, Value, value::RawValue}; +use std::fmt; + +pub fn parse(raw: &[u8]) -> Result { + let raw: Box = serde_json::from_slice(raw).context("invalid request JSON")?; + parse_value(&raw, 0) +} +fn parse_value(raw: &RawValue, depth: usize) -> Result { + let text = raw.get(); + if !matches!(text.as_bytes()[0], b'{' | b'[') { + if text.as_bytes()[0].is_ascii_digit() || text.starts_with('-') { + let number = if text.contains(['.', 'e', 'E']) { + let x: f64 = text.parse()?; + Number::from_f64(x).context("nonfinite or out-of-range JSON float")? + } else if text.starts_with('-') { + Number::from( + text.parse::() + .context("JSON integer outside i64/u64 range")?, + ) + } else { + Number::from( + text.parse::() + .context("JSON integer outside i64/u64 range")?, + ) + }; + return Ok(Value::Number(number)); + } + return Ok(serde_json::from_str(text)?); + } + ensure!(depth < 127, "JSON nesting exceeds supported depth"); + struct Container(usize); + impl<'de> Visitor<'de> for Container { + type Value = Value; + fn expecting(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.write_str("JSON array or object") + } + fn visit_seq>(self, mut seq: A) -> std::result::Result { + let mut out = Vec::new(); + while let Some(raw) = seq.next_element::>()? { + out.push(parse_value(&raw, self.0 + 1).map_err(de::Error::custom)?); + } + Ok(Value::Array(out)) + } + fn visit_map>(self, mut map: A) -> std::result::Result { + let mut out = Map::new(); + while let Some((key, raw)) = map.next_entry::>()? { + if out.contains_key(&key) { + return Err(de::Error::custom(format!("duplicate JSON key {key:?}"))); + } + out.insert( + key, + parse_value(&raw, self.0 + 1).map_err(de::Error::custom)?, + ); + } + Ok(Value::Object(out)) + } + } + use serde::Deserializer; + Ok(serde_json::Deserializer::from_str(text).deserialize_any(Container(depth))?) +} +pub fn object(value: &Value) -> Result<&Map> { + value.as_object().context("expected JSON object") +} +pub fn text(value: &Value) -> String { + value + .as_str() + .map(str::to_owned) + .unwrap_or_else(|| dumps(value)) +} +pub fn dumps(value: &Value) -> String { + match value { + Value::String(s) => serde_json::to_string(s).expect("string serialization"), + Value::Array(a) => format!("[{}]", a.iter().map(dumps).collect::>().join(", ")), + Value::Object(m) => format!( + "{{{}}}", + m.iter() + .map(|(k, v)| format!( + "{}: {}", + serde_json::to_string(k).expect("key serialization"), + dumps(v) + )) + .collect::>() + .join(", ") + ), + Value::Number(n) => { + let s = n.to_string(); + if s.contains(['.', 'e', 'E']) { + float_repr(n.as_f64().expect("validated float")) + } else { + s + } + } + _ => value.to_string(), + } +} +// Python repr notation thresholds, exponent sign/padding, and integer-valued floats. +// Uses Rust's shortest round-trip digits; rare shortest-decimal tie choices may differ. +fn float_repr(x: f64) -> String { + if x == 0.0 { + return if x.is_sign_negative() { "-0.0" } else { "0.0" }.into(); + } + let sci = format!("{x:e}"); + let (m, e) = sci.split_once('e').expect("scientific exponent"); + let exponent: i32 = e.parse().expect("integer exponent"); + let neg = m.starts_with('-'); + let digits = m.trim_start_matches('-').replace('.', ""); + let decpt = exponent + 1; + let sign = if neg { "-" } else { "" }; + if decpt <= -4 || decpt > 16 { + let fraction = if digits.len() > 1 { + format!(".{}", &digits[1..]) + } else { + String::new() + }; + format!( + "{sign}{}{fraction}e{}{:02}", + &digits[..1], + if exponent < 0 { '-' } else { '+' }, + exponent.unsigned_abs() + ) + } else if decpt <= 0 { + format!("{sign}0.{}{}", "0".repeat((-decpt) as usize), digits) + } else if decpt as usize >= digits.len() { + format!( + "{sign}{}{}.0", + digits, + "0".repeat(decpt as usize - digits.len()) + ) + } else { + format!( + "{sign}{}.{}", + &digits[..decpt as usize], + &digits[decpt as usize..] + ) + } +} +pub fn required<'a>(m: &'a Map, key: &str) -> Result<&'a Value> { + m.get(key).with_context(|| format!("missing {key}")) +} +pub fn default_mode(m: &Map, key: &str, expected: bool) -> Result<()> { + if let Some(v) = m.get(key) + && v.as_bool() != Some(expected) + { + bail!("unsupported {key}: expected {expected}"); + } + Ok(()) +} + +// Python float(string) accepts surrounding Python whitespace, digit separators, +// and Unicode decimal (Nd) digits. Normalize these before Rust's numeric parser. +pub fn python_float(key: &str) -> Result { + let normalized: String = key + .trim_matches(crate::text::float_whitespace) + .chars() + .map(|c| crate::text::decimal(c).unwrap_or(c)) + .collect(); + let bytes = normalized.as_bytes(); + for (i, c) in bytes.iter().enumerate() { + if *c == b'_' { + ensure!( + i > 0 + && i + 1 < bytes.len() + && bytes[i - 1].is_ascii_digit() + && bytes[i + 1].is_ascii_digit(), + "invalid digit separator in numeric text" + ); + } + } + normalized + .replace('_', "") + .parse::() + .context("numeric text must parse as float") +} diff --git a/src/models/decider/native/src/lib.rs b/src/models/decider/native/src/lib.rs new file mode 100644 index 00000000..a63c9cf2 --- /dev/null +++ b/src/models/decider/native/src/lib.rs @@ -0,0 +1,40 @@ +//! Decider-2B v11 native CPU processing, eager CUDA execution and HTTP serving. +//! +//! Reference: Mapika/decider 50d0be0 (decider-ai 1.8.1), checkpoint 533964d. +//! Only plain state-first, independent questions and isolated Score levels are supported. +//! The executor must return one candidate-logit vector per complete unpadded row, +//! in prepared order, after the reference BF16 LM projection output rounding. +//! Calibration and whole-question normalization belong to processing, independently of execution. +//! +//! Native input policy: duplicate keys, depth >=128, nonfinite numbers and integer +//! literals outside i64/u64 are rejected. Integer -0 renders as 0; float -0.0 is retained. +//! Choice list aliases accept strings; canonical maps support arbitrary JSON descriptions. +//! Score legend map keys accept Python whitespace, digit separators and Unicode decimal +//! digits, but nonfinite keys are rejected. Scalar numeric-string temperatures are accepted; +//! per-type values must be numbers. Extreme temperatures must remain positive/finite in FP32. +//! Rust shortest round-trip float rendering may differ on rare shortest-decimal ties; +//! CPU FP32 softmax exp/summation can differ from torch by a few ulps. Golden parity is +//! evidence for the recorded corpus, not universal bit-identical numerical output. +//! The CPU Processor checks config/tokenizer compatibility; Engine additionally verifies fixed +//! checkpoint/config/calibration hashes before CUDA loading and a real readiness warmup. +mod config; +mod contract; +mod json; +mod math; +mod processing; +mod text; + +pub use config::Config; +pub use contract::Kind; +pub use processing::{Label, Limits, PreparedRequest, Processor, ResponseContext, RowInput}; +pub const MODEL_ID: &str = "decider-2b-v11"; +pub const RUNTIME_REVISION: &str = "50d0be0d7cb43d2066965ce5fa7f3fe4e489a60f"; +pub const CHECKPOINT_REVISION: &str = "533964dae8be954c5b5e19fa4948e48408094c1e"; + +pub mod batching; +mod checkpoint; +pub mod engine; +pub mod executor; +pub mod serve; + +pub mod options; diff --git a/src/models/decider/native/src/main.rs b/src/models/decider/native/src/main.rs new file mode 100644 index 00000000..983e786b --- /dev/null +++ b/src/models/decider/native/src/main.rs @@ -0,0 +1,27 @@ +use anyhow::{Context, Result}; +use omni_decider_native::{engine::Engine, serve}; +use omni_qwen3_5_native::cuda; +use std::{path::PathBuf, sync::Arc}; + +#[tokio::main] +async fn main() -> Result<()> { + let model = std::env::var_os("DECIDER_MODEL") + .map(PathBuf::from) + .context("set DECIDER_MODEL to pinned Decider-2B v11 directory")?; + let library = std::env::var_os("DECIDER_CUDA_LIB") + .map(PathBuf::from) + .map_or_else(cuda::default_library, Ok)?; + let host = std::env::var("DECIDER_HOST").unwrap_or_else(|_| "127.0.0.1".into()); + let port: u16 = std::env::var("DECIDER_PORT") + .map_or(Ok(8000), |v| v.parse()) + .context("DECIDER_PORT must be a valid port")?; + let engine = Arc::new(Engine::load(&model, &library).await?); + let listener = tokio::net::TcpListener::bind((host.as_str(), port)).await?; + println!("native Decider worker listening on {host}:{port}"); + axum::serve(listener, serve::router(engine)) + .with_graceful_shutdown(async { + let _ = tokio::signal::ctrl_c().await; + }) + .await?; + Ok(()) +} diff --git a/src/models/decider/native/src/math.rs b/src/models/decider/native/src/math.rs new file mode 100644 index 00000000..6ffe7455 --- /dev/null +++ b/src/models/decider/native/src/math.rs @@ -0,0 +1,26 @@ +use anyhow::{Result, ensure}; +/// Divide, subtract, exponentiate and normalize in FP32, then expose those values to Python-style FP64 finishing. +/// CPU exp/summation order can differ from torch's vectorized FP32 softmax by a few ulps. +pub(crate) fn softmax(logits: &[f32], temperature: f32) -> Result> { + ensure!( + logits.len() >= 2 && logits.iter().all(|v| v.is_finite()), + "invalid candidate logits" + ); + let scaled: Vec = logits.iter().map(|v| v / temperature).collect(); + ensure!( + scaled.iter().all(|v| v.is_finite()), + "temperature-scaled logits overflow FP32" + ); + let max = scaled.iter().copied().fold(f32::NEG_INFINITY, f32::max); + let mut exps: Vec = scaled.iter().map(|v| (v - max).exp()).collect(); + let total: f32 = exps.iter().sum(); + exps.iter_mut().for_each(|v| *v /= total); + Ok(exps.into_iter().map(f64::from).collect()) +} +/// Fixed decimal formatting rounds the original binary64 value to even. +/// Multiplying by 10^d before rounding changes cases such as Python round(2.675,2). +pub(crate) fn round(x: f64, digits: usize) -> f64 { + format!("{x:.digits$}") + .parse() + .expect("finite decimal result") +} diff --git a/src/models/decider/native/src/options.rs b/src/models/decider/native/src/options.rs new file mode 100644 index 00000000..d280b680 --- /dev/null +++ b/src/models/decider/native/src/options.rs @@ -0,0 +1,17 @@ +//! Model-owned execution switches, independent from other workers' environment. +use anyhow::{Result, bail}; +pub fn graph_value(value: Option<&str>) -> Result { + match value { + None | Some("0") => Ok(false), + Some("1") => Ok(true), + Some(_) => bail!("DECIDER_GRAPH must be 0 or 1"), + } +} +pub fn graph_env() -> Result { + let value = match std::env::var("DECIDER_GRAPH") { + Ok(value) => Some(value), + Err(std::env::VarError::NotPresent) => None, + Err(error) => return Err(error.into()), + }; + graph_value(value.as_deref()) +} diff --git a/src/models/decider/native/src/processing.rs b/src/models/decider/native/src/processing.rs new file mode 100644 index 00000000..3a910a3d --- /dev/null +++ b/src/models/decider/native/src/processing.rs @@ -0,0 +1,314 @@ +//! Exact segment-aware token compilation and request-local response context. +use crate::{ + Config, + contract::{self, Kind, Question}, + json, +}; +use anyhow::{Context, Result, ensure}; +use serde_json::Value; +use sha2::{Digest, Sha256}; +use std::path::Path; +use tokenizers::Tokenizer; + +const TOKENIZER_SHA256: &str = "06b9509352d2af50381ab2247e083b80d32d5c0aba91c272ca9ff729b6a0e523"; +const MAX_CONTEXT: usize = 32768; +#[derive(Clone, Copy, Debug)] +pub struct Limits { + pub max_rows: usize, + pub max_row_tokens: usize, + /// Sum of complete row lengths, not unique-prefix API usage. + pub max_request_tokens: usize, + pub max_request_bytes: usize, +} +impl Default for Limits { + fn default() -> Self { + Self { + max_rows: 1024, + max_row_tokens: 36864, + max_request_tokens: 1048576, + max_request_bytes: 8 * 1024 * 1024, + } + } +} +#[derive(Clone, Debug)] +pub struct Label { + pub name: String, + pub id: u32, +} +pub struct Processor { + tokenizer: Tokenizer, + labels: Vec