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/README.md b/README.md index 62c30ed0..7ade21c2 100644 --- a/README.md +++ b/README.md @@ -166,6 +166,7 @@ Python screenshot worker, and CLM has a stub-encoder contract recipe: | LAYA | [External worker](recipe/laya/README.md); [Python worker on Apple Silicon (MPS) and CPU](recipe/laya/apple-silicon.md); [CPU checkpoint reader](src/models/laya/README.md); [native Rust/CUDA worker on Hopper](recipe/laya/native/README.md) | | Cua-S1 4B 0.2 (`text` adapter) | [Python worker](recipe/cua_s1/text.md); [native worker](recipe/cua_s1/native.md), CUDA, run on sm_89 | | Cua-S1 4B 0.2 (`multimodal` adapter) | [Python CUDA worker](src/frontend/cua_s1.py); one PNG/JPEG screenshot, `choice`; native screenshot execution remains in progress | +| Decider-2B v11 | [Native Rust/CUDA text worker](recipe/decider/README.md); Choice/Noul/isolated Score; opt-in batching, Graph and request-local prefixes; [validation scope](recipe/decider/validation.md) | | Open-Jev-27B-v1.1 | [Native Rust/CUDA worker](recipe/open_jev/native.md); eager independent text candidates; [H200 validation](recipe/open_jev/validation.md) | | Open-Jev-9B | The same [native Rust/CUDA worker](recipe/open_jev/native.md); eager independent text candidates; [reference comparison on sm_89](recipe/open_jev/validation-9b.md) | | CLM-v0.1-8B | [External worker with a CPU stub encoder](recipe/clm/README.md); contract checks only, real Qwen3-8B decisions unverified by this recipe | diff --git a/docs/supported-models.md b/docs/supported-models.md index 9aba8d39..ba2b6286 100644 --- a/docs/supported-models.md +++ b/docs/supported-models.md @@ -16,6 +16,7 @@ Models that are being added are also tracked in issues labeled [new model](https | Cua-S1 4B 0.2, `multimodal` adapter | Reference worker on Transformers and PEFT, [`src/frontend/cua_s1.py`](../src/frontend/cua_s1.py); no recipe yet | Not supported | Validated ([#17](https://github.com/ThinkFlowLab/system1-omni/pull/17), [#18](https://github.com/ThinkFlowLab/system1-omni/pull/18)) | Not supported | The state is one PNG or JPEG image; upstream's `weights.lock.json` next to the base weights | | Open-Jev-27B-v1.1 | [Native Rust/CUDA worker](../recipe/open_jev/native.md) on the shared Qwen3.5/3.8 executor | Not supported | Validated on H200 (sm_90) for the [74 single-candidate workload](../recipe/open_jev/validation.md) | Not supported | Compute capability 8.0 or newer, CUDA toolkit to build, exported merged weights and trained head | | Open-Jev-9B | The same [native Rust/CUDA worker](../recipe/open_jev/native.md) | Not supported | Validated on compute capability 8.9 against the reference for [253 requests](../recipe/open_jev/validation-9b.md) | Not supported | Compute capability 8.0 or newer, CUDA toolkit to build, exported merged weights and trained head | +| Decider-2B v11 | [Native Rust/CUDA worker](../recipe/decider/README.md) | Contract checks only; no CPU inference | RTX 4090 (`sm_89`) synthetic parity corpus; [limits and provenance](../recipe/decider/validation.md) | Not supported | Pinned original BF16 checkpoint; ABI 7 library and consumers rebuilt together | | CLM-v0.1-8B | [External `clm-serve` recipe](../recipe/clm/README.md) with a CPU stub embeddings server | **Stub-encoder contract checks only** ([#23](https://github.com/ThinkFlowLab/system1-omni/pull/23)); not real Qwen3-8B decisions | Real encoder unverified by the merged recipe | Unverified | Python, upstream CLM and head checkpoint; a real encoder requires a separate embeddings server | - **Validated:** covered by the recipe on `main` or by the checks in the linked merged pull request. diff --git a/mkdocs.yml b/mkdocs.yml index e276a583..417ba202 100644 --- a/mkdocs.yml +++ b/mkdocs.yml @@ -68,5 +68,7 @@ nav: - Open-Jev native text worker: recipe/open_jev/native.md - Open-Jev-27B H200 validation: recipe/open_jev/validation.md - Open-Jev-9B validation: recipe/open_jev/validation-9b.md + - Decider native text worker: recipe/decider/README.md + - Decider validation: recipe/decider/validation.md - CLM stub-encoder contract: recipe/clm/README.md - Contributing: CONTRIBUTING.md diff --git a/recipe/README.md b/recipe/README.md index 8909eb94..a01d54cb 100644 --- a/recipe/README.md +++ b/recipe/README.md @@ -17,6 +17,8 @@ For a first real decision, follow the [complete CPU walkthrough](../docs/getting backbone and trained decision head of Open-Jev-27B-v1.1 or Open-Jev-9B, then serve with Rust and CUDA. [Open-Jev-9B validation](open_jev/validation-9b.md) compares the 9B worker with the reference. +- [Decider-2B v11 native worker](decider/README.md): pinned BF16 text decisions, + optional request-local batching, Graph and prefix reuse; [validation](decider/validation.md). - [CLM behind the frontend](clm/README.md): run CLM's own server behind the frontend on CPU with a stub encoder, and what the response comparison has to allow for. diff --git a/recipe/decider/README.md b/recipe/decider/README.md new file mode 100644 index 00000000..0e0c56cc --- /dev/null +++ b/recipe/decider/README.md @@ -0,0 +1,207 @@ +# Native Decider-2B v11 text decisions + +Run these commands from the repository root on Linux with Rust, an NVIDIA GPU, +CUDA toolkit/nvcc and cuBLASLt. Rebuild the library and every native Qwen consumer together for ABI **7**. +The existing Qwen backend targets compute capability +8.0 or newer; hardware validation is described in [validation.md](validation.md). +The worker itself uses no Python/PyTorch serving process. Python 3.10+ is used only +for downloading and optional reference validation. Keep at least 8 GB disk free for +weights and Rust build products; reference Python packages need additional space. +Device memory includes about 3.76 GB of BF16 checkpoint data plus model/head buffers +and length-dependent activation/workspace storage. + +## Download and build + +```sh +python3 recipe/decider/download_weights.py /path/to/models/decider-2b-v11 +cargo build --release --locked -p omni-decider-native -p omni-jev +NVCC=/usr/local/cuda/bin/nvcc \ + bash src/backends/cuda/qwen3_5/build.sh target/release 89 +``` + +Replace `89` with the target GPU's compute capability. The download helper pins +`Mapika/decider-2b` revision `533964dae8be954c5b5e19fa4948e48408094c1e` and checks +all five downloaded files; it never downloads unpinned remote Python code. Its +`--endpoint` option can select an accessible Hugging Face mirror; fixed hashes +remain required. Interrupted downloads leave no usable unverified final file. + +Startup independently hashes the required checkpoint/config/calibration/tokenizer +and rejects changed artifacts before loading CUDA. Use the original single +`model.safetensors`, with no `model.safetensors.index.json`; no export or LoRA merge +is needed. Keep files immutable while loading and serving. + +## Start the worker and frontend + +Worker terminal: + +```sh +DECIDER_GRAPH=0 \ +DECIDER_MODEL=/path/to/models/decider-2b-v11 \ +DECIDER_CUDA_LIB="$PWD/target/release/libqwen3_5_cuda.so" \ +DECIDER_HOST=127.0.0.1 DECIDER_PORT=8000 \ + ./target/release/omni-decider +``` + +Decider defaults to eager. `DECIDER_GRAPH=1` opts the backbone into CUDA Graph +capture/replay; only 0/1 are accepted and an unset switch is off. `CUA_S1_GRAPH` +has no effect on Decider; other workers retain their existing switch. +`DECIDER_CUDA_LIB` defaults to the library beside the binary. The worker uses CUDA +device 0; use `CUDA_VISIBLE_DEVICES` to select a physical device. + +The socket binds only after artifact validation, CUDA loading and a real warmup +request. Startup failure does not advertise readiness. In another terminal: + +```sh +curl --fail http://127.0.0.1:8000/health +OMNI_JEV_BIND=127.0.0.1:8080 \ +OMNI_JEV_BACKEND_URL=http://127.0.0.1:8000 \ + ./target/release/omni-jev +``` + +Health returns the model identity, checkpoint/reference revisions, BF16/effective execution mode +and effective per-type temperatures. A retired executor returns HTTP 503 and +`status: unavailable`. The frontend exposes its existing health behavior; see the +[frontend configuration](../../src/frontend/README.md) for its own limits/timeouts. + +## First decision + +```sh +curl --fail http://127.0.0.1:8080/v1/systemone \ + -H 'Content-Type: application/json' \ + --data-binary @recipe/decider/example-request.json +``` + +The example asks a Choice, Noul and three-level Score question about one text +state. The response has model `decider-2b-v11`, ordered typed answers and usage. +Question identities are omitted from model text. The exact output depends on the +checkpoint computation; this example is not a task-quality guarantee. + +Errors use JSON `detail`: HTTP 415 for unsupported Content-Type, 413 for a body +above 8 MiB, 422 for invalid/unsupported requests, and 503 when the model becomes +unavailable after an execution failure. Request-preparation task failures return +500. Invalid inputs do not retire the executor. Empty questions return HTTP 200 +with empty answers and zero input/output usage. No partial answers are returned. + +Supported inputs are plain text/JSON state, independent Choice (2–255 options), +Noul and isolated Score (2–10 levels). State-prefix truncation is 32,768 tokens; +complete rows can include question suffixes up to 36,864 tokens. At most 1,024 +expanded rows and 1,048,576 processed tokens are accepted per request. These bounds +do not bound all waiting requests together; runtime pending-queue limits and +cross-request batching remain separate work. + +The pinned worker requires released calibration. Images, video, chat/schema-first, +packed questions, cross-request prefix caching and quantized/CPU/Metal execution +are unsupported. Request-local prefix reuse is opt-in through `DECIDER_PREFIX=1` +or the conservative `DECIDER_PREFIX=auto` policy; both require Graph off. +CUDA Graph replay is opt-in through `DECIDER_GRAPH=1`, as described below. See the [model contract](../../src/models/decider/README.md) +for native JSON restrictions and the processing/execution boundary. + +## Validation and diagnostics + +```sh +cargo fmt --all --check +cargo clippy --workspace --locked --all-targets -- -D warnings +cargo test --workspace --locked +cargo build --workspace --release --locked +DECIDER_MODEL=/path/to/models/decider-2b-v11 \ + cargo test --release --locked -p omni-decider-native --test contract -- --ignored +DECIDER_MODEL=/path/to/models/decider-2b-v11 \ +DECIDER_CUDA_LIB="$PWD/target/release/libqwen3_5_cuda.so" DECIDER_GRAPH=0 \ + cargo test --release --locked -p omni-decider-native --test gpu -- --ignored +DECIDER_CUDA_LIB="$PWD/target/release/libqwen3_5_cuda.so" \ + cargo test --release --locked -p omni-decider-native --lib \ + bf16_head_rounds_before_fp32_calibration -- --ignored +``` + +The CPU contract target requires only tokenizer/config; the GPU target requires the +complete pinned checkpoint. The projection test needs a built backend and a GPU. + +For exact prepared rows and raw BF16-rounded logits, the diagnostic binary reads +one JSON request per line and emits rows, logits and complete response: + +```sh +cargo build --release --locked -p omni-decider-native --example decider-run +python3 -c 'import json; print(json.dumps(json.load(open("recipe/decider/example-request.json"))))' > /tmp/decider-request.jsonl +DECIDER_GRAPH=0 ./target/release/examples/decider-run \ + /path/to/models/decider-2b-v11 "$PWD/target/release/libqwen3_5_cuda.so" \ + < /tmp/decider-request.jsonl +``` + +See [validation.md](validation.md) for full-checkpoint reference comparison, +HTTP/frontend checks, prerequisites and the recorded scope. + +## Optional request-local packing + +Set `DECIDER_BATCH_MAX_ROWS=4 DECIDER_BATCH_MAX_TOKENS=4096` on `omni-decider` +to pack independent rows within one admitted request. Rows default to 1; accepted +limits are 1–4 rows and 1–4096 packed tokens. Longer complete rows execute alone. +Malformed or out-of-range settings fail startup. These limits do not change HTTP +admission, token usage, or the maximum complete-row length. + +The selected-label head uses persistent buffers sized to the row limit and projects +each batch in one BF16 GEMM. No cross-request batching is introduced. Keep the +single-row default until numerical and performance results fit your workload. + +## Optional CUDA Graph replay + +`DECIDER_GRAPH=1` works with the request-local batch limits above. The backbone +caches at most 64 captures by the **ordered sequence-length vector**, uploads +current token IDs before replay, and clears captures after synchronizing before +scratch growth. A shape miss runs eagerly once and records the warmed layer loop; +it does not launch that capture on the already-computed residual. Label projection +and response assembly remain outside the graph. + +Health's `graph` object and the diagnostic runner's outer `graph` record report +requested/enabled mode, captures, replays, fallbacks, invalidations and cached shapes. +Decision-response fields remain unchanged. Capture failure disables Graph for the +worker lifetime and keeps the completed eager result; health then reports `eager`. +Launch/inference failures still retire the worker. Capture and startup costs must +be measured separately from warm replay. + +## Optional request-local shared prefixes (pending integration) + +This integration depends on the pending Qwen continuation/fixed-GEMM/executor +stack (#97/#98/#99). Build this complete branch with ABI7; an ABI5 library is +rejected. `DECIDER_PREFIX=1 DECIDER_GRAPH=0` plans exact shared token prefixes +within the admitted request and restores attention KV, GDN recurrent state and +three convolution-history rows for every branch. Contiguous rows for one question +can also share a longer question prefix. The executor ends reusable spans on +64-token boundaries and repeats their remainder in each suffix. Every branch +retains at least one token and original row order. Buffers persist but their +contents are valid only within one call; there is no cross-request cache. + +Prefix execution uses the dependency's fixed GEMM algorithms. Set `DECIDER_FIXED=1` +with prefix/Graph off for the independent full-row control. With no usable prefix +(less than one reusable 64-token chunk), prefix mode executes independent fixed +rows. Both modes preserve the batch adapter's head grouping; backbone branches +run eagerly. Combining prefix/fixed with Graph or enabling prefix and fixed +together rejects startup. All new switches default off. Health reports effective +`shared`/`fixed` mode and cumulative `prefix.shared_requests`/`saved_tokens`. +The diagnostic outer record includes the same counters; decision output and usage +do not count execution savings differently. + +## Conservative automatic prefix selection + +`DECIDER_PREFIX=auto` is a default-off alternative to forced `DECIDER_PREFIX=1`. +The worker uses shared fixed-GEMM execution only when the existing aligned prefix +plan saves at least 4096 tokens and at least one third of the original row-token +work. Other requests use the normal independent/packed eager path, including +single-row and short/no-sharing requests. This avoids forcing short work through +the slower fixed-GEMM path. These are conservative host-policy thresholds, not a +guarantee of faster execution on every GPU or workload. + +`DECIDER_PREFIX=1` and `DECIDER_FIXED=1` retain their existing forced semantics. +Auto, forced prefix and fixed modes reject `DECIDER_GRAPH=1`; all switches default +off. Auto retains request-local state only. Input/output contracts and token usage +are unchanged, but normal and fixed GEMM algorithms can differ numerically: check +the original reference gates for both selected paths. Health reports `prefix_mode: +auto`, total shared requests/saved tokens and `auto_independent_requests`; counts +include the readiness warmup. + +```sh +DECIDER_PREFIX=auto DECIDER_GRAPH=0 \ +DECIDER_BATCH_MAX_ROWS=4 DECIDER_BATCH_MAX_TOKENS=4096 \ +DECIDER_MODEL=/path/to/models/decider-2b-v11 \ +DECIDER_CUDA_LIB="$PWD/target/release/libqwen3_5_cuda.so" \ + target/release/omni-decider +``` 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/recipe/decider/example-request.json b/recipe/decider/example-request.json new file mode 100644 index 00000000..fbec40a6 --- /dev/null +++ b/recipe/decider/example-request.json @@ -0,0 +1 @@ +{"state":"A customer received a damaged item and requests a refund.","questions":{"route":{"type":"choice","instructions":"Choose the responsible team.","criteria":{"billing":"Refund and billing","shipping":"Delivery status","sales":"New purchases"}},"refund":{"type":"noul","instructions":"The customer requests a refund."},"severity":{"type":"score","instructions":"Severity of the damage","criteria":["Minor","Moderate","Severe"]}}} diff --git a/recipe/decider/validation.md b/recipe/decider/validation.md new file mode 100644 index 00000000..5340c397 --- /dev/null +++ b/recipe/decider/validation.md @@ -0,0 +1,260 @@ +# Decider-2B v11 validation on RTX 4090 + +Validation date: 2026-10-07. This records one BF16 CUDA device and a frozen synthetic +contract corpus. It does not establish task quality, every supported GPU, performance +improvement or universal bit-identical agreement with Transformers. + +## Controls and fidelity gates + +Checkpoint/tokenizer/calibration: Mapika/decider-2b +`533964dae8be954c5b5e19fa4948e48408094c1e`. Checkpoint SHA-256: +`acaef2228b134dcdc20cad4ee79219482c927ec819aa3687b9b8a575c338817f`. +Native checkpoint/config/tokenizer checks are in `src/models/decider/native/src/checkpoint.rs`. +Reference runtime: Mapika/decider `50d0be0d7cb43d2066965ce5fa7f3fe4e489a60f`, +decider-ai 1.8.1. The comparison harness verifies fixed hashes of all five imported +reference files before running either implementation and saves them in its protocol. + +RTX 4090, sm_89, 24 GiB; driver 595.71.05, CUDA toolkit 13.0.88, Rust 1.98.1, +Python 3.12, Torch 2.14.0+cu130 and Transformers 5.17.0. The original eager baseline used one complete unpadded row at a time and ABI5. +This dependent prefix branch rebuilds the backend and consumers for ABI7; followup +controls below cover unchanged eager/Graph behavior and fixed/shared execution. The reference uses `DecisionModel.slot_logits` +with unpadded independent rows and disabled cache. Its missing causal-conv1d and +flash-linear-attention packages use Transformers' PyTorch reference fallbacks. +No optimized-reference speed comparison is claimed. + +Before execution, the harness fixes these gates: + +- Maximum absolute answer-probability and per-level fit drift: **0.02**. +- Maximum Score expectation drift: **0.1**. +- Identical Choice decisions when reference top-two probability margin is at least + **0.05**; report all lower-margin disagreements. +- Exact prepared tokens, candidate IDs/order, readout positions, usage and answer + identities. Complete responses retain all model-specific fields. + +## Recorded eager baseline results + +The 18-request corpus has **41 independent inference rows**. It covers Choice with +2/3/10/11/26/255 options, Score with 2/3/10 isolated levels, true/false and +criteria-only Noul, structured Unicode and long arrays, mixed questions, reversed +question order, repeated input, empty questions and long context. The longest +reference-compared row is 4,140 tokens. + +Both the initial run and the source-hash-verified rerun pass the same gates: +maximum probability/per-level-fit drift **0.0119**, maximum Score drift **0.01**, +and **zero Choice disagreements**, including lower-margin cases. Raw logit drift +is reported separately (maximum 0.125); logits are not bit-identical. These are +observed results for this corpus, not a general numerical-error bound. + +CPU opt-in tests check 12 pinned token/row/usage fixtures, every label through JT, +12 response assemblies including zero-fit Score, and state-prefix truncation. +Normal CPU tests cover selected-row ordering/padding, invalid shapes/IDs/nonfinite +weights, file content/size hashes, alternate index rejection and complete-row +validation before dispatch. A real GPU microtest checks the BF16 projection tie +`1 + 0.00390625 -> 1`, before FP32 calibration, and excludes the padded label. + +Real model regression checks verify mixed answer order, exact repeated output, +validation-error recovery and empty requests without inference. The separate native GPU boundary test executes a **36,864-token complete row** +twice, rejects the next token before execution, and then reproduces an earlier short +request exactly after scratch growth. It does not expand the reference-parity claim +beyond 4,140 tokens. + +HTTP testing passes **27 checks**: 18 direct/native diagnostic/frontend comparisons, +six invalid-input cases, the worker body limit, concurrent requests and recovery +health. The concurrent check makes 12 requests across direct and frontend paths +with four client threads; all returned bytes match the serial baseline. The first +HTTP harness run stopped on a Python variable-shadowing error after its earlier +checks; the corrected full run passes. The worker binaries were unchanged. + +The test-only CUDA wrapper lets the genuine warmup projection complete, then fails +only the second `1 x 256 x 2048` head GEMM. The worker returns 503, health becomes +unavailable, two later requests are refused, and exactly two projections are seen. +Sampled per-process GPU residency drops from **4,192 MiB to 394 MiB**, demonstrating +retirement of the model allocation; these are two lifecycle samples, not peak-memory +or performance measurements. The wrapper delegates all other operations to the +real library and is never used by normal serving. + +Cancellation semantics inherit the tested shared SerialScheduler. Decider-specific +caller cancellation and injected Rust panic are not separately hardware-tested. +CPU/Metal inference, other GPUs/precisions, dynamic batching, prefix reuse, Decider +Graph replay, model-quality datasets, throughput and cold/warm latency comparisons +are outside the original eager campaign; the followup checks below cover batching, +Graph and request-local prefix reuse. + +## Reproduce + +Build and prepare the checkpoint using [the native recipe](README.md). For the +reference environment, use Python 3.12 and the recorded Torch/Transformers versions; +install Transformers 5.17.0, tokenizers 0.22.2 and Torch 2.14.0 built for CUDA 13.0. +The native worker's preparation and serving require none of these Python packages. + +Obtain the pinned reference in a separate directory: + +```sh +git clone https://github.com/Mapika/decider.git /path/to/decider-reference +git -C /path/to/decider-reference checkout 50d0be0d7cb43d2066965ce5fa7f3fe4e489a60f +``` + +From the System1-Omni repository root, with a reserved/available NVIDIA GPU: + +```sh +cargo build --release --locked -p omni-decider-native --example decider-run +DECIDER_GRAPH=0 python tests/decider/verify_reference.py \ + --model /path/to/models/decider-2b-v11 \ + --reference /path/to/decider-reference \ + --binary "$PWD/target/release/examples/decider-run" \ + --library "$PWD/target/release/libqwen3_5_cuda.so" \ + --output /path/to/evidence/parity + +python tests/decider/verify_http.py \ + --worker "$PWD/target/release/omni-decider" \ + --frontend "$PWD/target/release/omni-jev" \ + --model /path/to/models/decider-2b-v11 \ + --library "$PWD/target/release/libqwen3_5_cuda.so" \ + --parity-output /path/to/evidence/parity \ + --output /path/to/evidence/http +``` + +The HTTP runner starts local sockets on 18110/18111, checks readiness and shuts its +processes down even on failure. It retains worker/frontend logs and complete results. +Do not run it over unrelated processes already using these ports. + +For the isolated failure test, build a wrapper beside the real library so its +runtime search path resolves the actual backend: + +```sh +gcc -shared -fPIC -Wall -Wextra -I src/backends/cuda/qwen3_5 \ + tests/decider/fail_head.c -L target/release -Wl,--no-as-needed \ + -lqwen3_5_cuda -Wl,-rpath,'$ORIGIN' -ldl \ + -o target/release/decider-fail-head.so +python tests/decider/verify_retirement.py \ + --worker "$PWD/target/release/omni-decider" \ + --model /path/to/models/decider-2b-v11 \ + --wrapper "$PWD/target/release/decider-fail-head.so" \ + --output /path/to/evidence/retirement +``` + +This runner uses port 18112 and `nvidia-smi` per-process memory reporting. The recipe +also lists the explicit CPU and real-GPU Cargo tests. Ordinary workspace tests +compile/discover their ignored cases but do not establish hardware validation. + +## Evidence + +The task evidence archive retains both parity runs, the original HTTP harness +failure and its successful rerun, native/reference rows/logits/responses, protocol, +source/binary/library hashes, environment inventory, GPU tests and retirement logs. +Generated run output and binaries are kept out of the source tree; only maintained +fixtures and reproduction helpers are committed. The implementation PR links the +immutable evidence archive and its SHA-256. All requests are synthetic and carry +no private/customer content. A video is inapplicable to this text API integration. + +## Request-local batching comparison + +The followup runner `tests/decider/verify_modes.py` accepts a frozen JSON plan with +`model`, `library`, `parity_output` (the pinned reference run above), +`repetitions: 2`, `timed_cases` and `modes`. Each mode supplies `name`, an absolute +`binary` path and an `env` mapping. It hashes binaries, library, reference records +and comparison scripts before execution; retains complete records; checks prepared +rows exactly and the same response gates. Time is the diagnostic JSON-lines +roundtrip at concurrency 1, including processing, synchronization, serialization +and IPC. Startup to the empty diagnostic response and one workload warmup are +reported/excluded separately. These measurements do not establish HTTP latency, +peak memory or production throughput. Device memory is a whole-device sample at +the end of each mode, not a peak measurement. + +On the RTX 4090 followup campaign, row limits 2 and 4 passed all 18 reference +requests. Row limit 8 failed the unchanged probability gate on structured Unicode +input (0.0208 > 0.02), so supported packing is capped at **4 rows**. The failed +campaign and narrowed-protocol rerun are retained in the PR evidence archive. +Defaults remain one row; do not extrapolate these small workloads to arbitrary +requests or hardware. Batch fault injection also verifies a two-row head failure +retires the complete worker after real startup warmup. + +## Graph checks + +`DECIDER_GRAPH=1` enables explicit backbone capture/replay. The ignored root +`tests/qwen3_5/graph.rs` test uses the pinned Decider checkpoint and compares eager +and graph hidden states exactly for ordered length vectors, token changes under +the same shape, scratch growth, replay and FIFO eviction beyond 64 shapes. Run: + +```sh +DECIDER_MODEL=/path/to/models/decider-2b-v11 \ +DECIDER_CUDA_LIB="$PWD/target/release/libqwen3_5_cuda.so" CUA_S1_GRAPH=1 \ + cargo test --release --locked -p omni-qwen3-5-native --test graph -- --ignored --test-threads=1 +``` + +Set `require_replays: true` on each Graph mode in the comparison plan to require +actual captures and warm replays with zero fallbacks. The diagnostic runner preserves +these counters separately from the decision response. `tests/decider/fail_graph.c` +is a test-only wrapper that rejects capture startup while real CUDA eager inference +continues; `verify_graph_fallback.py` checks all pinned response cases and proves +one fallback, zero captures/replays, and effective eager execution. Build the wrapper +like `fail_head.c` above, then pass `--binary`, `--model`, `--wrapper`, +`--parity-output` and `--output` to that runner. + +## Shared-prefix controls + +The dependent prefix branch consumes the continuation, fixed-GEMM and executor +commits from Qwen PRs #97/#98/#99 while preserving main's packed prefill and Graph. +Rebuild all consumers with ABI7. The unchanged default GEMM handle serves the +normal/Graph path; fixed/shared modes use a separate no-split-K handle selected +at reference M64. Model tests compare shared branches bit for bit with independent +`forward_fixed` rows, including prefixes around 64-token boundaries, growth, +repetition, reordering, invalid IDs and interleaved Graph calls. Kernel tests now +include Decider's 2048 hidden width, 16 GDN heads and 8/2 attention heads alongside +the existing 4B/9B/27B shapes. Run ignored tests serially with `CUA_S1_CUDA_LIB`, +`QWEN3_5_MODEL`, `DECIDER_MODEL` and `DECIDER_CUDA_LIB` set to the ABI7 library +and pinned Decider checkpoint. + +For response controls, put fixed4 before shared4 in the frozen modes plan and set +`exact_control: "fixed4"` and `require_shared: true` on shared4. Both use +`DECIDER_BATCH_MAX_ROWS=4` and Graph off; fixed4 sets `DECIDER_FIXED=1`, shared4 +sets `DECIDER_PREFIX=1`. This verifies exact raw logits and complete responses +between fixed/shared as well as the original Transformers fidelity gates, and +requires actual request-local token savings. Keep normal eager4 in the plan to +measure algorithm-selection cost separately from prefix reuse. + +The cached stock Qwen3.5-4B artifact has FP32 A_log and is not a valid merged BF16 +worker export; its attempted legacy-loader check fails before inference and is +retained as unavailable coverage. The same constructor/explicit-mode regression +passes on Decider2B with the inherited CUA switch. Full Cua-S1/Open-Jev decisions +are not revalidated by the Decider campaign; they require their proper merged +exports and corresponding hardware. + +The maintained `tests/decider/data/prefix-workloads.json` adds two targeted long +shared-context requests (mixed types and ten Score levels). Generate their pinned +reference records with `verify_reference.py --cases tests/decider/data/prefix-workloads.json` +and the same model/reference/binary/library/output arguments above. Then point +`verify_modes.py` at that parity directory, time exactly those two names, and keep +two repetitions per mode. This corpus is independent of the original 18-case run; +it resolves long-prefix reuse performance and does not replace broader fidelity +checks. The fixed/shared exact-logit gate also applies to measured repetitions. + +## Followup measurements (diagnostic roundtrip) + +The final bounded-row campaign measured `mixed` at 25.69 ms for single rows and +15.75 ms for 4-row packing; `score_10` at 51.03→29.44 ms. Independent Graph4 measured +`mixed` at 15.80→15.12 ms relative to matched eager4, while long-context single-row +latency was 89.76→89.73 ms. Each value is the median of two warm samples on the +same RTX 4090. These are narrow synthetic diagnostic timings, not HTTP latency +or a general speedup. Startup to the empty diagnostic record was measured +separately (approximately 5.5–6.4 seconds), with actual captures/replays recorded. + +The ABI 7 shared-prefix comparison retains the fixed-algorithm control separately: + +| Workload | Normal eager4 (ms) | Independent fixed4 (ms) | Shared4 (ms) | +| --- | ---: | ---: | ---: | +| Original mixed, short state |15.73|31.27|31.51| +| Original structured Unicode |32.43|47.57|41.44| +| Targeted long mixed, 5 rows / 17683 processed tokens |414.47|484.32|144.18| +| Targeted long Score10, 10 rows / 24460 processed tokens |576.17|672.89|151.96| + +Every cell uses two warmed repetitions with raw min/max retained in the archive. +The long shared runs are bit-identical in raw candidate logits and complete +responses to their independent fixed controls, and satisfy the original reference +fidelity gates. The two long cases show 2.87× and 3.79× diagnostic-roundtrip speedup +against matched normal eager4; they do not establish task quality or a workload-wide +improvement. Short states can regress due to fixed GEMM selection and branch +launch overhead, so all optimizations remain opt-in. Shared and Graph are tested +as separate modes, not combined. End-of-mode whole-device memory samples for the +long corpus were 4626 MiB eager4, 4658 MiB fixed4 and 4746 MiB shared4; peak memory +and sustained throughput were not measured. diff --git a/src/backends/cuda/qwen3_5/README.md b/src/backends/cuda/qwen3_5/README.md index e142176e..ddd033a5 100644 --- a/src/backends/cuda/qwen3_5/README.md +++ b/src/backends/cuda/qwen3_5/README.md @@ -11,12 +11,23 @@ The norm, elementwise and q/k preparation kernels round to bfloat16 where Transf `cs1_attention_gated` fuses the sigmoid gate into the attention epilogue, preserving the BF16 rounding of both attention and sigmoid before multiplication. The native workers use this entry point; the separate operations remain available for kernel -comparisons. Rebuild the library and workers together for ABI version 5, which -includes the shared vision and CUDA Graph entry points alongside gated attention. +comparisons. Rebuild the library and workers together for ABI version 7, which +includes shared vision and CUDA Graph entry points, request-local continuation +state/KV operations and fixed-algorithm GEMM selection alongside gated attention. Gated DeltaNet preparation stores converted TF32 operands in three-byte component planes, preserves the original four-term TF32 accumulation, and writes U/W fragments directly as bfloat16. Dynamic shared memory is 72 KiB per block. The [H200 comparison](../../../../benchmarks/gdn/README.md) records complete GDN call latency, numerical checks, and the small end-to-end change measured with the Open-Jev worker from PR #55. The shared Rust model can pack independent sequences for input and gate/up GEMMs. Output/down GEMMs retain each prompt's original shape and reduction order; -attention, convolution and GDN calls remain sequence-local. The CUDA ABI is -unchanged. Open-Jev uses this path within requests; Cua-S1 keeps single-prompt calls. +attention, convolution and GDN calls remain sequence-local. Packed prefill preserves per-sequence state. Open-Jev uses this path within +requests; Cua-S1 keeps single-prompt calls. Decider optionally packs complete +question rows within one admitted request. + +ABI7 adds `cs1_copy_rows`, `cs1_gdn_conv_history`, `cs1_gdn_prefill_state`, +`cs1_attention_gated_cached` and `cs1_gemm_create_fixed` for request-local prefix +continuations. Shared spans end on 64-token GDN chunk boundaries. Fixed GEMM +selection fails explicitly if a requested shape cannot use the selected algorithm; +M-independent output equality is a device-tested requirement, not a portable +cuBLAS guarantee. See the [Decider recipe](../../../../recipe/decider/README.md). +Rebuild `libqwen3_5_cuda.so` and restart every Qwen consumer together; the Rust +loader rejects an older ABI at startup with a rebuild hint. diff --git a/src/backends/cuda/qwen3_5/attention.cu b/src/backends/cuda/qwen3_5/attention.cu index 20ffceed..162943db 100644 --- a/src/backends/cuda/qwen3_5/attention.cu +++ b/src/backends/cuda/qwen3_5/attention.cu @@ -5,7 +5,9 @@ // four warps of 16 rows each, and walks the keys up to its last query in tiles of // 32, keeping the output and the online softmax in registers. The probabilities are // rounded to bfloat16 for the P*V product, as in flash attention; the running sums -// stay float32. +// stay float32. The queries can be the last Tq of Tk positions (cached keys before +// them); key tiles always start at position 0, so a query sees the same tiles in the +// same order wherever the queries start, and tiles past its own position change nothing. #include "common.cuh" #include "mma.cuh" #include "ops.h" @@ -81,7 +83,7 @@ constexpr int SMEM_BYTES = (BM + 2 * BN) * LDS * 2; template __global__ void __launch_bounds__(THREADS) flash_kernel(const bf16* __restrict__ q, const bf16* __restrict__ k, const bf16* __restrict__ v, int ldv, - const bf16* __restrict__ gate, bf16* __restrict__ out, int T, int Hq, int Hk, + const bf16* __restrict__ gate, bf16* __restrict__ out, int Tq, int Tk, int Hq, int Hk, float scale_log2) { extern __shared__ __align__(16) unsigned char smem[]; bf16* qs = reinterpret_cast(smem); @@ -92,10 +94,11 @@ __global__ void __launch_bounds__(THREADS) const int tid = threadIdx.x, warp = tid / 32, lane = tid % 32; const int g = lane / 4, t = lane % 4; const int row0 = q0 + warp * 16; // this warp's first query + const int offset = Tk - Tq; // the position of query 0 for (int c = tid; c < BM * (D / 8); c += THREADS) { const int r = c / (D / 8), col = (c % (D / 8)) * 8, row = q0 + r; - cp_async16(qs + r * LDS + col, q + ((size_t)min(row, T - 1) * Hq + h) * D + col, row < T); + cp_async16(qs + r * LDS + col, q + ((size_t)min(row, Tq - 1) * Hq + h) * D + col, row < Tq); } cp_async_commit(); @@ -104,23 +107,24 @@ __global__ void __launch_bounds__(THREADS) for (int n = 0; n < D / 8; n++) o[n][0] = o[n][1] = o[n][2] = o[n][3] = 0.f; float m[2] = {-INFINITY, -INFINITY}, l[2] = {0.f, 0.f}; - const int kv_end = min(T, q0 + BM); + const int kv_end = min(Tk, offset + q0 + BM); for (int k0 = 0; k0 < kv_end; k0 += BN) { for (int c = tid; c < BN * (D / 8); c += THREADS) { const int r = c / (D / 8), col = (c % (D / 8)) * 8, s = k0 + r; - cp_async16(ks + r * LDS + col, k + ((size_t)min(s, T - 1) * Hk + hk) * D + col, s < T); + cp_async16(ks + r * LDS + col, k + ((size_t)min(s, Tk - 1) * Hk + hk) * D + col, s < Tk); } cp_async_commit(); for (int c = tid; c < BN * (D / 8); c += THREADS) { const int r = c / (D / 8), col = (c % (D / 8)) * 8, s = k0 + r; - cp_async16(vs + r * LDS + col, v + (size_t)min(s, T - 1) * ldv + (size_t)hk * D + col, s < T); + cp_async16(vs + r * LDS + col, v + (size_t)min(s, Tk - 1) * ldv + (size_t)hk * D + col, s < Tk); } cp_async_commit(); cp_async_wait<1>(); // Q and K __syncthreads(); - // keys past every query of this warp contribute nothing - const bool active = k0 <= row0 + 15; + // keys past every query of this warp contribute nothing, and rows past the + // last query are never stored + const bool active = row0 < Tq && k0 <= offset + row0 + 15; float sc[BN / 8][4]; #pragma unroll for (int n = 0; n < BN / 8; n++) sc[n][0] = sc[n][1] = sc[n][2] = sc[n][3] = 0.f; @@ -146,8 +150,8 @@ __global__ void __launch_bounds__(THREADS) for (int n = 0; n < BN / 8; n++) { #pragma unroll for (int e = 0; e < 4; e++) { - const int key = k0 + n * 8 + 2 * t + (e & 1), row = row0 + g + (e >> 1) * 8; - sc[n][e] = (key <= row && key < T) ? sc[n][e] * scale_log2 : -INFINITY; + const int key = k0 + n * 8 + 2 * t + (e & 1), pos = offset + row0 + g + (e >> 1) * 8; + sc[n][e] = (key <= pos && key < Tk) ? sc[n][e] * scale_log2 : -INFINITY; mx[e >> 1] = fmaxf(mx[e >> 1], sc[n][e]); } } @@ -213,7 +217,7 @@ __global__ void __launch_bounds__(THREADS) #pragma unroll for (int r = 0; r < 2; r++) { const int row = row0 + g + r * 8; - if (row >= T) continue; + if (row >= Tq) continue; bf16* dst = out + ((size_t)row * Hq + h) * D + 2 * t; #pragma unroll for (int n = 0; n < D / 8; n++) { @@ -232,20 +236,20 @@ __global__ void __launch_bounds__(THREADS) template int launch(const void* q, const void* k, const void* v, int ldv, const void* gate, void* out, - int T, int Hq, int Hk, int Dh, float scale, void* stream) { - if (Dh != D || Hk <= 0 || Hq <= 0 || Hq % Hk != 0 || ldv % 8 != 0 || ldv < Hk * Dh || T < 0) + int Tq, int Tk, int Hq, int Hk, int Dh, float scale, void* stream) { + if (Dh != D || Hk <= 0 || Hq <= 0 || Hq % Hk != 0 || ldv % 8 != 0 || ldv < Hk * Dh || Tq < 0 || Tk < Tq) return cudaErrorInvalidValue; - if (T == 0) return cudaSuccess; + if (Tq == 0) return cudaSuccess; if (Gated && gate == nullptr) return cudaErrorInvalidValue; // Once per specialization (for the device current at the first call). static const cudaError_t configured = cudaFuncSetAttribute( flash_kernel, cudaFuncAttributeMaxDynamicSharedMemorySize, SMEM_BYTES); if (configured != cudaSuccess) return configured; constexpr float LOG2E = 1.4426950408889634f; - flash_kernel<<<<(stream)>>>( static_cast(q), static_cast(k), static_cast(v), ldv, - static_cast(gate), static_cast(out), T, Hq, Hk, scale * LOG2E); + static_cast(gate), static_cast(out), Tq, Tk, Hq, Hk, scale * LOG2E); return cudaGetLastError(); } @@ -274,10 +278,16 @@ extern "C" int cs1_attn_prep(const void* qg, const void* kr, int ld, const void* extern "C" int cs1_attention(const void* q, const void* k, const void* v, int ldv, void* out, int T, int Hq, int Hk, int Dh, float scale, void* stream) { - return flash::launch(q, k, v, ldv, nullptr, out, T, Hq, Hk, Dh, scale, stream); + return flash::launch(q, k, v, ldv, nullptr, out, T, T, Hq, Hk, Dh, scale, stream); } extern "C" int cs1_attention_gated(const void* q, const void* k, const void* v, int ldv, const void* gate, void* out, int T, int Hq, int Hk, int Dh, float scale, void* stream) { - return flash::launch(q, k, v, ldv, gate, out, T, Hq, Hk, Dh, scale, stream); + return flash::launch(q, k, v, ldv, gate, out, T, T, Hq, Hk, Dh, scale, stream); +} + +extern "C" int cs1_attention_gated_cached(const void* q, const void* k, const void* v, int ldv, const void* gate, + void* out, int Tq, int Tk, int Hq, int Hk, int Dh, float scale, + void* stream) { + return flash::launch(q, k, v, ldv, gate, out, Tq, Tk, Hq, Hk, Dh, scale, stream); } diff --git a/src/backends/cuda/qwen3_5/elementwise.cu b/src/backends/cuda/qwen3_5/elementwise.cu index 97ae9c5a..2d9881ed 100644 --- a/src/backends/cuda/qwen3_5/elementwise.cu +++ b/src/backends/cuda/qwen3_5/elementwise.cu @@ -16,10 +16,13 @@ __global__ void embed_kernel(const int32_t* __restrict__ ids, const Pack8* __res } // F.conv1d in bfloat16 (float32 accumulation, rounded), then SiLU (rounded again), -// written to three contiguous outputs. +// written to three contiguous outputs. Positions before the first row come from +// history [3, channels] when there is one; without it they are skipped, as the zero +// padding of a sequence's start adds nothing. __global__ void gdn_conv_kernel(const bf16* __restrict__ qkv, int ld, const bf16* __restrict__ w, - bf16* __restrict__ q, bf16* __restrict__ k, bf16* __restrict__ v, int T, - int key_dim, int value_dim) { + const bf16* __restrict__ history, bf16* __restrict__ q, + bf16* __restrict__ k, bf16* __restrict__ v, int T, int key_dim, + int value_dim) { const int channels = 2 * key_dim + value_dim; const size_t idx = (size_t)blockIdx.x * blockDim.x + threadIdx.x; if (idx >= (size_t)T * channels) return; @@ -28,7 +31,10 @@ __global__ void gdn_conv_kernel(const bf16* __restrict__ qkv, int ld, const bf16 #pragma unroll for (int j = 0; j < 4; j++) { const int s = t - 3 + j; - if (s >= 0) acc = fmaf(f32(w[c * 4 + j]), f32(qkv[(size_t)s * ld + c]), acc); + if (s >= 0) + acc = fmaf(f32(w[c * 4 + j]), f32(qkv[(size_t)s * ld + c]), acc); + else if (history) + acc = fmaf(f32(w[c * 4 + j]), f32(history[(size_t)(3 + s) * channels + c]), acc); } const bf16 y = to_bf16(silu(round_bf16(acc))); if (c < key_dim) @@ -39,6 +45,16 @@ __global__ void gdn_conv_kernel(const bf16* __restrict__ qkv, int ld, const bf16 v[(size_t)t * value_dim + c - 2 * key_dim] = y; } +// The conv inputs of the last three positions, from qkv or, before its first row, from +// history (zeros without one). +__global__ void gdn_conv_history_kernel(const bf16* __restrict__ qkv, int ld, const bf16* __restrict__ history, + bf16* __restrict__ out, int T, int channels) { + const int idx = blockIdx.x * blockDim.x + threadIdx.x; + if (idx >= 3 * channels) return; + const int r = idx / channels, c = idx % channels, s = T - 3 + r; + out[idx] = s >= 0 ? qkv[(size_t)s * ld + c] : history ? history[(3 + s) * channels + c] : to_bf16(0.f); +} + // beta = sigmoid(b) in bfloat16; g = -exp(A_log) * softplus(a + dt_bias) in float32 // (F.softplus with threshold 20). __global__ void gdn_gates_kernel(const bf16* __restrict__ b, const bf16* __restrict__ a, int ld, @@ -101,17 +117,30 @@ extern "C" int cs1_embed(const int32_t* ids, const void* table, void* out, int T return cudaGetLastError(); } -extern "C" int cs1_gdn_conv(const void* qkv, int ld, const void* w, void* q, void* k, void* v, int T, int key_dim, - int value_dim, void* stream) { +extern "C" int cs1_gdn_conv_history(const void* qkv, int ld, const void* w, const void* history, void* history_out, + void* q, void* k, void* v, int T, int key_dim, int value_dim, void* stream) { if (T < 0 || key_dim < 0 || value_dim < 0 || ld < 2 * key_dim + value_dim) return cudaErrorInvalidValue; - const size_t n = (size_t)T * (2 * key_dim + value_dim); - if (n == 0) return cudaSuccess; - gdn_conv_kernel<<(stream)>>>( - static_cast(qkv), ld, static_cast(w), static_cast(q), - static_cast(k), static_cast(v), T, key_dim, value_dim); + const int channels = 2 * key_dim + value_dim; + const size_t n = (size_t)T * channels; + const cudaStream_t st = static_cast(stream); + if (n > 0) { + gdn_conv_kernel<<>>( + static_cast(qkv), ld, static_cast(w), static_cast(history), + static_cast(q), static_cast(k), static_cast(v), T, key_dim, value_dim); + } + if (history_out && channels > 0) { + gdn_conv_history_kernel<<>>( + static_cast(qkv), ld, static_cast(history), static_cast(history_out), + T, channels); + } return cudaGetLastError(); } +extern "C" int cs1_gdn_conv(const void* qkv, int ld, const void* w, void* q, void* k, void* v, int T, int key_dim, + int value_dim, void* stream) { + return cs1_gdn_conv_history(qkv, ld, w, nullptr, nullptr, q, k, v, T, key_dim, value_dim, stream); +} + extern "C" int cs1_gdn_gates(const void* b, const void* a, int ld, const void* A_log, const void* dt_bias, void* beta, float* g, int T, int H, void* stream) { if (T < 0 || H < 0 || ld < H) return cudaErrorInvalidValue; diff --git a/src/backends/cuda/qwen3_5/gdn_prefill.cu b/src/backends/cuda/qwen3_5/gdn_prefill.cu index 676ce446..67624dbb 100644 --- a/src/backends/cuda/qwen3_5/gdn_prefill.cu +++ b/src/backends/cuda/qwen3_5/gdn_prefill.cu @@ -8,7 +8,8 @@ // (bfloat16 with float32 accumulation). Results are stored as bfloat16. // 2. gdn_chunk_state, per (head, 32 value columns), over the chunks in order: keeps // the state S in float32 registers, stores it as bfloat16 before each chunk, and -// computes v_new = u - w S and S = decay S + kd^T v_new with mma.sync. +// computes v_new = u - w S and S = decay S + kd^T v_new with mma.sync. S starts +// from zero or from a given float32 state, and the final one can be written out. // 3. gdn_chunk_out, per (chunk, head): o = qd S + P v_new with mma.sync. // Transformers computes all of this in float32. Keeping the intermediate results in // bfloat16, as flash-linear-attention does, makes this kernel less precise than that @@ -337,7 +338,8 @@ constexpr int SS_LD = BVS + 8; // bfloat16 row stride of the S copy and v_n constexpr int STAGE = C * WS_LD; // elements of one staged w or kd constexpr size_t SMEM2_BYTES = (4 * STAGE + K * SS_LD + C * SS_LD) * 2; -__global__ void __launch_bounds__(ST_THREADS) gdn_chunk_state(Work ws, int NC) { +__global__ void __launch_bounds__(ST_THREADS) + gdn_chunk_state(Work ws, int NC, const float* s0, float* s_out) { // s0 and s_out may alias extern __shared__ __align__(128) unsigned char sm[]; bf16* wbuf = reinterpret_cast(sm); // [2][C][WS_LD] bf16* kbuf = wbuf + 2 * STAGE; // [2][C][WS_LD] @@ -356,12 +358,22 @@ __global__ void __launch_bounds__(ST_THREADS) gdn_chunk_state(Work ws, int NC) { cs1::cp_async_commit(); }; - // S rows warp * 32 + mt * 16 + {g, g + 8}, columns nt * 8 + {2t, 2t + 1} + // S rows warp * 32 + mt * 16 + {g, g + 8}, columns nt * 8 + {2t, 2t + 1}; the given + // state is float [H, K, V], read and written by the block that owns its columns float st[2][4][4]; + const size_t s_at = (size_t)h * K * V + vb0; #pragma unroll for (int mt = 0; mt < 2; mt++) #pragma unroll - for (int nt = 0; nt < 4; nt++) st[mt][nt][0] = st[mt][nt][1] = st[mt][nt][2] = st[mt][nt][3] = 0.f; + for (int nt = 0; nt < 4; nt++) +#pragma unroll + for (int r = 0; r < 2; r++) { + const int row = warp * 32 + mt * 16 + g + r * 8, col = nt * 8 + 2 * t; + const float2 x = s0 ? *reinterpret_cast(s0 + s_at + (size_t)row * V + col) + : make_float2(0.f, 0.f); + st[mt][nt][2 * r] = x.x; + st[mt][nt][2 * r + 1] = x.y; + } load(0, 0); for (int c = 0; c < NC; c++) { @@ -438,6 +450,18 @@ __global__ void __launch_bounds__(ST_THREADS) gdn_chunk_state(Work ws, int NC) { } } } + if (s_out) { +#pragma unroll + for (int mt = 0; mt < 2; mt++) +#pragma unroll + for (int nt = 0; nt < 4; nt++) +#pragma unroll + for (int r = 0; r < 2; r++) { + const int row = warp * 32 + mt * 16 + g + r * 8, col = nt * 8 + 2 * t; + *reinterpret_cast(s_out + s_at + (size_t)row * V + col) = + make_float2(st[mt][nt][2 * r], st[mt][nt][2 * r + 1]); + } + } } // ---- kernel 3: the output, per chunk ---- @@ -535,10 +559,19 @@ extern "C" { size_t cs1_gdn_workspace_floats(int T, int H) { return (Layout(T, H).total + 3) / 4; } -int cs1_gdn_prefill(const void* q, const void* k, const void* v, const float* g, const void* beta, - void* o, float* workspace, int T, int H, int HK, float scale, void* stream) { - if (T < 0 || HK <= 0 || H % HK != 0) return cudaErrorInvalidValue; - if (T == 0) return cudaSuccess; +int cs1_gdn_prefill_state(const void* q, const void* k, const void* v, const float* g, const void* beta, + void* o, float* workspace, const float* initial_state, float* final_state, int T, + int H, int HK, float scale, void* stream) { + if (T < 0 || H < 0 || HK <= 0 || H % HK != 0) return cudaErrorInvalidValue; + if ((reinterpret_cast(initial_state) | reinterpret_cast(final_state)) & 7) + return cudaErrorInvalidValue; + if (T == 0) { + if (!final_state || final_state == initial_state) return cudaSuccess; + const size_t bytes = (size_t)H * K * V * sizeof(float); + cudaStream_t st = static_cast(stream); + return initial_state ? cudaMemcpyAsync(final_state, initial_state, bytes, cudaMemcpyDeviceToDevice, st) + : cudaMemsetAsync(final_state, 0, bytes, st); + } // once per process (for the device current at the first call) static const cudaError_t configured = [] { cudaError_t e = cudaFuncSetAttribute(gdn_chunk_prep, cudaFuncAttributeMaxDynamicSharedMemorySize, @@ -558,9 +591,14 @@ int cs1_gdn_prefill(const void* q, const void* k, const void* v, const float* g, gdn_chunk_prep<<>>( static_cast(q), static_cast(k), static_cast(v), g, static_cast(beta), ws, T, H, HK, scale); - gdn_chunk_state<<>>(ws, NC); + gdn_chunk_state<<>>(ws, NC, initial_state, final_state); gdn_chunk_out<<>>(ws, static_cast(o), T, H); return cudaGetLastError(); } +int cs1_gdn_prefill(const void* q, const void* k, const void* v, const float* g, const void* beta, + void* o, float* workspace, int T, int H, int HK, float scale, void* stream) { + return cs1_gdn_prefill_state(q, k, v, g, beta, o, workspace, nullptr, nullptr, T, H, HK, scale, stream); +} + } // extern "C" diff --git a/src/backends/cuda/qwen3_5/gemm.cu b/src/backends/cuda/qwen3_5/gemm.cu index 154ad698..f0febca4 100644 --- a/src/backends/cuda/qwen3_5/gemm.cu +++ b/src/backends/cuda/qwen3_5/gemm.cu @@ -6,7 +6,14 @@ // // Each shape uses cuBLASLt's first heuristic choice, excluding split-K reductions that // accumulate into the output in place, since their order, and so the rounding, is not -// fixed. +// fixed (vision's FP32 and biased GEMMs exclude all split-K). That choice depends on M, +// so a row's result can change with the number of rows in the call. A handle from +// cs1_gemm_create_fixed instead keeps one algorithm per weight shape (N, K, ldy, and for +// vision's GEMMs the data type and bias) for every M: the heuristic's first choice at a +// reference M among algorithms without split-K. Each output row then takes the same +// path whatever the other rows, so its result does not depend on M or on its row index. +// cuBLASLt does not document this; tests/qwen3_5/kernels.rs checks it on the GPU it +// runs on. #include #include @@ -28,7 +35,9 @@ struct Gemm { cublasLtHandle_t handle = nullptr; void* workspace = nullptr; size_t workspace_bytes = 0; + int reference_m = 0; // > 0: one algorithm per weight shape, chosen at this M std::map plans; + std::map, cublasLtMatmulAlgo_t> fixed; // N, K, ldy, FP32, bias }; int status(cublasStatus_t s) { return s == CUBLAS_STATUS_SUCCESS ? 0 : 1000 + (int)s; } @@ -58,15 +67,13 @@ int describe(int M, int N, int K, int ldy, Plan& p, bool fp32, bool bias) { return 0; } -// The heuristic's first choice. Vision disables all split-K to avoid BF16 -// intermediate reductions; existing language GEMMs exclude only in-place reductions. -int first_choice(Gemm& g, Plan& p, bool vision) { +// The heuristic's first choice among the given reduction schemes. +int first_choice(Gemm& g, Plan& p, uint32_t schemes) { cublasLtMatmulPreference_t pref; cublasStatus_t s = cublasLtMatmulPreferenceCreate(&pref); if (s != CUBLAS_STATUS_SUCCESS) return status(s); cublasLtMatmulPreferenceSetAttribute(pref, CUBLASLT_MATMUL_PREF_MAX_WORKSPACE_BYTES, &g.workspace_bytes, sizeof(g.workspace_bytes)); - const uint32_t schemes = vision ? CUBLASLT_REDUCTION_SCHEME_NONE : (CUBLASLT_REDUCTION_SCHEME_MASK & ~CUBLASLT_REDUCTION_SCHEME_INPLACE); cublasLtMatmulPreferenceSetAttribute(pref, CUBLASLT_MATMUL_PREF_REDUCTION_SCHEME_MASK, &schemes, sizeof(schemes)); cublasLtMatmulHeuristicResult_t r{}; @@ -79,6 +86,40 @@ int first_choice(Gemm& g, Plan& p, bool vision) { return 0; } +// The weight shape's algorithm, chosen at the reference M on first use. An M it cannot +// serve is an error, not a reason to switch algorithms. +int fixed_choice(Gemm& g, int N, int K, int ldy, bool fp32, bool bias, Plan& p) { + const std::tuple key{N, K, ldy, fp32, bias}; + auto it = g.fixed.find(key); + if (it == g.fixed.end()) { + Plan r; + int rc = describe(g.reference_m, N, K, ldy, r, fp32, bias); + if (rc == 0) rc = first_choice(g, r, CUBLASLT_REDUCTION_SCHEME_NONE); + if (rc == 0) { + // Check that the choice really has no split-K reduction. + int32_t splits = 0; + uint32_t scheme = 0; + size_t written = 0; + const bool read = + cublasLtMatmulAlgoConfigGetAttribute(&r.algo, CUBLASLT_ALGO_CONFIG_SPLITK_NUM, &splits, + sizeof(splits), &written) == CUBLAS_STATUS_SUCCESS && + cublasLtMatmulAlgoConfigGetAttribute(&r.algo, CUBLASLT_ALGO_CONFIG_REDUCTION_SCHEME, &scheme, + sizeof(scheme), &written) == CUBLAS_STATUS_SUCCESS; + if (!read || splits != 1 || scheme != CUBLASLT_REDUCTION_SCHEME_NONE) + rc = status(CUBLAS_STATUS_NOT_SUPPORTED); + } + if (rc == 0) it = g.fixed.emplace(key, r.algo).first; + destroy(r); + if (rc != 0) return rc; + } + p.algo = it->second; + cublasLtMatmulHeuristicResult_t check{}; + const cublasStatus_t s = cublasLtMatmulAlgoCheck(g.handle, p.op, p.a, p.b, p.c, p.c, &p.algo, &check); + if (s != CUBLAS_STATUS_SUCCESS) return status(s); + if (check.workspaceSize > g.workspace_bytes) return status(CUBLAS_STATUS_NOT_SUPPORTED); + return 0; +} + // The plan for a shape, created on first use. int plan_for(Gemm& g, int M, int N, int K, int ldy, Plan*& out, bool fp32 = false, bool bias = false) { const Key key{M, N, K, ldy, fp32, bias}; @@ -89,7 +130,14 @@ int plan_for(Gemm& g, int M, int N, int K, int ldy, Plan*& out, bool fp32 = fals } Plan p; int rc = describe(M, N, K, ldy, p, fp32, bias); - if (rc == 0) rc = first_choice(g, p, fp32 || bias); + if (rc == 0 && g.reference_m > 0) { + rc = fixed_choice(g, N, K, ldy, fp32, bias, p); + } else if (rc == 0) { + // Vision's FP32 and biased GEMMs exclude all split-K to avoid BF16 intermediate + // reductions; the language GEMMs exclude only in-place reductions. + rc = first_choice(g, p, fp32 || bias ? CUBLASLT_REDUCTION_SCHEME_NONE + : CUBLASLT_REDUCTION_SCHEME_MASK & ~CUBLASLT_REDUCTION_SCHEME_INPLACE); + } if (rc != 0) { destroy(p); return rc; @@ -100,6 +148,13 @@ int plan_for(Gemm& g, int M, int N, int K, int ldy, Plan*& out, bool fp32 = fals } // namespace +extern "C" void* cs1_gemm_create_fixed(size_t workspace_bytes, int reference_m) { + if (reference_m <= 0) return nullptr; + Gemm* g = static_cast(cs1_gemm_create(workspace_bytes)); + if (g) g->reference_m = reference_m; + return g; +} + extern "C" void* cs1_gemm_create(size_t workspace_bytes) { Gemm* g = new Gemm(); if (cublasLtCreate(&g->handle) != CUBLAS_STATUS_SUCCESS || diff --git a/src/backends/cuda/qwen3_5/ops.h b/src/backends/cuda/qwen3_5/ops.h index 37a08a77..61412b25 100644 --- a/src/backends/cuda/qwen3_5/ops.h +++ b/src/backends/cuda/qwen3_5/ops.h @@ -13,7 +13,7 @@ #include // Bumped whenever the required interface below changes. -#define CS1_ABI_VERSION 5 +#define CS1_ABI_VERSION 7 #ifdef __cplusplus extern "C" { @@ -36,6 +36,10 @@ int cs1_graph_destroy(void* exec); // Copy and wait for the copy. int cs1_upload(void* dst, const void* src, size_t bytes, void* stream); int cs1_download(void* dst, const void* src, size_t bytes, void* stream); +// Queue a device-to-device copy of `rows` rows of `row_bytes` bytes, with row +// pitches in bytes; nothing waits for it. +int cs1_copy_rows(void* dst, size_t dst_pitch, const void* src, size_t src_pitch, size_t row_bytes, + int rows, void* stream); // ---- operations ---- @@ -59,6 +63,14 @@ int cs1_gated_rms_norm(const void* x, const void* z, int ldz, const void* w, voi int cs1_gdn_conv(const void* qkv, int ld, const void* w, void* q, void* k, void* v, int T, int key_dim, int value_dim, void* stream); +// The same conv continuing a sequence: history [3, key_dim*2 + value_dim] holds the conv +// inputs of the three positions before qkv's first row, oldest first (null: zeros, as at +// the start of a sequence). If history_out is not null, it receives the inputs of the +// last three positions in the same layout, taken from history where T < 3. history_out +// must not overlap history or qkv. Each output equals the unsplit conv's bit for bit. +int cs1_gdn_conv_history(const void* qkv, int ld, const void* w, const void* history, void* history_out, + void* q, void* k, void* v, int T, int key_dim, int value_dim, void* stream); + // beta = sigmoid(b) (bfloat16) and g = -exp(A_log) * softplus(a + dt_bias) (float32), [T, H]; // b and a are [T, H] in rows of ld. int cs1_gdn_gates(const void* b, const void* a, int ld, const void* A_log, const void* dt_bias, @@ -70,6 +82,16 @@ size_t cs1_gdn_workspace_floats(int T, int H); int cs1_gdn_prefill(const void* q, const void* k, const void* v, const float* g, const void* beta, void* o, float* workspace, int T, int H, int HK, float scale, void* stream); +// The same prefill continuing a sequence from initial_state (float [H, 128, 128], key by +// value per head; null: zeros) and, if final_state is not null, writing the state after +// the last token there in the same layout. Both are 8-byte aligned, and either the same +// buffer or not overlapping. With T = 0, final_state receives initial_state. When every +// split falls on a multiple of 64 tokens, the outputs match the unsplit prefill bit for +// bit; elsewhere the chunks fall differently. +int cs1_gdn_prefill_state(const void* q, const void* k, const void* v, const float* g, const void* beta, + void* o, float* workspace, const float* initial_state, float* final_state, int T, + int H, int HK, float scale, void* stream); + // Attention inputs: q and gate from qg [T, Hq, 2*Dh], k from kr [T, Hk, Dh], both in rows // of ld; per-head zero-centred RMSNorm, then rotary embedding on the first 2*half dims // using cos/sin [T, half] (bfloat16). Writes q [T, Hq, Dh], gate [T, Hq*Dh], k [T, Hk, Dh]. @@ -88,6 +110,14 @@ int cs1_attention(const void* q, const void* k, const void* v, int ldv, void* ou int cs1_attention_gated(const void* q, const void* k, const void* v, int ldv, const void* gate, void* out, int T, int Hq, int Hk, int Dh, float scale, void* stream); +// Gated attention for the last Tq of Tk positions: k [Tk, Hk, Dh] and v (Tk rows of ldv) +// cover all positions, while q, gate and out [Tq, Hq, Dh] are the queries at positions +// Tk - Tq onwards, each attending to the keys up to its own position. Keys are visited +// in the same order as in cs1_attention_gated, so each output row matches the unsplit call. +int cs1_attention_gated_cached(const void* q, const void* k, const void* v, int ldv, const void* gate, + void* out, int Tq, int Tk, int Hq, int Hk, int Dh, float scale, + void* stream); + // x = x * sigmoid(gate), n elements. int cs1_sigmoid_gate(void* x, const void* gate, size_t n, void* stream); @@ -97,6 +127,12 @@ int cs1_silu_mul(const void* gate_up, int ld, void* out, int T, int I, void* str // y [M, N] (rows of ldy) = x [M, K] * w [N, K]^T through cuBLASLt, float32 accumulation, // with cuBLASLt's first heuristic choice for each shape (see gemm.cu). void* cs1_gemm_create(size_t workspace_bytes); +// A handle that keeps one algorithm per weight shape (N, K, ldy) for every M, chosen at +// reference_m rows among algorithms without split-K, so that a row's result depends +// neither on M nor on its row index (checked by tests/qwen3_5/kernels.rs on the GPU it +// runs on). Null if reference_m <= 0 or setup fails. cs1_gemm with this handle returns a +// cuBLAS status for an M the algorithm cannot serve, rather than switching algorithms. +void* cs1_gemm_create_fixed(size_t workspace_bytes, int reference_m); void cs1_gemm_destroy(void* gemm); int cs1_gemm(void* gemm, const void* x, const void* w, void* y, int M, int N, int K, int ldy, void* stream); diff --git a/src/backends/cuda/qwen3_5/runtime.cu b/src/backends/cuda/qwen3_5/runtime.cu index 697c9db1..1a17b616 100644 --- a/src/backends/cuda/qwen3_5/runtime.cu +++ b/src/backends/cuda/qwen3_5/runtime.cu @@ -36,6 +36,14 @@ int cs1_download(void* dst, const void* src, size_t bytes, void* stream) { return e != cudaSuccess ? e : cudaStreamSynchronize(st); } +int cs1_copy_rows(void* dst, size_t dst_pitch, const void* src, size_t src_pitch, size_t row_bytes, + int rows, void* stream) { + if (rows < 0) return cudaErrorInvalidValue; + if (rows == 0 || row_bytes == 0) return cudaSuccess; + return cudaMemcpy2DAsync(dst, dst_pitch, src, src_pitch, row_bytes, rows, cudaMemcpyDeviceToDevice, + static_cast(stream)); +} + int cs1_graph_begin(void* stream) { const cudaError_t e = cudaStreamBeginCapture(static_cast(stream), cudaStreamCaptureModeThreadLocal); if (e != cudaSuccess) (void)cudaGetLastError(); diff --git a/src/models/decider/README.md b/src/models/decider/README.md new file mode 100644 index 00000000..6905153c --- /dev/null +++ b/src/models/decider/README.md @@ -0,0 +1,94 @@ +# Decider-2B v11 native worker + +The native Rust/CUDA worker serves text `choice`, `noul` and isolated-level `score` +through `POST /v1/systemone`. Follow the [native recipe](../../../recipe/decider/README.md). +The [validation protocol](../../../recipe/decider/validation.md) distinguishes +CPU contract checks, real checkpoint inference and HTTP checks. + +## Pinned artifacts and contract + +- Model/tokenizer/config: Mapika/decider-2b + `533964dae8be954c5b5e19fa4948e48408094c1e` (released 2B v11). +- Reference: Mapika/decider `50d0be0d7cb43d2066965ce5fa7f3fe4e489a60f`, decider-ai 1.8.1. +- BF16, 24 Qwen3.5 layers, hidden width 2048, 18 Gated DeltaNet and six full-attention + layers. Output embeddings are tied; there is no adapter merge or separate trained head. +- Plain state-first, independent questions; Score levels become separate no/yes rows. + Choice supports 2–255 single-token labels A through JT; Score supports 2–10 levels. +- Per-type temperatures are Choice 1.164, Noul 1.624 and Score 1.124, read from the + verified released configuration. The worker accepts only the original calibration + artifact. The CPU library's configuration API also supports explicit calibration + overrides, but those are outside the pinned worker's supported scope. +- Preserve ordered structured state/question rendering, long-array annotations, + per-question output ordering, confidence/certainty/legend/level-fit fields and + Python-style response rounding. `usage.input_tokens` counts the common prefix once; + `usage.output_tokens` is zero. Empty questions return empty answers without execution. + +State truncation caps the tokenized `Context:` prefix at 32,768 tokens; the question +suffix is additional. Native admission bounds complete rows to 36,864 tokens, +expanded rows to 1,024, total processed row tokens to 1,048,576 and raw bodies to +8 MiB. Unsupported modes, duplicate JSON keys, depth >=128, nonfinite literals +and integers outside i64/u64 are rejected. No image/video, chat/schema-first, +packed-question execution, neutralization, quantization, CPU/Metal inference, +cross-request prefix cache is included. CUDA Graph replay and request-local prefix +reuse are optional and default off. Prefix integration is a pending integration until +Qwen PRs #97/#98/#99 merge. + +## Ownership and execution + +`processing.rs` compiles the entire request before device work. Prepared rows own +unpadded token IDs, final-position readout, candidate IDs and original identities; +`ResponseContext` owns whole-question calibration and answer reconstruction. + +`checkpoint.rs` verifies fixed sizes and SHA-256 hashes for the checkpoint, +configuration, calibration and tokenizer before CUDA initialization. It rejects +sharded indices that could redirect the shared loader away from the verified file. +Keep artifacts immutable for the entire worker lifetime. Selected BF16 output rows +are loaded from the pinned tied input embedding; all 255 rows are retained, with a +256th zero row for aligned CUDA GEMM storage. + +`executor.rs` owns the shared Qwen model and a separate CUDA projection stream, +GEMM handle and persistent head buffers. Each complete request gets one shared +`SerialScheduler` permit. Default execution runs rows eagerly in order. Optional request-local packing +uses `DECIDER_BATCH_MAX_ROWS` (1–4, default 1) and +`DECIDER_BATCH_MAX_TOKENS` (1–4096, default 4096). Contiguous complete rows are +packed with the shared backbone `forward_batch`; each sequence resets its own +GDN/attention state. A row above the packing budget runs alone under the unchanged +complete-row limit. A persistent head projects all batch hidden rows in one GEMM. +Score levels can cross batch boundaries; response normalization still uses the +complete question. This does not combine separate requests or change admission. The BF16 selected projection uses existing backend GEMM, then converts +its BF16 outputs to FP32 before CPU temperature scaling and normalization. Padding +never participates in softmax or token accounting. Current GEMM/reduction algorithms +can differ from Transformers; numerical agreement is measured rather than assumed. + +Decider uses an explicit shared-model loader option controlled only by +`DECIDER_GRAPH`; existing workers retain their `CUA_S1_GRAPH` constructor behavior. +The backbone caches up to 64 ordered sequence-length shapes, clears them before +scratch growth, and falls back permanently to eager if capture fails. Effective +mode and cumulative capture/replay counters appear in health/diagnostic metadata. +The selected head and calibrated response remain outside the graph. + +`prefix.rs` independently plans exact request/question token prefixes from complete +prepared rows. `DECIDER_PREFIX=1` uses shared Qwen continuation and fixed GEMM +selection; `DECIDER_FIXED=1` provides independent fixed full rows for a matched +control. Short/no-sharing cases run independent fixed rows. These paths require +Graph off. The selected head keeps identical batch grouping; calibrated responses +and unique-prefix usage remain unchanged. Prefix snapshots/KV are overwritten per +request and never reused across calls. The backend/consumers require ABI7 rebuilds. + +Both streams synchronize on completion, errors and caught execution panics before +admission is released. An execution failure retires the loaded model/head and makes +health unavailable. Validation errors preserve readiness. Queued cancellation +removes the caller; after dispatch, resources and admission remain held until work +completes. FIFO waiting is serial and inherits the current runtime's unbounded +pending queue; no new queue/token scheduling policy is claimed. + +`engine.rs` assembles the processor/executor and performs a real warmup before +binding a socket. `serve.rs` handles transport, status codes and worker health. +HTTP CPU preparation uses a blocking task outside GPU admission and the model lock. +The existing Rust frontend forwards bytes to this separately running worker. + +## Tests + +All test bodies and fixtures live under root `tests/decider`. Normal Cargo tests +require neither CUDA nor downloads. Explicit checkpoint/accelerator tests are +registered and ignored until their prerequisites are supplied; see the recipe. diff --git a/src/models/decider/native/Cargo.toml b/src/models/decider/native/Cargo.toml new file mode 100644 index 00000000..ca279132 --- /dev/null +++ b/src/models/decider/native/Cargo.toml @@ -0,0 +1,51 @@ +[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" + +[[test]] +name = "prefix" +path = "../../../../tests/decider/prefix.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..3d89dd00 --- /dev/null +++ b/src/models/decider/native/src/engine.rs @@ -0,0 +1,125 @@ +//! 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 { + let graph = crate::options::graph_env()?; + Self::load_with_modes( + dir, + library, + BatchLimits::from_env()?, + graph, + crate::options::prefix_env(graph)?, + ) + .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 { + Self::load_with_modes( + dir, + library, + batching, + graph, + crate::prefix::PrefixMode::Off, + ) + .await + } + pub async fn load_with_modes( + dir: &Path, + library: &Path, + batching: BatchLimits, + graph: bool, + prefix_mode: crate::prefix::PrefixMode, + ) -> Result { + anyhow::ensure!( + !graph || prefix_mode == crate::prefix::PrefixMode::Off, + "shared/fixed execution requires Graph off" + ); + 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, + prefix_mode, + ) + .await?; + metadata["batch_limits"] = + json!({"rows":batching.max_rows(),"tokens":batching.max_tokens()}); + metadata["prefix_mode"] = json!(prefix_mode.name()); + 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["prefix"] = json!(self.executor.prefix_stats()); + if self.executor.is_ready() && metadata["prefix_mode"] != "off" { + metadata["execution"] = metadata["prefix_mode"].clone(); + } + 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..5d73735b --- /dev/null +++ b/src/models/decider/native/src/executor.rs @@ -0,0 +1,373 @@ +//! Eager complete-row prefills and a selected tied-embedding CUDA readout. +use crate::{ + Limits, RowInput, + batching::BatchLimits, + checkpoint::{Checkpoint, HIDDEN, LABELS, PADDED_LABELS, VOCAB}, + prefix::{PrefixMode, PrefixPlan, PrefixStats}, +}; +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, + prefix_mode: PrefixMode, + prefix_stats: PrefixStats, +} +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 plan = match self.prefix_mode { + PrefixMode::Shared => PrefixPlan::new(rows), + PrefixMode::Auto => PrefixPlan::new(rows).filter(|plan| plan.worth_auto(rows)), + _ => None, + }; + if self.prefix_mode == PrefixMode::Fixed + || self.prefix_mode == PrefixMode::Shared + || plan.is_some() + { + let hidden = if let Some(plan) = plan { + let hidden = self + .model + .forward_shared(&plan.prompts)? + .into_iter() + .flatten() + .collect::>(); + self.prefix_stats.shared_requests += 1; + self.prefix_stats.saved_tokens += plan.saved_tokens as u64; + hidden + } else { + rows.iter() + .map(|row| self.model.forward_fixed(&row.ids)) + .collect::>>()? + }; + ensure!( + hidden.len() == rows.len(), + "shared backbone output count mismatch" + ); + let mut logits = Vec::with_capacity(rows.len()); + for range in self.batching.ranges(&lengths) { + let batch = &rows[range.clone()]; + let hidden: Vec<&[f32]> = hidden[range].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)?); + } + return Ok(logits); + } + 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)?); + } + } + if self.prefix_mode == PrefixMode::Auto { + self.prefix_stats.auto_independent_requests += 1; + } + 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>, + prefix_stats: Arc>, +} +impl Executor { + pub(crate) async fn load( + dir: &Path, + library: &Path, + checkpoint: Checkpoint, + labels: Vec, + batching: BatchLimits, + graph: bool, + prefix_mode: PrefixMode, + ) -> Result { + ensure!( + !graph || prefix_mode == PrefixMode::Off, + "shared/fixed execution requires Graph off" + ); + 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, + prefix_mode, + prefix_stats: PrefixStats::default(), + }) + }) + .await??; + let graph_stats = Arc::new(Mutex::new(loaded.model.graph_stats())); + Ok(Self { + graph_stats, + prefix_stats: Arc::new(Mutex::new(PrefixStats::default())), + 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 prefix_stats(&self) -> PrefixStats { + *self + .prefix_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(), + ); + let prefix_stats = self.prefix_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); + *prefix_stats + .lock() + .unwrap_or_else(|error| error.into_inner()) = loaded.prefix_stats; + 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..3fd6067d --- /dev/null +++ b/src/models/decider/native/src/lib.rs @@ -0,0 +1,42 @@ +//! 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; + +pub mod prefix; 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..8d9df440 --- /dev/null +++ b/src/models/decider/native/src/options.rs @@ -0,0 +1,62 @@ +//! 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()) +} + +pub fn prefix_values( + prefix: Option<&str>, + fixed: Option<&str>, + graph: bool, +) -> Result { + use crate::prefix::PrefixMode; + let switch = |name, value| match value { + None | Some("0") => Ok(false), + Some("1") => Ok(true), + Some(_) => bail!("{name} must be 0 or 1"), + }; + let prefix = match prefix { + None | Some("0") => PrefixMode::Off, + Some("1") => PrefixMode::Shared, + Some("auto") => PrefixMode::Auto, + Some(_) => bail!("DECIDER_PREFIX must be 0, 1 or auto"), + }; + let fixed = switch("DECIDER_FIXED", fixed)?; + anyhow::ensure!( + !(prefix != PrefixMode::Off && fixed), + "select only one of DECIDER_PREFIX and DECIDER_FIXED" + ); + anyhow::ensure!( + !(graph && (prefix != PrefixMode::Off || fixed)), + "shared/fixed execution requires DECIDER_GRAPH=0" + ); + Ok(if prefix != PrefixMode::Off { + prefix + } else if fixed { + PrefixMode::Fixed + } else { + PrefixMode::Off + }) +} +pub fn prefix_env(graph: bool) -> Result { + let get = |name| match std::env::var(name) { + Ok(value) => Ok(Some(value)), + Err(std::env::VarError::NotPresent) => Ok(None), + Err(error) => Err(error), + }; + let prefix = get("DECIDER_PREFIX")?; + let fixed = get("DECIDER_FIXED")?; + prefix_values(prefix.as_deref(), fixed.as_deref(), graph) +} diff --git a/src/models/decider/native/src/prefix.rs b/src/models/decider/native/src/prefix.rs new file mode 100644 index 00000000..ea407cd4 --- /dev/null +++ b/src/models/decider/native/src/prefix.rs @@ -0,0 +1,90 @@ +//! Request-local exact token prefix planning; no state is cached across requests. +use crate::RowInput; +use omni_qwen3_5_native::model::{PromptGroup, SharedPrompts}; +pub struct PrefixPlan<'a> { + pub prompts: SharedPrompts<'a>, + pub saved_tokens: usize, +} +fn common(rows: &[RowInput], skip: usize) -> usize { + let first = &rows[0].ids; + let limit = rows + .iter() + .map(|row| row.ids.len().saturating_sub(1)) + .min() + .unwrap(); + (skip..limit) + .take_while(|&i| rows.iter().all(|row| row.ids[i] == first[i])) + .count() +} +impl<'a> PrefixPlan<'a> { + /// Conservative opt-in policy; savings are aligned token work, not a speed guarantee. + pub fn worth_auto(&self, rows: &[RowInput]) -> bool { + let total = rows.iter().map(|row| row.ids.len()).sum::(); + self.saved_tokens >= 4096 && self.saved_tokens >= total.div_ceil(3) + } + + pub fn new(rows: &'a [RowInput]) -> Option { + if rows.len() < 2 || rows.iter().any(|row| row.ids.is_empty()) { + return None; + } + let request = common(rows, 0); + let p = request / 64 * 64; + let mut saved_tokens = (rows.len() - 1) * p; + let mut groups = Vec::new(); + let mut start = 0; + while start < rows.len() { + let mut end = start + 1; + while end < rows.len() && rows[end].question_id == rows[start].question_id { + end += 1; + } + let group = &rows[start..end]; + let prefix = if group.len() > 1 { + common(group, request) + } else { + 0 + }; + let q = (request + prefix) / 64 * 64; + saved_tokens += (group.len() - 1) * (q - p); + groups.push(PromptGroup { + prefix: &group[0].ids[request..request + prefix], + branches: group + .iter() + .map(|row| &row.ids[request + prefix..]) + .collect(), + }); + start = end; + } + (saved_tokens > 0).then(|| Self { + prompts: SharedPrompts { + prefix: &rows[0].ids[..request], + groups, + }, + saved_tokens, + }) + } +} + +#[derive(Clone, Copy, Debug, Default, serde::Serialize)] +pub struct PrefixStats { + pub shared_requests: u64, + pub auto_independent_requests: u64, + pub saved_tokens: u64, +} +#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)] +pub enum PrefixMode { + #[default] + Off, + Fixed, + Shared, + Auto, +} +impl PrefixMode { + pub fn name(self) -> &'static str { + match self { + Self::Off => "off", + Self::Fixed => "fixed", + Self::Shared => "shared", + Self::Auto => "auto", + } + } +} 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