Spaces:
Running
Running
refactor(decision): use direct serving mode without legacy aliases
Browse filesSigned-off-by: Xunzhuo <Xunzhuo@users.noreply.huggingface.co>
- DEPLOYMENT.md +21 -23
- Dockerfile +2 -2
- README.md +1 -1
- SYSTEMONE_API.md +4 -4
- SYSTEM_ONE_MAPPING.md +3 -3
- app.py +40 -40
- drun_contract.py → direct_contract.py +1 -1
- drun_gateway.py → direct_gateway.py +49 -49
- static/app.js +2 -2
- tests/studio_contract.test.mjs +2 -2
- tests/{test_drun_gateway.py → test_direct_gateway.py} +46 -33
- tests/test_public_api_contract.py +0 -2
- tetris_arena.py +4 -4
DEPLOYMENT.md
CHANGED
|
@@ -1,47 +1,45 @@
|
|
| 1 |
# Decision Gateway deployment
|
| 2 |
|
| 3 |
-
## Direct
|
| 4 |
|
| 5 |
-
`DECISION_BACKEND=
|
| 6 |
|
| 7 |
- `POST /v1/systemone`
|
| 8 |
- `POST /v1/systemone/batches`
|
| 9 |
|
| 10 |
Public requests must include one exact canonical model ID advertised by `/v1/models`. The Studio-only `/api/evaluate` route requires a configured short wire ID or canonical ID; the Gateway converts it to the canonical ID before forwarding and returns the strict backend response unchanged. Direct mode advertises no default model in `/v1/models`; the Studio editor chooses its model locally and sends it explicitly. There is no default-model inference and no fallback to another model, the legacy native engine, or a queue worker.
|
| 11 |
|
| 12 |
-
Set the backend, the pinned `DECISION_MODEL_REGISTRY_V2`, and a server-only JSON map whose keys cover that registry exactly. Every direct-mode registry entry must include the full immutable 40-character Hub `revision` and 64-character `manifest_sha256`; it may also pin the runtime's 64-character `content_sha256`. These values come from each released artifact and may change through deployment configuration without a Gateway source release.
|
| 13 |
|
| 14 |
-
`/v1/models` reads presentation dates from the validated `MODEL_RELEASES.json` metadata shipped with the Gateway image. Its top-level `release` is the default date; a model row can set its own ISO `release_date`. A date is published only when that row's canonical ID, revision, and manifest digest match the active registry. An artifact with newer deployment pins still serves normally when the presentation file has not caught up; its date is omitted rather than reported incorrectly. If the file is missing or invalid, `/v1/models` still lists every configured model but omits all release dates; direct serving is unaffected. Operators can call `load_model_releases()` explicitly to fail on invalid metadata. Update the JSON metadata and rebuild the Gateway image when the public date should change. This file does not allowlist model files or decide which artifacts
|
| 15 |
|
| 16 |
```text
|
| 17 |
-
DECISION_BACKEND=
|
| 18 |
-
|
| 19 |
-
|
| 20 |
```
|
| 21 |
|
| 22 |
-
The example assumes a two-model registry; the endpoint map must include every model in the active registry, and every model must have a distinct origin. A
|
| 23 |
|
| 24 |
-
- Run the Gateway with host networking on the same trusted Linux host, keep every
|
| 25 |
-
- Attach the Gateway and model containers to a dedicated private network, bind each
|
| 26 |
- For a remote Space, establish an authenticated private overlay or tunnel first and use its private IP literals. Direct mode is not ready if the Space can reach the models only through public origins.
|
| 27 |
|
| 28 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
| 29 |
|
| 30 |
-
```
|
| 31 |
-
DECISION_BACKEND=drun
|
| 32 |
-
DECISION_DRUN_ENDPOINTS={"llm-semantic-router/Decision-1.0-Kai-0.6B":"http://10.90.0.11:8101","llm-semantic-router/Decision-1.0-Lux-9B":"http://10.90.0.12:8106"}
|
| 33 |
-
DECISION_DRUN_TIMEOUT_SECONDS=45
|
| 34 |
-
```
|
| 35 |
-
|
| 36 |
-
Store `DECISION_DRUN_ENDPOINTS` as a private Space secret or equivalent server-only setting because it describes private topology. Values must be bare `http` or `https` origins with an explicit port and a literal loopback, RFC 1918, or IPv6 unique-local address. DNS names, public addresses, link-local addresses, URL credentials, paths, queries, and fragments fail startup. `DECISION_DRUN_TIMEOUT_SECONDS` is bounded to 0.1–300 seconds.
|
| 37 |
|
| 38 |
-
The direct transport ignores ambient proxy variables, does not follow redirects, requests identity encoding, reuses private keep-alive connections, bounds single requests at 256 KiB, batch requests at 2 MiB, and inference responses at 8 MiB. Bind each
|
| 39 |
|
| 40 |
-
Each successful private inference response carries `X-Decision-Artifact-Model`, `X-Decision-Artifact-Revision`, `X-Decision-Artifact-Manifest-Sha256`, and `X-Decision-Artifact-Content-Sha256` headers from the resident
|
| 41 |
|
| 42 |
`GET /api/ready` probes every configured model and returns HTTP 200 only when all are loaded with the pinned artifact; otherwise it returns HTTP 503 with public model IDs and Boolean readiness, never private origins. The Docker health check uses this aggregate route in direct mode, while legacy queue deployments keep their existing `/api/status` health check. The Studio menu reports an unavailable model distinctly from one still loading. Tetris configuration probes live Decision readiness before offering a local competitor, and a race start rechecks each selected local model so a disconnected instance returns 503 before the runners start. These control probes occur at page load and race admission, not on every game turn.
|
| 43 |
|
| 44 |
-
For rollout, start and warm every pinned
|
| 45 |
|
| 46 |
The public reverse proxy must route `/tetris/` and `/api/tetris/*` to the same pinned Gateway release. Preserve the `/tetris` path when forwarding: the new Gateway serves the page at `/tetris/`, while an older standalone Tetris service used a stripped prefix. Keep `/v1/systemone` bounded at 256 KiB and admit up to 2 MiB at `/v1/systemone/batches`; a site-wide 256 KiB body limit would reject valid batch traffic before it reaches this application. Align the proxy's response-header timeout with the selected synchronous runtime timeout. Validate the proxy configuration before reload, retain its exact prior version, and test both paths before retiring the old page.
|
| 47 |
|
|
@@ -108,7 +106,7 @@ Public single-state request bodies remain bounded at 256 KiB; public batch bodie
|
|
| 108 |
|
| 109 |
`/tetris/` is a passive SSE client for server-authoritative races. Starting a race creates two independent asynchronous runners immediately; connecting, reconnecting, or slowing the event stream does not pause either model. A disconnect starts a bounded reconnection lease, after which an abandoned race is cancelled. Only the seeded tetromino sequence is shared. The server retains separate board snapshots, request states, model actions, and traces at `/api/tetris/races/{id}/trace`. The browser paints each side from its own bounded step queue; a short 40-step burst keeps every step, while a longer burst may skip intermediate visual frames without changing the server race, final board, or result. The server records side completion order in the SSE result and snapshot so recovery preserves the real finishing order. An illegal model choice ends that side and is never replaced by a heuristic choice or another model.
|
| 110 |
|
| 111 |
-
All endpoints and credentials are server-only Space variables or secrets. In `DECISION_BACKEND=
|
| 112 |
|
| 113 |
Configure Jev Cloud with `TETRIS_JEV_API_URL`, the `TETRIS_JEV_API_KEY` secret, and optional `TETRIS_JEV_MODEL`. Configure the alternate SystemOne provider with `TETRIS_SYSTEMONE_API_URL`, the `TETRIS_SYSTEMONE_API_KEY` secret, and optional `TETRIS_SYSTEMONE_MODEL`. Existing `JEV_API_*`, `JEV_MIRROR_API_*`, and `LOCAL_API_URL` variables remain accepted while the deployment migrates. The `jev-latest` alias accepts only itself or a versioned `jev-x.y.z` response; an explicitly pinned Jev model must match exactly. `/api/tetris/config` exposes only readiness and display metadata; it never returns endpoint URLs or credentials.
|
| 114 |
|
|
@@ -116,6 +114,6 @@ Step races accept 1 through `TETRIS_MAX_STEPS` pieces. The default server limit
|
|
| 116 |
|
| 117 |
Paid Cloud competitors additionally share one process-wide call budget across all clients. `TETRIS_CLOUD_CALLS_PER_WINDOW` defaults to 480 attempted calls, with a `TETRIS_CLOUD_CALL_WINDOW_SECONDS` default of 600 seconds. Admission reserves one call per possible step for each selected Cloud side (two for Cloud-vs-Cloud); a completed, failed, or cancelled race settles unused calls back to the budget. A race that cannot reserve its full possible Cloud cost receives HTTP 429 before any provider call. A local-vs-local race reserves no paid Cloud calls, including at the 2,000-step limit. Keep the default 120-step ceiling for public Cloud races unless the paid-provider budget is deliberately raised. The ledger is in memory and protects only this one Gateway process: deployments with multiple replicas or processes need a shared atomic budget store or an upstream account-level rate limit before claiming a fleet-wide cap.
|
| 118 |
|
| 119 |
-
`TETRIS_RACE_TIMEOUT_SECONDS` sets a 1–3,600 second whole-race deadline (default 900); raise it deliberately for slow long runs. `TETRIS_STREAM_GRACE_SECONDS` sets the 1–300 second reconnection lease after the last event stream disconnects (default 30); reconnecting resumes retained SSE events, and the browser can recover the latest board and terminal result from the race snapshot. An abandoned race is cancelled. `TETRIS_COMPLETED_TTL_SECONDS` controls completed-race retention from 30–3,600 seconds (default 300). Shutdown cancels and awaits every active race before closing the provider transport. A failed, timed-out, or otherwise incomplete race does not award a score winner. The speed percentage is available only when both sides complete the same step target; it compares the sums of complete server-observed adapter calls. For
|
| 120 |
|
| 121 |
`TETRIS_REQUEST_TIMEOUT_SECONDS` bounds each server-to-provider call to 0.1–300 seconds (default 45). The adapter sends optional bearer credentials only from server configuration, rejects redirects, ignores ambient proxy settings, bounds response bytes, and never returns provider URLs or credentials to the browser. Validate each configured endpoint with its selected canonical model and expected response identity before enabling that competitor; publishing a Space commit rebuilds the live Space.
|
|
|
|
| 1 |
# Decision Gateway deployment
|
| 2 |
|
| 3 |
+
## Direct serving mode
|
| 4 |
|
| 5 |
+
`DECISION_BACKEND=direct` is the production direct-routing mode. The Space keeps no model weights: it validates the public request, resolves its exact canonical Hugging Face model ID, and sends it to that model's private `vllm-sr decision serve` origin. Both public inference routes keep the same path end to end:
|
| 6 |
|
| 7 |
- `POST /v1/systemone`
|
| 8 |
- `POST /v1/systemone/batches`
|
| 9 |
|
| 10 |
Public requests must include one exact canonical model ID advertised by `/v1/models`. The Studio-only `/api/evaluate` route requires a configured short wire ID or canonical ID; the Gateway converts it to the canonical ID before forwarding and returns the strict backend response unchanged. Direct mode advertises no default model in `/v1/models`; the Studio editor chooses its model locally and sends it explicitly. There is no default-model inference and no fallback to another model, the legacy native engine, or a queue worker.
|
| 11 |
|
| 12 |
+
Set the backend, the pinned `DECISION_MODEL_REGISTRY_V2`, and a server-only JSON map whose keys cover that registry exactly. Every direct-mode registry entry must include the full immutable 40-character Hub `revision` and 64-character `manifest_sha256`; it may also pin the runtime's 64-character `content_sha256`. These values come from each released artifact and may change through deployment configuration without a Gateway source release. The loopback example below works only when the Gateway and every Decision runtime listener share the host network namespace:
|
| 13 |
|
| 14 |
+
`/v1/models` reads presentation dates from the validated `MODEL_RELEASES.json` metadata shipped with the Gateway image. Its top-level `release` is the default date; a model row can set its own ISO `release_date`. A date is published only when that row's canonical ID, revision, and manifest digest match the active registry. An artifact with newer deployment pins still serves normally when the presentation file has not caught up; its date is omitted rather than reported incorrectly. If the file is missing or invalid, `/v1/models` still lists every configured model but omits all release dates; direct serving is unaffected. Operators can call `load_model_releases()` explicitly to fail on invalid metadata. Update the JSON metadata and rebuild the Gateway image when the public date should change. This file does not allowlist model files or decide which artifacts the Decision runtime may load.
|
| 15 |
|
| 16 |
```text
|
| 17 |
+
DECISION_BACKEND=direct
|
| 18 |
+
DECISION_DIRECT_ENDPOINTS={"llm-semantic-router/Decision-1.0-Kai-0.6B":"http://127.0.0.1:8101","llm-semantic-router/Decision-1.0-Lux-9B":"http://127.0.0.1:8106"}
|
| 19 |
+
DECISION_DIRECT_TIMEOUT_SECONDS=45
|
| 20 |
```
|
| 21 |
|
| 22 |
+
The example assumes a two-model registry; the endpoint map must include every model in the active registry, and every model must have a distinct origin. A Decision runtime instance binds to host loopback by default. `127.0.0.1` inside a separate Gateway or Tetris container is that container itself, **not** the model host, so loopback endpoints do not work across ordinary container namespaces. Use one of these explicit deployment boundaries:
|
| 23 |
|
| 24 |
+
- Run the Gateway with host networking on the same trusted Linux host, keep every Decision runtime listener on host loopback, and retain the loopback map above.
|
| 25 |
+
- Attach the Gateway and model containers to a dedicated private network, bind each Decision runtime listener intentionally to its assigned private interface, and admit only the Gateway with host firewall/container-network policy. Never publish those ports on a public interface.
|
| 26 |
- For a remote Space, establish an authenticated private overlay or tunnel first and use its private IP literals. Direct mode is not ready if the Space can reach the models only through public origins.
|
| 27 |
|
| 28 |
+
For a private-network deployment, keep the same JSON map structure and replace
|
| 29 |
+
each loopback origin with the private IP origin assigned to that model in your
|
| 30 |
+
environment. Configure the active registry with the same model IDs. Store the
|
| 31 |
+
map only in the server-side secret setting; no private-network addresses belong
|
| 32 |
+
in this repository.
|
| 33 |
|
| 34 |
+
Store `DECISION_DIRECT_ENDPOINTS` as a private Space secret or equivalent server-only setting because it describes private topology. Values must be bare `http` or `https` origins with an explicit port and a literal loopback, RFC 1918, or IPv6 unique-local address. DNS names, public addresses, link-local addresses, URL credentials, paths, queries, and fragments fail startup. `DECISION_DIRECT_TIMEOUT_SECONDS` is bounded to 0.1–300 seconds.
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 35 |
|
| 36 |
+
The direct transport ignores ambient proxy variables, does not follow redirects, requests identity encoding, reuses private keep-alive connections, bounds single requests at 256 KiB, batch requests at 2 MiB, and inference responses at 8 MiB. Bind each Decision runtime service to loopback or enforce a network policy that admits only the Gateway; direct mode intentionally does not forward browser authorization. Neither backend origins nor private configuration are returned by `/api/status`, `/v1/models`, error responses, or browser configuration.
|
| 37 |
|
| 38 |
+
Each successful private inference response carries `X-Decision-Artifact-Model`, `X-Decision-Artifact-Revision`, `X-Decision-Artifact-Manifest-Sha256`, and `X-Decision-Artifact-Content-Sha256` headers from the resident Decision runtime process. The Gateway requires exactly one well-formed value for each, matches the canonical model, pinned Hub revision and manifest digest, and matches the content digest when configured. Missing, duplicate, malformed, or mismatched identity discards that response. It also verifies the exact requested model, strict official JSON field set, answer IDs and types, state order, probability distributions, Choice winner, Score legend/value, type-aware Decision confidence, and batch usage aggregation. Confidence is accepted only from the validated runtime response and is never replaced with `max(probabilities)`. Explicit `/api/status?model=<wire-id>` probes report readiness plus the configured and live artifact digests without disclosing private origins. A Gateway using response attestation requires the matching Decision runtime version; older runtimes fail closed.
|
| 39 |
|
| 40 |
`GET /api/ready` probes every configured model and returns HTTP 200 only when all are loaded with the pinned artifact; otherwise it returns HTTP 503 with public model IDs and Boolean readiness, never private origins. The Docker health check uses this aggregate route in direct mode, while legacy queue deployments keep their existing `/api/status` health check. The Studio menu reports an unavailable model distinctly from one still loading. Tetris configuration probes live Decision readiness before offering a local competitor, and a race start rechecks each selected local model so a disconnected instance returns 503 before the runners start. These control probes occur at page load and race admission, not on every game turn.
|
| 41 |
|
| 42 |
+
For rollout, start and warm every pinned Decision runtime instance first, verify reachability **from the Gateway network namespace**, and verify that each private origin's live artifact report matches its assigned model and deployment pins. Then switch the Space to `DECISION_BACKEND=direct`. Keep the prior queue workers and their exact settings available but inactive until direct-mode probes and representative single/batch calls pass. A response identity mismatch is a deployment failure, never a reason to replay against another backend.
|
| 43 |
|
| 44 |
The public reverse proxy must route `/tetris/` and `/api/tetris/*` to the same pinned Gateway release. Preserve the `/tetris` path when forwarding: the new Gateway serves the page at `/tetris/`, while an older standalone Tetris service used a stripped prefix. Keep `/v1/systemone` bounded at 256 KiB and admit up to 2 MiB at `/v1/systemone/batches`; a site-wide 256 KiB body limit would reject valid batch traffic before it reaches this application. Align the proxy's response-header timeout with the selected synchronous runtime timeout. Validate the proxy configuration before reload, retain its exact prior version, and test both paths before retiring the old page.
|
| 45 |
|
|
|
|
| 106 |
|
| 107 |
`/tetris/` is a passive SSE client for server-authoritative races. Starting a race creates two independent asynchronous runners immediately; connecting, reconnecting, or slowing the event stream does not pause either model. A disconnect starts a bounded reconnection lease, after which an abandoned race is cancelled. Only the seeded tetromino sequence is shared. The server retains separate board snapshots, request states, model actions, and traces at `/api/tetris/races/{id}/trace`. The browser paints each side from its own bounded step queue; a short 40-step burst keeps every step, while a longer burst may skip intermediate visual frames without changing the server race, final board, or result. The server records side completion order in the SSE result and snapshot so recovery preserves the real finishing order. An illegal model choice ends that side and is never replaced by a heuristic choice or another model.
|
| 108 |
|
| 109 |
+
All endpoints and credentials are server-only Space variables or secrets. In `DECISION_BACKEND=direct` mode, Tetris selects local models from the same private `DECISION_DIRECT_ENDPOINTS` map and sends each turn through the Gateway's attested direct client. It ignores legacy `TETRIS_LOCAL_*_API_URL` and local bearer variables. In the legacy modes, configure the shared Decision gateway with `TETRIS_LOCAL_API_URL` pointing to its `/v1/systemone` route, or configure per-model `TETRIS_LOCAL_LUX_API_URL`, `TETRIS_LOCAL_NOX_API_URL`, `TETRIS_LOCAL_SOL_API_URL`, `TETRIS_LOCAL_EOS_API_URL`, `TETRIS_LOCAL_KAI_API_URL`, and `TETRIS_LOCAL_LEX_API_URL`; optional local bearer secrets use the corresponding `_API_KEY` name or the shared `TETRIS_LOCAL_API_KEY`. A Decision response must return the exact Hugging Face canonical model ID selected by the player, such as `llm-semantic-router/Decision-1.0-Lux-9B`; a different or missing identity ends that side instead of being displayed under the wrong model name. This strict response identity is a deployment prerequisite. Legacy wire IDs such as `decision-lux` are not compatibility aliases here and deliberately fail closed during integration with an outdated Gateway.
|
| 110 |
|
| 111 |
Configure Jev Cloud with `TETRIS_JEV_API_URL`, the `TETRIS_JEV_API_KEY` secret, and optional `TETRIS_JEV_MODEL`. Configure the alternate SystemOne provider with `TETRIS_SYSTEMONE_API_URL`, the `TETRIS_SYSTEMONE_API_KEY` secret, and optional `TETRIS_SYSTEMONE_MODEL`. Existing `JEV_API_*`, `JEV_MIRROR_API_*`, and `LOCAL_API_URL` variables remain accepted while the deployment migrates. The `jev-latest` alias accepts only itself or a versioned `jev-x.y.z` response; an explicitly pinned Jev model must match exactly. `/api/tetris/config` exposes only readiness and display metadata; it never returns endpoint URLs or credentials.
|
| 112 |
|
|
|
|
| 114 |
|
| 115 |
Paid Cloud competitors additionally share one process-wide call budget across all clients. `TETRIS_CLOUD_CALLS_PER_WINDOW` defaults to 480 attempted calls, with a `TETRIS_CLOUD_CALL_WINDOW_SECONDS` default of 600 seconds. Admission reserves one call per possible step for each selected Cloud side (two for Cloud-vs-Cloud); a completed, failed, or cancelled race settles unused calls back to the budget. A race that cannot reserve its full possible Cloud cost receives HTTP 429 before any provider call. A local-vs-local race reserves no paid Cloud calls, including at the 2,000-step limit. Keep the default 120-step ceiling for public Cloud races unless the paid-provider budget is deliberately raised. The ledger is in memory and protects only this one Gateway process: deployments with multiple replicas or processes need a shared atomic budget store or an upstream account-level rate limit before claiming a fleet-wide cap.
|
| 116 |
|
| 117 |
+
`TETRIS_RACE_TIMEOUT_SECONDS` sets a 1–3,600 second whole-race deadline (default 900); raise it deliberately for slow long runs. `TETRIS_STREAM_GRACE_SECONDS` sets the 1–300 second reconnection lease after the last event stream disconnects (default 30); reconnecting resumes retained SSE events, and the browser can recover the latest board and terminal result from the race snapshot. An abandoned race is cancelled. `TETRIS_COMPLETED_TTL_SECONDS` controls completed-race retention from 30–3,600 seconds (default 300). Shutdown cancels and awaits every active race before closing the provider transport. A failed, timed-out, or otherwise incomplete race does not award a score winner. The speed percentage is available only when both sides complete the same step target; it compares the sums of complete server-observed adapter calls. For a private Decision runtime, that includes its inference POST, response attestation, and response validation. Cloud competitors include the provider POST, response read, and response validation. Local `inference_ms`, when returned, is retained per turn in the trace for diagnosis and never enters that comparison. Whole-run durations remain in the race summary as diagnostics; browser rendering delay does not enter the speed percentage.
|
| 118 |
|
| 119 |
`TETRIS_REQUEST_TIMEOUT_SECONDS` bounds each server-to-provider call to 0.1–300 seconds (default 45). The adapter sends optional bearer credentials only from server configuration, rejects redirects, ignores ambient proxy settings, bounds response bytes, and never returns provider URLs or credentials to the browser. Validate each configured endpoint with its selected canonical model and expected response identity before enabling that competitor; publishing a Space commit rebuilds the live Space.
|
Dockerfile
CHANGED
|
@@ -6,10 +6,10 @@ RUN useradd --create-home --uid 1000 studio
|
|
| 6 |
WORKDIR /app
|
| 7 |
COPY requirements.txt /app/requirements.txt
|
| 8 |
RUN pip install -r /app/requirements.txt
|
| 9 |
-
COPY --chown=studio:studio app.py systemone_api.py contract.py engine.py relay.py model_registry.py
|
| 10 |
COPY --chown=studio:studio static /app/static
|
| 11 |
USER studio
|
| 12 |
EXPOSE 7860
|
| 13 |
-
HEALTHCHECK --interval=30s --timeout=5s --start-period=20s CMD python -c "import os, urllib.request; path = '/api/ready' if os.getenv('DECISION_BACKEND') == '
|
| 14 |
# One process preserves the legacy in-memory queue/lease rollback mode.
|
| 15 |
CMD ["python", "-m", "uvicorn", "app:app", "--host", "0.0.0.0", "--port", "7860", "--workers", "1", "--no-access-log"]
|
|
|
|
| 6 |
WORKDIR /app
|
| 7 |
COPY requirements.txt /app/requirements.txt
|
| 8 |
RUN pip install -r /app/requirements.txt
|
| 9 |
+
COPY --chown=studio:studio app.py systemone_api.py contract.py engine.py relay.py model_registry.py direct_contract.py direct_gateway.py tetris_arena.py examples.json MODEL_RELEASES.json LICENSE /app/
|
| 10 |
COPY --chown=studio:studio static /app/static
|
| 11 |
USER studio
|
| 12 |
EXPOSE 7860
|
| 13 |
+
HEALTHCHECK --interval=30s --timeout=5s --start-period=20s CMD python -c "import os, urllib.request; path = '/api/ready' if os.getenv('DECISION_BACKEND') == 'direct' else '/api/status'; urllib.request.urlopen('http://127.0.0.1:7860' + path, timeout=3).close()"
|
| 14 |
# One process preserves the legacy in-memory queue/lease rollback mode.
|
| 15 |
CMD ["python", "-m", "uvicorn", "app:app", "--host", "0.0.0.0", "--port", "7860", "--workers", "1", "--no-access-log"]
|
README.md
CHANGED
|
@@ -40,4 +40,4 @@ Complete input includes context, instructions, candidates and special tokens. No
|
|
| 40 |
|
| 41 |
The editor runs on this CPU Space. Production direct mode maps each canonical model ID to its own private `vllm-sr decision serve` instance; the authenticated pull-queue deployment remains an explicit rollback mode. Switching clears old results, and an unavailable or misidentified model never silently falls back to another. Requests and results are not application-logged.
|
| 42 |
|
| 43 |
-
[Historical presentation metadata](MODEL_RELEASES.json) · [Request contract](SYSTEM_ONE_MAPPING.md) · [Deployment](DEPLOYMENT.md). Active revision and manifest pins come from the deployment registry and are checked against each live
|
|
|
|
| 40 |
|
| 41 |
The editor runs on this CPU Space. Production direct mode maps each canonical model ID to its own private `vllm-sr decision serve` instance; the authenticated pull-queue deployment remains an explicit rollback mode. Switching clears old results, and an unavailable or misidentified model never silently falls back to another. Requests and results are not application-logged.
|
| 42 |
|
| 43 |
+
[Historical presentation metadata](MODEL_RELEASES.json) · [Request contract](SYSTEM_ONE_MAPPING.md) · [Deployment](DEPLOYMENT.md). Active revision and manifest pins come from the deployment registry and are checked against each live Decision runtime service; the presentation file is not a serving allowlist. Studio code is MIT licensed. Models have their own Apache-2.0 project licenses and retained third-party notices. This repository contains no model weights or worker credentials.
|
SYSTEMONE_API.md
CHANGED
|
@@ -87,7 +87,7 @@ with TypeSafeClient(
|
|
| 87 |
print(result.scores["urgency"].score)
|
| 88 |
```
|
| 89 |
|
| 90 |
-
Select any exact Hugging Face repository ID returned as `id` by `/v1/models`. The model is required on public inference requests; short names, wire aliases, and case variants are rejected. The response repeats the selected canonical ID. The Studio editor uses a same-origin endpoint; in direct
|
| 91 |
|
| 92 |
Direct mode does not advertise a default model. The Studio editor selects a model locally and includes it explicitly in every `/api/evaluate` request; a missing model receives HTTP 422. Legacy rollback modes may retain their previous discovery default.
|
| 93 |
|
|
@@ -107,9 +107,9 @@ Direct mode does not advertise a default model. The Studio editor selects a mode
|
|
| 107 |
- Noul returns P(true); Choice returns its winning label and probabilities; Score returns the expected ordinal level and its distribution.
|
| 108 |
- Choice `confidence` is the top-two probability margin; Score `confidence` is the ordinal concentration around its expected level, normalized against uniform variance. These are versioned Decision-owned statistics, not calibrated probabilities of correctness or a claim of equivalence to TypeSafe's undisclosed calculation. The direct Gateway validates and forwards the runtime value unchanged; it never substitutes `max(probabilities)`. Public responses contain only `model`, `answers`, and `usage`.
|
| 109 |
- The hosted endpoint accepts one or more named questions without a fixed question-count cap. In direct mode, each model instance owns its physical microbatch setting. Choice accepts 2–255 options, Score accepts 2–10 ordered levels, and JSON requests are limited to 256 KiB. Expanded state/question input is limited to 16 MiB so large Cartesian workloads stay bounded. Provide explicit nonempty instructions. State plus question plus all candidates must fit the model limit; overflowing requests are rejected without truncation.
|
| 110 |
-
- In direct
|
| 111 |
-
- `POST /v1/systemone/batches` applies one shared `questions` map to ordered `states: [{id, state}, ...]`. It requires the same explicit canonical `model`, unique state IDs, at most 1,024 states, questions, and total decisions, a 2 MiB JSON body, and the same 16 MiB expanded-input and per-question token budgets. Its atomic response contains `model`, ordered `results: [{id, answers, usage}, ...]`, and aggregate `usage`. This Decision extension is separate from the SDK's single-state method.
|
| 112 |
|
| 113 |
See [SYSTEM_ONE_MAPPING.md](SYSTEM_ONE_MAPPING.md) for the complete adapter contract. The standard single-state call, all three answer types, and model discovery are validated with the official SDK.
|
| 114 |
|
| 115 |
-
In legacy pull-queue mode, completed asynchronous results are kept for up to 120 seconds, with the oldest completed records evicted when the 64-result cache fills. Direct
|
|
|
|
| 87 |
print(result.scores["urgency"].score)
|
| 88 |
```
|
| 89 |
|
| 90 |
+
Select any exact Hugging Face repository ID returned as `id` by `/v1/models`. The model is required on public inference requests; short names, wire aliases, and case variants are rejected. The response repeats the selected canonical ID. The Studio editor uses a same-origin endpoint; in direct serving mode it receives the same strict response envelope, while legacy modes retain their internal queue/diagnostic contracts.
|
| 91 |
|
| 92 |
Direct mode does not advertise a default model. The Studio editor selects a model locally and includes it explicitly in every `/api/evaluate` request; a missing model receives HTTP 422. Legacy rollback modes may retain their previous discovery default.
|
| 93 |
|
|
|
|
| 107 |
- Noul returns P(true); Choice returns its winning label and probabilities; Score returns the expected ordinal level and its distribution.
|
| 108 |
- Choice `confidence` is the top-two probability margin; Score `confidence` is the ordinal concentration around its expected level, normalized against uniform variance. These are versioned Decision-owned statistics, not calibrated probabilities of correctness or a claim of equivalence to TypeSafe's undisclosed calculation. The direct Gateway validates and forwards the runtime value unchanged; it never substitutes `max(probabilities)`. Public responses contain only `model`, `answers`, and `usage`.
|
| 109 |
- The hosted endpoint accepts one or more named questions without a fixed question-count cap. In direct mode, each model instance owns its physical microbatch setting. Choice accepts 2–255 options, Score accepts 2–10 ordered levels, and JSON requests are limited to 256 KiB. Expanded state/question input is limited to 16 MiB so large Cartesian workloads stay bounded. Provide explicit nonempty instructions. State plus question plus all candidates must fit the model limit; overflowing requests are rejected without truncation.
|
| 110 |
+
- In direct serving mode, the synchronous wait is configurable and each instance owns its own concurrency, queue, and physical-batch limits; `/v1/models` does not invent fixed values for them. A busy or timed-out service may return 529/503/504 (legacy queue mode may return 429). Completed requests do not consume admission capacity. Do not treat transport errors as predictions. SDK example retries are disabled to avoid duplicate inference.
|
| 111 |
+
- `POST /v1/systemone/batches` applies one shared `questions` map to ordered `states: [{id, state}, ...]`. It requires the same explicit canonical `model`, unique state IDs, at most 1,024 states, questions, and total decisions, a 2 MiB JSON body, and the same 16 MiB expanded-input and per-question token budgets. Its atomic response contains `model`, ordered `results: [{id, answers, usage}, ...]`, and aggregate `usage`. This Decision extension is separate from the SDK's single-state method.
|
| 112 |
|
| 113 |
See [SYSTEM_ONE_MAPPING.md](SYSTEM_ONE_MAPPING.md) for the complete adapter contract. The standard single-state call, all three answer types, and model discovery are validated with the official SDK.
|
| 114 |
|
| 115 |
+
In legacy pull-queue mode, completed asynchronous results are kept for up to 120 seconds, with the oldest completed records evicted when the 64-result cache fills. Direct serving mode does not create Gateway jobs or retain completion receipts. `GET /v1/models` publishes the service limits. The original Space origin remains available.
|
SYSTEM_ONE_MAPPING.md
CHANGED
|
@@ -16,13 +16,13 @@ Studio exposes a **SystemOne-format Decision service** with official Python SDK
|
|
| 16 |
|
| 17 |
The public single-state endpoint requires exactly `model`, `state`, and `questions`, with a 256 KiB JSON body. The separate batch endpoint requires exactly `model`, `states`, and `questions`, with a 2 MiB JSON body and at most 1,024 states, questions, and total decisions. Both accept 2–255 Choice options, 2–10 Score levels, and up to 16 MiB of expanded context/question input. The direct runtime uses a per-model configurable physical microbatch size rather than a fixed eight rows. The complete-input limit is 1,024 tokens for Kai/Lex and 16,384 for Eos/Sol/Nox/Lux **per question**, including state and all candidate descriptions. Unsupported fields and models are rejected.
|
| 18 |
|
| 19 |
-
Single-context responses preserve named `answers`, Choice `choice/probabilities`, Score `score/legend/probabilities`, and Noul `noul`. `usage.output_tokens=0` because these models do not generate text. The public inference envelopes contain only `model`, answers or results, and usage. In direct
|
| 20 |
|
| 21 |
The public inference responses supply versioned, type-aware Decision confidence: Choice uses the top-two probability margin, while Score uses ordinal concentration around the expected level relative to uniform variance. It is not calibrated correctness or a claim of equivalence to TypeSafe’s undisclosed confidence formula. The direct Gateway validates and preserves the runtime value rather than recalculating it as `max(probabilities)`. Noul returns its original probability of yes.
|
| 22 |
|
| 23 |
This mapping is independently testable with small metadata fixtures. It does not change the model's packing, native graph, weights, scoring arithmetic, or export format.
|
| 24 |
|
| 25 |
-
Sol and Nox use their published `question_row` and `predict_rows` path, the same numerical path used by `DecisionModel.decide`. Shipped temperatures, candidate order, true-probability orientation and ordinal Score arithmetic are preserved. Legacy Studio modes retain exact per-question input-token diagnostics; direct
|
| 26 |
|
| 27 |
## A shared question across contexts
|
| 28 |
|
|
@@ -46,6 +46,6 @@ Sol and Nox use their published `question_row` and `predict_rows` path, the same
|
|
| 46 |
}
|
| 47 |
```
|
| 48 |
|
| 49 |
-
Submit this payload to `POST /v1/systemone/batches`. Batch responses replace top-level `answers` with ordered `results: [{id, answers, usage}, ...]`. Context order and shared-question order are retained. Top-level `usage` sums all pairs. Legacy native/pull-queue Studio responses also report `contexts`, `questions`, and `decisions` in `timing`; direct
|
| 50 |
|
| 51 |
The worker admits every complete context/question input before its first forward. Any overflow rejects the whole request: no truncation or partial result. Admitted pairs are flattened into one predictor call with physical batches of up to eight. This is actual GPU batching, not browser-side request fanout or shared-state caching. The native numerical engine, weights, candidate order and published temperatures remain unchanged.
|
|
|
|
| 16 |
|
| 17 |
The public single-state endpoint requires exactly `model`, `state`, and `questions`, with a 256 KiB JSON body. The separate batch endpoint requires exactly `model`, `states`, and `questions`, with a 2 MiB JSON body and at most 1,024 states, questions, and total decisions. Both accept 2–255 Choice options, 2–10 Score levels, and up to 16 MiB of expanded context/question input. The direct runtime uses a per-model configurable physical microbatch size rather than a fixed eight rows. The complete-input limit is 1,024 tokens for Kai/Lex and 16,384 for Eos/Sol/Nox/Lux **per question**, including state and all candidate descriptions. Unsupported fields and models are rejected.
|
| 18 |
|
| 19 |
+
Single-context responses preserve named `answers`, Choice `choice/probabilities`, Score `score/legend/probabilities`, and Noul `noul`. `usage.output_tokens=0` because these models do not generate text. The public inference envelopes contain only `model`, answers or results, and usage. In direct serving mode, Studio's `/api/evaluate` route returns that same strict envelope; legacy native and pull-queue modes retain their diagnostic response fields for rollback.
|
| 20 |
|
| 21 |
The public inference responses supply versioned, type-aware Decision confidence: Choice uses the top-two probability margin, while Score uses ordinal concentration around the expected level relative to uniform variance. It is not calibrated correctness or a claim of equivalence to TypeSafe’s undisclosed confidence formula. The direct Gateway validates and preserves the runtime value rather than recalculating it as `max(probabilities)`. Noul returns its original probability of yes.
|
| 22 |
|
| 23 |
This mapping is independently testable with small metadata fixtures. It does not change the model's packing, native graph, weights, scoring arithmetic, or export format.
|
| 24 |
|
| 25 |
+
Sol and Nox use their published `question_row` and `predict_rows` path, the same numerical path used by `DecisionModel.decide`. Shipped temperatures, candidate order, true-probability orientation and ordinal Score arithmetic are preserved. Legacy Studio modes retain exact per-question input-token diagnostics; direct serving mode returns only strict aggregate usage. Neither path reinterprets a Decision statistic as TypeSafe confidence.
|
| 26 |
|
| 27 |
## A shared question across contexts
|
| 28 |
|
|
|
|
| 46 |
}
|
| 47 |
```
|
| 48 |
|
| 49 |
+
Submit this payload to `POST /v1/systemone/batches`. Batch responses replace top-level `answers` with ordered `results: [{id, answers, usage}, ...]`. Context order and shared-question order are retained. Top-level `usage` sums all pairs. Legacy native/pull-queue Studio responses also report `contexts`, `questions`, and `decisions` in `timing`; direct responses stay strict. A row can expand to show the original probabilities. Score is an expected ordinal index; Noul is P(Yes), not an independent confidence estimate.
|
| 50 |
|
| 51 |
The worker admits every complete context/question input before its first forward. Any overflow rejects the whole request: no truncation or partial result. Admitted pairs are flattened into one predictor call with physical batches of up to eight. This is actual GPU batching, not browser-side request fanout or shared-state caching. The native numerical engine, weights, candidate order and published temperatures remain unchanged.
|
app.py
CHANGED
|
@@ -14,7 +14,7 @@ from fastapi.staticfiles import StaticFiles
|
|
| 14 |
from starlette.concurrency import run_in_threadpool
|
| 15 |
|
| 16 |
from contract import MODEL, to_records
|
| 17 |
-
from
|
| 18 |
from engine import Busy, Engine, Unavailable
|
| 19 |
from relay import Relay, RelayError
|
| 20 |
from systemone_api import (
|
|
@@ -112,13 +112,13 @@ def create_app(
|
|
| 112 |
relay=None,
|
| 113 |
registry=None,
|
| 114 |
relays=None,
|
| 115 |
-
|
| 116 |
tetris_adapter=None,
|
| 117 |
tetris_manager=None,
|
| 118 |
):
|
| 119 |
mode = mode or os.getenv("DECISION_BACKEND", "native")
|
| 120 |
-
if mode not in {"native", "pull_queue", "
|
| 121 |
-
raise ValueError("DECISION_BACKEND must be native, pull_queue, or
|
| 122 |
from model_registry import model_registry
|
| 123 |
if relay is not None and registry is None:
|
| 124 |
registry = [{'id': MODEL, 'label': 'Kai', 'version': '1.0', 'manifest_sha256': relay.manifest}]
|
|
@@ -126,16 +126,16 @@ def create_app(
|
|
| 126 |
default_model = next(iter(registry))
|
| 127 |
canonical_models = tuple(item['repo_id'] for item in registry.values())
|
| 128 |
expected_artifacts = {item['repo_id']: item for item in registry.values()}
|
| 129 |
-
|
| 130 |
-
if mode == "
|
| 131 |
if any('revision' not in item for item in registry.values()):
|
| 132 |
-
raise ValueError("Direct
|
| 133 |
-
|
| 134 |
-
if getattr(
|
| 135 |
raise ValueError(
|
| 136 |
-
"The
|
| 137 |
)
|
| 138 |
-
gateway_pins = getattr(
|
| 139 |
if gateway_pins is not None and any(
|
| 140 |
gateway_pins.get(model) != {
|
| 141 |
"revision": item["revision"],
|
|
@@ -145,9 +145,9 @@ def create_app(
|
|
| 145 |
}
|
| 146 |
for model, item in expected_artifacts.items()
|
| 147 |
):
|
| 148 |
-
raise ValueError("The
|
| 149 |
-
elif
|
| 150 |
-
raise ValueError("A
|
| 151 |
trusted_proxies = trusted_proxy_networks(
|
| 152 |
os.getenv("TETRIS_TRUSTED_PROXY_CIDRS", "")
|
| 153 |
)
|
|
@@ -156,9 +156,9 @@ def create_app(
|
|
| 156 |
if tetris_manager is None:
|
| 157 |
tetris_adapter = tetris_adapter or HTTPDecisionAdapter.from_environment(
|
| 158 |
local_endpoints=(
|
| 159 |
-
|
| 160 |
),
|
| 161 |
-
direct_gateway=
|
| 162 |
)
|
| 163 |
tetris_manager = RaceManager.from_environment(tetris_adapter)
|
| 164 |
|
|
@@ -173,8 +173,8 @@ def create_app(
|
|
| 173 |
if owns_tetris_adapter:
|
| 174 |
await tetris_adapter.close()
|
| 175 |
finally:
|
| 176 |
-
if
|
| 177 |
-
await
|
| 178 |
|
| 179 |
api = FastAPI(
|
| 180 |
title="Decision Studio",
|
|
@@ -200,7 +200,7 @@ def create_app(
|
|
| 200 |
api.state.relays = relays
|
| 201 |
api.state.relay = relays.get(MODEL) # backward-compatible default for local tooling
|
| 202 |
api.state.registry = registry
|
| 203 |
-
api.state.
|
| 204 |
|
| 205 |
api.state.tetris = tetris_manager
|
| 206 |
|
|
@@ -210,10 +210,10 @@ def create_app(
|
|
| 210 |
return relays.get(model)
|
| 211 |
|
| 212 |
async def direct_readiness(models):
|
| 213 |
-
"""Probe distinct configured
|
| 214 |
selected = tuple(dict.fromkeys(models))
|
| 215 |
observations = await asyncio.gather(
|
| 216 |
-
*(
|
| 217 |
return_exceptions=True,
|
| 218 |
)
|
| 219 |
return {
|
|
@@ -222,7 +222,7 @@ def create_app(
|
|
| 222 |
}
|
| 223 |
|
| 224 |
async def tetris_direct_readiness(competitor_ids=None):
|
| 225 |
-
if mode != "
|
| 226 |
return {}
|
| 227 |
competitors = (
|
| 228 |
api.state.tetris.catalog()
|
|
@@ -247,7 +247,7 @@ def create_app(
|
|
| 247 |
limit=(
|
| 248 |
2 * 1024 * 1024
|
| 249 |
if public_contract == 'batch'
|
| 250 |
-
or (mode == '
|
| 251 |
else 256 * 1024
|
| 252 |
),
|
| 253 |
)
|
|
@@ -260,7 +260,7 @@ def create_app(
|
|
| 260 |
resolve_canonical_model(requested_model, registry)
|
| 261 |
if public_contract is not None
|
| 262 |
else resolve_model(
|
| 263 |
-
requested_model if mode == "
|
| 264 |
registry,
|
| 265 |
)
|
| 266 |
)
|
|
@@ -275,7 +275,7 @@ def create_app(
|
|
| 275 |
records = to_records(payload, model=model)
|
| 276 |
except (ValueError, TypeError, UnicodeError, RecursionError) as exc:
|
| 277 |
if (
|
| 278 |
-
mode == "
|
| 279 |
and "Expanded context/question input exceeds" in str(exc)
|
| 280 |
):
|
| 281 |
raise HTTPException(413, str(exc)) from None
|
|
@@ -287,8 +287,8 @@ def create_app(
|
|
| 287 |
return JSONResponse({"detail": exc.message}, status_code=exc.code,
|
| 288 |
headers={"Retry-After": "2"} if exc.code in {429, 503} else None)
|
| 289 |
|
| 290 |
-
@api.exception_handler(
|
| 291 |
-
async def
|
| 292 |
headers = {"Retry-After": exc.retry_after} if exc.retry_after else None
|
| 293 |
return JSONResponse({"detail": exc.message}, status_code=exc.code, headers=headers)
|
| 294 |
|
|
@@ -296,8 +296,8 @@ def create_app(
|
|
| 296 |
async def status(model: str = default_model):
|
| 297 |
model = resolve_model(model, registry)
|
| 298 |
selected = select(model)
|
| 299 |
-
if mode == "
|
| 300 |
-
observed = await
|
| 301 |
loaded = observed['loaded']
|
| 302 |
artifact = observed['artifact'] or {}
|
| 303 |
return {
|
|
@@ -312,7 +312,7 @@ def create_app(
|
|
| 312 |
"loaded_content_sha256": artifact.get('content_sha256'),
|
| 313 |
"complete_input_tokens": registry[model]['complete_input_tokens'],
|
| 314 |
"context_batch": loaded,
|
| 315 |
-
"backend": "
|
| 316 |
"live_inference": loaded,
|
| 317 |
"running": False,
|
| 318 |
"queued": 0,
|
|
@@ -322,7 +322,7 @@ def create_app(
|
|
| 322 |
@api.get("/api/ready")
|
| 323 |
async def ready():
|
| 324 |
"""Return HTTP 503 when any configured model is unavailable."""
|
| 325 |
-
if mode == '
|
| 326 |
loaded = await direct_readiness(canonical_models)
|
| 327 |
elif relays:
|
| 328 |
observations = await asyncio.gather(
|
|
@@ -366,7 +366,7 @@ def create_app(
|
|
| 366 |
@api.post("/api/tetris/races", status_code=201)
|
| 367 |
async def create_tetris_race(request: Request):
|
| 368 |
payload = await read_json(request)
|
| 369 |
-
if mode == '
|
| 370 |
selected_ids = {
|
| 371 |
value for side in ('left', 'right')
|
| 372 |
if isinstance(value := payload.get(side), str)
|
|
@@ -443,8 +443,8 @@ def create_app(
|
|
| 443 |
@api.get("/v1/models")
|
| 444 |
def models():
|
| 445 |
synchronous_wait = (
|
| 446 |
-
getattr(
|
| 447 |
-
if mode == "
|
| 448 |
else 45
|
| 449 |
)
|
| 450 |
limits = {
|
|
@@ -452,17 +452,17 @@ def create_app(
|
|
| 452 |
"expanded_input_bytes": 16 * 1024 * 1024,
|
| 453 |
"questions": None,
|
| 454 |
"contexts": None,
|
| 455 |
-
#
|
| 456 |
# retired Gateway's fixed eight-request/eight-row profile.
|
| 457 |
-
"active_requests_per_model": None if mode == "
|
| 458 |
-
"gpu_microbatch": None if mode == "
|
| 459 |
"synchronous_wait_seconds": synchronous_wait,
|
| 460 |
}
|
| 461 |
result = {
|
| 462 |
"models": public_models(registry),
|
| 463 |
"limits": limits,
|
| 464 |
}
|
| 465 |
-
if mode != "
|
| 466 |
result["default"] = registry[default_model]['repo_id']
|
| 467 |
return result
|
| 468 |
|
|
@@ -482,8 +482,8 @@ def create_app(
|
|
| 482 |
payload, records, selected, public_model, canonical_payload, _ = await input_for(
|
| 483 |
request, public_contract=public_contract
|
| 484 |
)
|
| 485 |
-
if mode == "
|
| 486 |
-
return await
|
| 487 |
canonical_payload,
|
| 488 |
batch='states' in canonical_payload,
|
| 489 |
)
|
|
|
|
| 14 |
from starlette.concurrency import run_in_threadpool
|
| 15 |
|
| 16 |
from contract import MODEL, to_records
|
| 17 |
+
from direct_gateway import DirectGateway, DirectGatewayError
|
| 18 |
from engine import Busy, Engine, Unavailable
|
| 19 |
from relay import Relay, RelayError
|
| 20 |
from systemone_api import (
|
|
|
|
| 112 |
relay=None,
|
| 113 |
registry=None,
|
| 114 |
relays=None,
|
| 115 |
+
direct_gateway=None,
|
| 116 |
tetris_adapter=None,
|
| 117 |
tetris_manager=None,
|
| 118 |
):
|
| 119 |
mode = mode or os.getenv("DECISION_BACKEND", "native")
|
| 120 |
+
if mode not in {"native", "pull_queue", "direct"}:
|
| 121 |
+
raise ValueError("DECISION_BACKEND must be native, pull_queue, or direct")
|
| 122 |
from model_registry import model_registry
|
| 123 |
if relay is not None and registry is None:
|
| 124 |
registry = [{'id': MODEL, 'label': 'Kai', 'version': '1.0', 'manifest_sha256': relay.manifest}]
|
|
|
|
| 126 |
default_model = next(iter(registry))
|
| 127 |
canonical_models = tuple(item['repo_id'] for item in registry.values())
|
| 128 |
expected_artifacts = {item['repo_id']: item for item in registry.values()}
|
| 129 |
+
owns_direct_gateway = mode == "direct" and direct_gateway is None
|
| 130 |
+
if mode == "direct":
|
| 131 |
if any('revision' not in item for item in registry.values()):
|
| 132 |
+
raise ValueError("Direct serving mode requires a pinned Hub revision for every model")
|
| 133 |
+
direct_gateway = direct_gateway or DirectGateway.from_environment(expected_artifacts)
|
| 134 |
+
if getattr(direct_gateway, "models", None) != frozenset(canonical_models):
|
| 135 |
raise ValueError(
|
| 136 |
+
"The direct Gateway must match every configured canonical model"
|
| 137 |
)
|
| 138 |
+
gateway_pins = getattr(direct_gateway, "expected_artifacts", None)
|
| 139 |
if gateway_pins is not None and any(
|
| 140 |
gateway_pins.get(model) != {
|
| 141 |
"revision": item["revision"],
|
|
|
|
| 145 |
}
|
| 146 |
for model, item in expected_artifacts.items()
|
| 147 |
):
|
| 148 |
+
raise ValueError("The direct Gateway artifact pins must match the registry")
|
| 149 |
+
elif direct_gateway is not None:
|
| 150 |
+
raise ValueError("A direct Gateway may only be supplied in direct mode")
|
| 151 |
trusted_proxies = trusted_proxy_networks(
|
| 152 |
os.getenv("TETRIS_TRUSTED_PROXY_CIDRS", "")
|
| 153 |
)
|
|
|
|
| 156 |
if tetris_manager is None:
|
| 157 |
tetris_adapter = tetris_adapter or HTTPDecisionAdapter.from_environment(
|
| 158 |
local_endpoints=(
|
| 159 |
+
direct_gateway.local_systemone_endpoints() if mode == "direct" else None
|
| 160 |
),
|
| 161 |
+
direct_gateway=direct_gateway if mode == "direct" else None,
|
| 162 |
)
|
| 163 |
tetris_manager = RaceManager.from_environment(tetris_adapter)
|
| 164 |
|
|
|
|
| 173 |
if owns_tetris_adapter:
|
| 174 |
await tetris_adapter.close()
|
| 175 |
finally:
|
| 176 |
+
if owns_direct_gateway:
|
| 177 |
+
await direct_gateway.aclose()
|
| 178 |
|
| 179 |
api = FastAPI(
|
| 180 |
title="Decision Studio",
|
|
|
|
| 200 |
api.state.relays = relays
|
| 201 |
api.state.relay = relays.get(MODEL) # backward-compatible default for local tooling
|
| 202 |
api.state.registry = registry
|
| 203 |
+
api.state.direct = direct_gateway
|
| 204 |
|
| 205 |
api.state.tetris = tetris_manager
|
| 206 |
|
|
|
|
| 210 |
return relays.get(model)
|
| 211 |
|
| 212 |
async def direct_readiness(models):
|
| 213 |
+
"""Probe distinct configured Decision models without exposing their origins."""
|
| 214 |
selected = tuple(dict.fromkeys(models))
|
| 215 |
observations = await asyncio.gather(
|
| 216 |
+
*(direct_gateway.probe(model) for model in selected),
|
| 217 |
return_exceptions=True,
|
| 218 |
)
|
| 219 |
return {
|
|
|
|
| 222 |
}
|
| 223 |
|
| 224 |
async def tetris_direct_readiness(competitor_ids=None):
|
| 225 |
+
if mode != "direct":
|
| 226 |
return {}
|
| 227 |
competitors = (
|
| 228 |
api.state.tetris.catalog()
|
|
|
|
| 247 |
limit=(
|
| 248 |
2 * 1024 * 1024
|
| 249 |
if public_contract == 'batch'
|
| 250 |
+
or (mode == 'direct' and public_contract is None)
|
| 251 |
else 256 * 1024
|
| 252 |
),
|
| 253 |
)
|
|
|
|
| 260 |
resolve_canonical_model(requested_model, registry)
|
| 261 |
if public_contract is not None
|
| 262 |
else resolve_model(
|
| 263 |
+
requested_model if mode == "direct" else requested_model or default_model,
|
| 264 |
registry,
|
| 265 |
)
|
| 266 |
)
|
|
|
|
| 275 |
records = to_records(payload, model=model)
|
| 276 |
except (ValueError, TypeError, UnicodeError, RecursionError) as exc:
|
| 277 |
if (
|
| 278 |
+
mode == "direct"
|
| 279 |
and "Expanded context/question input exceeds" in str(exc)
|
| 280 |
):
|
| 281 |
raise HTTPException(413, str(exc)) from None
|
|
|
|
| 287 |
return JSONResponse({"detail": exc.message}, status_code=exc.code,
|
| 288 |
headers={"Retry-After": "2"} if exc.code in {429, 503} else None)
|
| 289 |
|
| 290 |
+
@api.exception_handler(DirectGatewayError)
|
| 291 |
+
async def direct_error(request, exc):
|
| 292 |
headers = {"Retry-After": exc.retry_after} if exc.retry_after else None
|
| 293 |
return JSONResponse({"detail": exc.message}, status_code=exc.code, headers=headers)
|
| 294 |
|
|
|
|
| 296 |
async def status(model: str = default_model):
|
| 297 |
model = resolve_model(model, registry)
|
| 298 |
selected = select(model)
|
| 299 |
+
if mode == "direct":
|
| 300 |
+
observed = await direct_gateway.probe(registry[model]['repo_id'])
|
| 301 |
loaded = observed['loaded']
|
| 302 |
artifact = observed['artifact'] or {}
|
| 303 |
return {
|
|
|
|
| 312 |
"loaded_content_sha256": artifact.get('content_sha256'),
|
| 313 |
"complete_input_tokens": registry[model]['complete_input_tokens'],
|
| 314 |
"context_batch": loaded,
|
| 315 |
+
"backend": "direct",
|
| 316 |
"live_inference": loaded,
|
| 317 |
"running": False,
|
| 318 |
"queued": 0,
|
|
|
|
| 322 |
@api.get("/api/ready")
|
| 323 |
async def ready():
|
| 324 |
"""Return HTTP 503 when any configured model is unavailable."""
|
| 325 |
+
if mode == 'direct':
|
| 326 |
loaded = await direct_readiness(canonical_models)
|
| 327 |
elif relays:
|
| 328 |
observations = await asyncio.gather(
|
|
|
|
| 366 |
@api.post("/api/tetris/races", status_code=201)
|
| 367 |
async def create_tetris_race(request: Request):
|
| 368 |
payload = await read_json(request)
|
| 369 |
+
if mode == 'direct' and isinstance(payload, dict):
|
| 370 |
selected_ids = {
|
| 371 |
value for side in ('left', 'right')
|
| 372 |
if isinstance(value := payload.get(side), str)
|
|
|
|
| 443 |
@api.get("/v1/models")
|
| 444 |
def models():
|
| 445 |
synchronous_wait = (
|
| 446 |
+
getattr(direct_gateway, "timeout_seconds", 45)
|
| 447 |
+
if mode == "direct"
|
| 448 |
else 45
|
| 449 |
)
|
| 450 |
limits = {
|
|
|
|
| 452 |
"expanded_input_bytes": 16 * 1024 * 1024,
|
| 453 |
"questions": None,
|
| 454 |
"contexts": None,
|
| 455 |
+
# direct capacity belongs to each live model instance, not the
|
| 456 |
# retired Gateway's fixed eight-request/eight-row profile.
|
| 457 |
+
"active_requests_per_model": None if mode == "direct" else 8,
|
| 458 |
+
"gpu_microbatch": None if mode == "direct" else 8,
|
| 459 |
"synchronous_wait_seconds": synchronous_wait,
|
| 460 |
}
|
| 461 |
result = {
|
| 462 |
"models": public_models(registry),
|
| 463 |
"limits": limits,
|
| 464 |
}
|
| 465 |
+
if mode != "direct":
|
| 466 |
result["default"] = registry[default_model]['repo_id']
|
| 467 |
return result
|
| 468 |
|
|
|
|
| 482 |
payload, records, selected, public_model, canonical_payload, _ = await input_for(
|
| 483 |
request, public_contract=public_contract
|
| 484 |
)
|
| 485 |
+
if mode == "direct":
|
| 486 |
+
return await direct_gateway.evaluate(
|
| 487 |
canonical_payload,
|
| 488 |
batch='states' in canonical_payload,
|
| 489 |
)
|
drun_contract.py → direct_contract.py
RENAMED
|
@@ -163,7 +163,7 @@ def _answers_for_questions(questions: object, answers: object, path: str) -> Non
|
|
| 163 |
def validate_response_for_request(
|
| 164 |
request: Mapping[str, object], response: object, *, batch: bool
|
| 165 |
) -> None:
|
| 166 |
-
"""Validate a strict
|
| 167 |
|
| 168 |
expected_top = {"model", "results", "usage"} if batch else {
|
| 169 |
"model",
|
|
|
|
| 163 |
def validate_response_for_request(
|
| 164 |
request: Mapping[str, object], response: object, *, batch: bool
|
| 165 |
) -> None:
|
| 166 |
+
"""Validate a strict direct response without changing any returned statistic."""
|
| 167 |
|
| 168 |
expected_top = {"model", "results", "usage"} if batch else {
|
| 169 |
"model",
|
drun_gateway.py → direct_gateway.py
RENAMED
|
@@ -1,4 +1,4 @@
|
|
| 1 |
-
"""Private
|
| 2 |
|
| 3 |
The browser never receives backend locations. Every inference response is
|
| 4 |
validated against the originating strict request before it crosses the public
|
|
@@ -19,7 +19,7 @@ from urllib.parse import urlsplit
|
|
| 19 |
|
| 20 |
import httpx
|
| 21 |
|
| 22 |
-
from
|
| 23 |
|
| 24 |
SINGLE_PATH = "/v1/systemone"
|
| 25 |
BATCH_PATH = "/v1/systemone/batches"
|
|
@@ -55,7 +55,7 @@ def _hex_digest(value: object, length: int) -> bool:
|
|
| 55 |
)
|
| 56 |
|
| 57 |
|
| 58 |
-
class
|
| 59 |
"""A sanitized direct-backend failure safe to return to a caller."""
|
| 60 |
|
| 61 |
def __init__(self, code: int, message: str, *, retry_after: str | None = None):
|
|
@@ -88,12 +88,12 @@ def _private_origin(value: object) -> str:
|
|
| 88 |
"""Accept an explicit private/loopback IP origin, never a URL path or DNS name."""
|
| 89 |
|
| 90 |
if not isinstance(value, str) or not value or len(value) > 2048:
|
| 91 |
-
raise ValueError("Every
|
| 92 |
try:
|
| 93 |
parsed = urlsplit(value)
|
| 94 |
port = parsed.port
|
| 95 |
except ValueError:
|
| 96 |
-
raise ValueError("Every
|
| 97 |
if (
|
| 98 |
parsed.scheme not in {"http", "https"}
|
| 99 |
or parsed.username is not None
|
|
@@ -106,15 +106,15 @@ def _private_origin(value: object) -> str:
|
|
| 106 |
or not 1 <= port <= 65535
|
| 107 |
):
|
| 108 |
raise ValueError(
|
| 109 |
-
"Every
|
| 110 |
)
|
| 111 |
if "%" in parsed.hostname:
|
| 112 |
-
raise ValueError("Scoped IPv6
|
| 113 |
try:
|
| 114 |
address = ipaddress.ip_address(parsed.hostname)
|
| 115 |
except ValueError:
|
| 116 |
raise ValueError(
|
| 117 |
-
"Every
|
| 118 |
) from None
|
| 119 |
if isinstance(address, ipaddress.IPv6Address) and address.ipv4_mapped:
|
| 120 |
address = address.ipv4_mapped
|
|
@@ -125,42 +125,42 @@ def _private_origin(value: object) -> str:
|
|
| 125 |
for network in PRIVATE_NETWORKS
|
| 126 |
)
|
| 127 |
):
|
| 128 |
-
raise ValueError("Every
|
| 129 |
host = f"[{address.compressed}]" if address.version == 6 else address.compressed
|
| 130 |
return f"{parsed.scheme}://{host}:{port}"
|
| 131 |
|
| 132 |
|
| 133 |
def endpoints_from_environment(expected_models: Collection[str]) -> dict[str, str]:
|
| 134 |
-
raw = os.getenv("
|
| 135 |
if not raw:
|
| 136 |
-
raise ValueError("
|
| 137 |
if len(raw.encode("utf-8")) > MAX_ENDPOINT_CONFIG_BYTES:
|
| 138 |
-
raise ValueError("
|
| 139 |
try:
|
| 140 |
value = _loads_strict(raw)
|
| 141 |
except (TypeError, ValueError, UnicodeError, RecursionError) as exc:
|
| 142 |
-
raise ValueError("
|
| 143 |
if not isinstance(value, dict) or set(value) != set(expected_models):
|
| 144 |
raise ValueError(
|
| 145 |
-
"
|
| 146 |
)
|
| 147 |
return {model: _private_origin(endpoint) for model, endpoint in value.items()}
|
| 148 |
|
| 149 |
|
| 150 |
def timeout_from_environment() -> float:
|
| 151 |
-
raw = os.getenv("
|
| 152 |
try:
|
| 153 |
timeout = float(raw)
|
| 154 |
except ValueError as exc:
|
| 155 |
-
raise ValueError("
|
| 156 |
if not math.isfinite(timeout) or not MIN_TIMEOUT_SECONDS <= timeout <= MAX_TIMEOUT_SECONDS:
|
| 157 |
raise ValueError(
|
| 158 |
-
"
|
| 159 |
)
|
| 160 |
return timeout
|
| 161 |
|
| 162 |
|
| 163 |
-
class
|
| 164 |
"""Keep-alive direct client with fail-closed model and response validation."""
|
| 165 |
|
| 166 |
def __init__(
|
|
@@ -179,14 +179,14 @@ class DrunGateway:
|
|
| 179 |
or any(not isinstance(model, str) or not model for model in expected)
|
| 180 |
or len(set(expected)) != len(expected)
|
| 181 |
):
|
| 182 |
-
raise ValueError("Configure a nonempty set of canonical
|
| 183 |
if set(endpoints) != set(expected):
|
| 184 |
-
raise ValueError("
|
| 185 |
pinned = {}
|
| 186 |
for model in expected:
|
| 187 |
item = expected_artifacts[model]
|
| 188 |
if not isinstance(item, Mapping):
|
| 189 |
-
raise ValueError("Pin each
|
| 190 |
revision = item.get("revision")
|
| 191 |
manifest = item.get("manifest_sha256")
|
| 192 |
content = item.get("content_sha256")
|
|
@@ -195,7 +195,7 @@ class DrunGateway:
|
|
| 195 |
or not _hex_digest(manifest, 64)
|
| 196 |
or (content is not None and not _hex_digest(content, 64))
|
| 197 |
):
|
| 198 |
-
raise ValueError("Pin each
|
| 199 |
pinned[model] = MappingProxyType(
|
| 200 |
{
|
| 201 |
"revision": revision,
|
|
@@ -214,15 +214,15 @@ class DrunGateway:
|
|
| 214 |
except OverflowError:
|
| 215 |
timeout_valid = False
|
| 216 |
if not timeout_valid:
|
| 217 |
-
raise ValueError("
|
| 218 |
if (
|
| 219 |
type(max_response_bytes) is not int
|
| 220 |
or not 1 <= max_response_bytes <= MAX_INFERENCE_RESPONSE_BYTES
|
| 221 |
):
|
| 222 |
-
raise ValueError("Invalid
|
| 223 |
normalized = {model: _private_origin(endpoints[model]) for model in expected}
|
| 224 |
if len(set(normalized.values())) != len(normalized):
|
| 225 |
-
raise ValueError("Each configured Decision model requires its own
|
| 226 |
self._endpoints = MappingProxyType(normalized)
|
| 227 |
self.expected_artifacts = MappingProxyType(pinned)
|
| 228 |
self.models = frozenset(expected)
|
|
@@ -239,7 +239,7 @@ class DrunGateway:
|
|
| 239 |
@classmethod
|
| 240 |
def from_environment(
|
| 241 |
cls, expected_artifacts: Mapping[str, Mapping[str, str]]
|
| 242 |
-
) ->
|
| 243 |
return cls(
|
| 244 |
endpoints_from_environment(expected_artifacts),
|
| 245 |
expected_artifacts=expected_artifacts,
|
|
@@ -268,9 +268,9 @@ class DrunGateway:
|
|
| 268 |
|
| 269 |
model = payload.get("model") if isinstance(payload, Mapping) else None
|
| 270 |
if not isinstance(model, str) or model not in self.models:
|
| 271 |
-
raise
|
| 272 |
if batch != ("states" in payload) or ("state" in payload and "states" in payload):
|
| 273 |
-
raise
|
| 274 |
try:
|
| 275 |
encoded = json.dumps(
|
| 276 |
payload,
|
|
@@ -279,10 +279,10 @@ class DrunGateway:
|
|
| 279 |
allow_nan=False,
|
| 280 |
).encode("utf-8")
|
| 281 |
except (TypeError, ValueError, UnicodeError, RecursionError) as exc:
|
| 282 |
-
raise
|
| 283 |
request_limit = MAX_BATCH_REQUEST_BYTES if batch else MAX_SINGLE_REQUEST_BYTES
|
| 284 |
if len(encoded) > request_limit:
|
| 285 |
-
raise
|
| 286 |
started = self._clock()
|
| 287 |
response = await self._request_json(
|
| 288 |
"POST",
|
|
@@ -296,7 +296,7 @@ class DrunGateway:
|
|
| 296 |
try:
|
| 297 |
validate_response_for_request(payload, response, batch=batch)
|
| 298 |
except (TypeError, ValueError, KeyError, OverflowError, RecursionError) as exc:
|
| 299 |
-
raise
|
| 300 |
502, "The selected Decision runtime returned an invalid response"
|
| 301 |
) from exc
|
| 302 |
return response, request_ms
|
|
@@ -327,7 +327,7 @@ class DrunGateway:
|
|
| 327 |
and self._artifact_matches(model, artifact)
|
| 328 |
)
|
| 329 |
return {"loaded": bool(loaded), "artifact": artifact}
|
| 330 |
-
except
|
| 331 |
return {"loaded": False, "artifact": None}
|
| 332 |
|
| 333 |
@staticmethod
|
|
@@ -390,7 +390,7 @@ class DrunGateway:
|
|
| 390 |
):
|
| 391 |
endpoint = self._endpoints.get(model)
|
| 392 |
if endpoint is None:
|
| 393 |
-
raise
|
| 394 |
request_timeout = timeout_seconds or self.timeout_seconds
|
| 395 |
try:
|
| 396 |
headers = {
|
|
@@ -409,7 +409,7 @@ class DrunGateway:
|
|
| 409 |
timeout=request_timeout,
|
| 410 |
) as response:
|
| 411 |
if 300 <= response.status_code < 400:
|
| 412 |
-
raise
|
| 413 |
502,
|
| 414 |
"The selected Decision runtime returned an unsupported redirect",
|
| 415 |
)
|
|
@@ -418,20 +418,20 @@ class DrunGateway:
|
|
| 418 |
if attest_model is not None and not self._artifact_matches(
|
| 419 |
attest_model, self._parse_response_artifact(response.headers)
|
| 420 |
):
|
| 421 |
-
raise
|
| 422 |
503, "The selected Decision runtime artifact is unavailable"
|
| 423 |
)
|
| 424 |
return await self._read_json_response(response, limit)
|
| 425 |
except asyncio.CancelledError:
|
| 426 |
raise
|
| 427 |
-
except
|
| 428 |
raise
|
| 429 |
except (TimeoutError, httpx.TimeoutException) as exc:
|
| 430 |
-
raise
|
| 431 |
504, "The selected Decision runtime did not respond in time"
|
| 432 |
) from exc
|
| 433 |
except httpx.HTTPError as exc:
|
| 434 |
-
raise
|
| 435 |
503, "The selected Decision runtime is unavailable"
|
| 436 |
) from exc
|
| 437 |
|
|
@@ -439,26 +439,26 @@ class DrunGateway:
|
|
| 439 |
def _raise_status(response: httpx.Response) -> None:
|
| 440 |
status = response.status_code
|
| 441 |
if status == 413:
|
| 442 |
-
raise
|
| 443 |
413, "Decision input exceeds the selected model limit"
|
| 444 |
)
|
| 445 |
if status in {429, 529}:
|
| 446 |
-
raise
|
| 447 |
status,
|
| 448 |
"The selected Decision runtime is temporarily overloaded",
|
| 449 |
retry_after=_safe_retry_after(response.headers.get("retry-after")),
|
| 450 |
)
|
| 451 |
if status == 503:
|
| 452 |
-
raise
|
| 453 |
503,
|
| 454 |
"The selected Decision runtime is unavailable",
|
| 455 |
retry_after=_safe_retry_after(response.headers.get("retry-after")),
|
| 456 |
)
|
| 457 |
if status == 504:
|
| 458 |
-
raise
|
| 459 |
504, "The selected Decision runtime did not respond in time"
|
| 460 |
)
|
| 461 |
-
raise
|
| 462 |
502, "The selected Decision runtime rejected a validated request"
|
| 463 |
)
|
| 464 |
|
|
@@ -468,12 +468,12 @@ class DrunGateway:
|
|
| 468 |
response.headers.get("content-type", "").split(";", 1)[0].strip().lower()
|
| 469 |
)
|
| 470 |
if content_type != "application/json":
|
| 471 |
-
raise
|
| 472 |
502, "The selected Decision runtime returned an invalid content type"
|
| 473 |
)
|
| 474 |
encoding = response.headers.get("content-encoding", "identity").strip().lower()
|
| 475 |
if encoding not in {"", "identity"}:
|
| 476 |
-
raise
|
| 477 |
502, "The selected Decision runtime returned an unsupported encoding"
|
| 478 |
)
|
| 479 |
declared = response.headers.get("content-length")
|
|
@@ -481,11 +481,11 @@ class DrunGateway:
|
|
| 481 |
try:
|
| 482 |
declared_size = int(declared)
|
| 483 |
except ValueError as exc:
|
| 484 |
-
raise
|
| 485 |
502, "The selected Decision runtime returned an invalid response"
|
| 486 |
) from exc
|
| 487 |
if declared_size < 0 or declared_size > limit:
|
| 488 |
-
raise
|
| 489 |
502, "The selected Decision runtime returned an oversized response"
|
| 490 |
)
|
| 491 |
body = bytearray()
|
|
@@ -494,14 +494,14 @@ class DrunGateway:
|
|
| 494 |
# live stream. Non-identity content encodings were rejected above.
|
| 495 |
async for chunk in response.aiter_bytes(chunk_size=64 * 1024):
|
| 496 |
if len(chunk) > limit - len(body):
|
| 497 |
-
raise
|
| 498 |
502, "The selected Decision runtime returned an oversized response"
|
| 499 |
)
|
| 500 |
body.extend(chunk)
|
| 501 |
try:
|
| 502 |
return _loads_strict(bytes(body))
|
| 503 |
except (TypeError, ValueError, UnicodeError, RecursionError) as exc:
|
| 504 |
-
raise
|
| 505 |
502, "The selected Decision runtime returned invalid JSON"
|
| 506 |
) from exc
|
| 507 |
|
|
|
|
| 1 |
+
"""Private transport from the public Gateway to per-model Decision runtimes.
|
| 2 |
|
| 3 |
The browser never receives backend locations. Every inference response is
|
| 4 |
validated against the originating strict request before it crosses the public
|
|
|
|
| 19 |
|
| 20 |
import httpx
|
| 21 |
|
| 22 |
+
from direct_contract import validate_response_for_request
|
| 23 |
|
| 24 |
SINGLE_PATH = "/v1/systemone"
|
| 25 |
BATCH_PATH = "/v1/systemone/batches"
|
|
|
|
| 55 |
)
|
| 56 |
|
| 57 |
|
| 58 |
+
class DirectGatewayError(RuntimeError):
|
| 59 |
"""A sanitized direct-backend failure safe to return to a caller."""
|
| 60 |
|
| 61 |
def __init__(self, code: int, message: str, *, retry_after: str | None = None):
|
|
|
|
| 88 |
"""Accept an explicit private/loopback IP origin, never a URL path or DNS name."""
|
| 89 |
|
| 90 |
if not isinstance(value, str) or not value or len(value) > 2048:
|
| 91 |
+
raise ValueError("Every direct endpoint must be a bounded URL string")
|
| 92 |
try:
|
| 93 |
parsed = urlsplit(value)
|
| 94 |
port = parsed.port
|
| 95 |
except ValueError:
|
| 96 |
+
raise ValueError("Every direct endpoint must be a valid private origin") from None
|
| 97 |
if (
|
| 98 |
parsed.scheme not in {"http", "https"}
|
| 99 |
or parsed.username is not None
|
|
|
|
| 106 |
or not 1 <= port <= 65535
|
| 107 |
):
|
| 108 |
raise ValueError(
|
| 109 |
+
"Every direct endpoint must be an http(s) origin with an explicit port"
|
| 110 |
)
|
| 111 |
if "%" in parsed.hostname:
|
| 112 |
+
raise ValueError("Scoped IPv6 direct endpoints are not supported")
|
| 113 |
try:
|
| 114 |
address = ipaddress.ip_address(parsed.hostname)
|
| 115 |
except ValueError:
|
| 116 |
raise ValueError(
|
| 117 |
+
"Every direct endpoint host must be a private or loopback IP literal"
|
| 118 |
) from None
|
| 119 |
if isinstance(address, ipaddress.IPv6Address) and address.ipv4_mapped:
|
| 120 |
address = address.ipv4_mapped
|
|
|
|
| 125 |
for network in PRIVATE_NETWORKS
|
| 126 |
)
|
| 127 |
):
|
| 128 |
+
raise ValueError("Every direct endpoint must stay on a private or loopback IP")
|
| 129 |
host = f"[{address.compressed}]" if address.version == 6 else address.compressed
|
| 130 |
return f"{parsed.scheme}://{host}:{port}"
|
| 131 |
|
| 132 |
|
| 133 |
def endpoints_from_environment(expected_models: Collection[str]) -> dict[str, str]:
|
| 134 |
+
raw = os.getenv("DECISION_DIRECT_ENDPOINTS", "")
|
| 135 |
if not raw:
|
| 136 |
+
raise ValueError("DECISION_DIRECT_ENDPOINTS is required for the direct backend")
|
| 137 |
if len(raw.encode("utf-8")) > MAX_ENDPOINT_CONFIG_BYTES:
|
| 138 |
+
raise ValueError("DECISION_DIRECT_ENDPOINTS is too large")
|
| 139 |
try:
|
| 140 |
value = _loads_strict(raw)
|
| 141 |
except (TypeError, ValueError, UnicodeError, RecursionError) as exc:
|
| 142 |
+
raise ValueError("DECISION_DIRECT_ENDPOINTS must be a strict JSON object") from exc
|
| 143 |
if not isinstance(value, dict) or set(value) != set(expected_models):
|
| 144 |
raise ValueError(
|
| 145 |
+
"DECISION_DIRECT_ENDPOINTS must cover exactly the configured canonical models"
|
| 146 |
)
|
| 147 |
return {model: _private_origin(endpoint) for model, endpoint in value.items()}
|
| 148 |
|
| 149 |
|
| 150 |
def timeout_from_environment() -> float:
|
| 151 |
+
raw = os.getenv("DECISION_DIRECT_TIMEOUT_SECONDS", "45")
|
| 152 |
try:
|
| 153 |
timeout = float(raw)
|
| 154 |
except ValueError as exc:
|
| 155 |
+
raise ValueError("DECISION_DIRECT_TIMEOUT_SECONDS must be numeric") from exc
|
| 156 |
if not math.isfinite(timeout) or not MIN_TIMEOUT_SECONDS <= timeout <= MAX_TIMEOUT_SECONDS:
|
| 157 |
raise ValueError(
|
| 158 |
+
"DECISION_DIRECT_TIMEOUT_SECONDS must be between 0.1 and 300"
|
| 159 |
)
|
| 160 |
return timeout
|
| 161 |
|
| 162 |
|
| 163 |
+
class DirectGateway:
|
| 164 |
"""Keep-alive direct client with fail-closed model and response validation."""
|
| 165 |
|
| 166 |
def __init__(
|
|
|
|
| 179 |
or any(not isinstance(model, str) or not model for model in expected)
|
| 180 |
or len(set(expected)) != len(expected)
|
| 181 |
):
|
| 182 |
+
raise ValueError("Configure a nonempty set of canonical Decision models")
|
| 183 |
if set(endpoints) != set(expected):
|
| 184 |
+
raise ValueError("direct endpoints must cover exactly the configured models")
|
| 185 |
pinned = {}
|
| 186 |
for model in expected:
|
| 187 |
item = expected_artifacts[model]
|
| 188 |
if not isinstance(item, Mapping):
|
| 189 |
+
raise ValueError("Pin each Decision model to a released artifact")
|
| 190 |
revision = item.get("revision")
|
| 191 |
manifest = item.get("manifest_sha256")
|
| 192 |
content = item.get("content_sha256")
|
|
|
|
| 195 |
or not _hex_digest(manifest, 64)
|
| 196 |
or (content is not None and not _hex_digest(content, 64))
|
| 197 |
):
|
| 198 |
+
raise ValueError("Pin each Decision model revision and manifest digest")
|
| 199 |
pinned[model] = MappingProxyType(
|
| 200 |
{
|
| 201 |
"revision": revision,
|
|
|
|
| 214 |
except OverflowError:
|
| 215 |
timeout_valid = False
|
| 216 |
if not timeout_valid:
|
| 217 |
+
raise ValueError("direct timeout must be between 0.1 and 300 seconds")
|
| 218 |
if (
|
| 219 |
type(max_response_bytes) is not int
|
| 220 |
or not 1 <= max_response_bytes <= MAX_INFERENCE_RESPONSE_BYTES
|
| 221 |
):
|
| 222 |
+
raise ValueError("Invalid direct response limit")
|
| 223 |
normalized = {model: _private_origin(endpoints[model]) for model in expected}
|
| 224 |
if len(set(normalized.values())) != len(normalized):
|
| 225 |
+
raise ValueError("Each configured Decision model requires its own direct origin")
|
| 226 |
self._endpoints = MappingProxyType(normalized)
|
| 227 |
self.expected_artifacts = MappingProxyType(pinned)
|
| 228 |
self.models = frozenset(expected)
|
|
|
|
| 239 |
@classmethod
|
| 240 |
def from_environment(
|
| 241 |
cls, expected_artifacts: Mapping[str, Mapping[str, str]]
|
| 242 |
+
) -> DirectGateway:
|
| 243 |
return cls(
|
| 244 |
endpoints_from_environment(expected_artifacts),
|
| 245 |
expected_artifacts=expected_artifacts,
|
|
|
|
| 268 |
|
| 269 |
model = payload.get("model") if isinstance(payload, Mapping) else None
|
| 270 |
if not isinstance(model, str) or model not in self.models:
|
| 271 |
+
raise DirectGatewayError(422, "The selected canonical model is not configured")
|
| 272 |
if batch != ("states" in payload) or ("state" in payload and "states" in payload):
|
| 273 |
+
raise DirectGatewayError(422, "The Decision request does not match its route")
|
| 274 |
try:
|
| 275 |
encoded = json.dumps(
|
| 276 |
payload,
|
|
|
|
| 279 |
allow_nan=False,
|
| 280 |
).encode("utf-8")
|
| 281 |
except (TypeError, ValueError, UnicodeError, RecursionError) as exc:
|
| 282 |
+
raise DirectGatewayError(422, "The Decision request is not valid JSON") from exc
|
| 283 |
request_limit = MAX_BATCH_REQUEST_BYTES if batch else MAX_SINGLE_REQUEST_BYTES
|
| 284 |
if len(encoded) > request_limit:
|
| 285 |
+
raise DirectGatewayError(413, "Request body exceeds the endpoint limit")
|
| 286 |
started = self._clock()
|
| 287 |
response = await self._request_json(
|
| 288 |
"POST",
|
|
|
|
| 296 |
try:
|
| 297 |
validate_response_for_request(payload, response, batch=batch)
|
| 298 |
except (TypeError, ValueError, KeyError, OverflowError, RecursionError) as exc:
|
| 299 |
+
raise DirectGatewayError(
|
| 300 |
502, "The selected Decision runtime returned an invalid response"
|
| 301 |
) from exc
|
| 302 |
return response, request_ms
|
|
|
|
| 327 |
and self._artifact_matches(model, artifact)
|
| 328 |
)
|
| 329 |
return {"loaded": bool(loaded), "artifact": artifact}
|
| 330 |
+
except DirectGatewayError:
|
| 331 |
return {"loaded": False, "artifact": None}
|
| 332 |
|
| 333 |
@staticmethod
|
|
|
|
| 390 |
):
|
| 391 |
endpoint = self._endpoints.get(model)
|
| 392 |
if endpoint is None:
|
| 393 |
+
raise DirectGatewayError(422, "The selected canonical model is not configured")
|
| 394 |
request_timeout = timeout_seconds or self.timeout_seconds
|
| 395 |
try:
|
| 396 |
headers = {
|
|
|
|
| 409 |
timeout=request_timeout,
|
| 410 |
) as response:
|
| 411 |
if 300 <= response.status_code < 400:
|
| 412 |
+
raise DirectGatewayError(
|
| 413 |
502,
|
| 414 |
"The selected Decision runtime returned an unsupported redirect",
|
| 415 |
)
|
|
|
|
| 418 |
if attest_model is not None and not self._artifact_matches(
|
| 419 |
attest_model, self._parse_response_artifact(response.headers)
|
| 420 |
):
|
| 421 |
+
raise DirectGatewayError(
|
| 422 |
503, "The selected Decision runtime artifact is unavailable"
|
| 423 |
)
|
| 424 |
return await self._read_json_response(response, limit)
|
| 425 |
except asyncio.CancelledError:
|
| 426 |
raise
|
| 427 |
+
except DirectGatewayError:
|
| 428 |
raise
|
| 429 |
except (TimeoutError, httpx.TimeoutException) as exc:
|
| 430 |
+
raise DirectGatewayError(
|
| 431 |
504, "The selected Decision runtime did not respond in time"
|
| 432 |
) from exc
|
| 433 |
except httpx.HTTPError as exc:
|
| 434 |
+
raise DirectGatewayError(
|
| 435 |
503, "The selected Decision runtime is unavailable"
|
| 436 |
) from exc
|
| 437 |
|
|
|
|
| 439 |
def _raise_status(response: httpx.Response) -> None:
|
| 440 |
status = response.status_code
|
| 441 |
if status == 413:
|
| 442 |
+
raise DirectGatewayError(
|
| 443 |
413, "Decision input exceeds the selected model limit"
|
| 444 |
)
|
| 445 |
if status in {429, 529}:
|
| 446 |
+
raise DirectGatewayError(
|
| 447 |
status,
|
| 448 |
"The selected Decision runtime is temporarily overloaded",
|
| 449 |
retry_after=_safe_retry_after(response.headers.get("retry-after")),
|
| 450 |
)
|
| 451 |
if status == 503:
|
| 452 |
+
raise DirectGatewayError(
|
| 453 |
503,
|
| 454 |
"The selected Decision runtime is unavailable",
|
| 455 |
retry_after=_safe_retry_after(response.headers.get("retry-after")),
|
| 456 |
)
|
| 457 |
if status == 504:
|
| 458 |
+
raise DirectGatewayError(
|
| 459 |
504, "The selected Decision runtime did not respond in time"
|
| 460 |
)
|
| 461 |
+
raise DirectGatewayError(
|
| 462 |
502, "The selected Decision runtime rejected a validated request"
|
| 463 |
)
|
| 464 |
|
|
|
|
| 468 |
response.headers.get("content-type", "").split(";", 1)[0].strip().lower()
|
| 469 |
)
|
| 470 |
if content_type != "application/json":
|
| 471 |
+
raise DirectGatewayError(
|
| 472 |
502, "The selected Decision runtime returned an invalid content type"
|
| 473 |
)
|
| 474 |
encoding = response.headers.get("content-encoding", "identity").strip().lower()
|
| 475 |
if encoding not in {"", "identity"}:
|
| 476 |
+
raise DirectGatewayError(
|
| 477 |
502, "The selected Decision runtime returned an unsupported encoding"
|
| 478 |
)
|
| 479 |
declared = response.headers.get("content-length")
|
|
|
|
| 481 |
try:
|
| 482 |
declared_size = int(declared)
|
| 483 |
except ValueError as exc:
|
| 484 |
+
raise DirectGatewayError(
|
| 485 |
502, "The selected Decision runtime returned an invalid response"
|
| 486 |
) from exc
|
| 487 |
if declared_size < 0 or declared_size > limit:
|
| 488 |
+
raise DirectGatewayError(
|
| 489 |
502, "The selected Decision runtime returned an oversized response"
|
| 490 |
)
|
| 491 |
body = bytearray()
|
|
|
|
| 494 |
# live stream. Non-identity content encodings were rejected above.
|
| 495 |
async for chunk in response.aiter_bytes(chunk_size=64 * 1024):
|
| 496 |
if len(chunk) > limit - len(body):
|
| 497 |
+
raise DirectGatewayError(
|
| 498 |
502, "The selected Decision runtime returned an oversized response"
|
| 499 |
)
|
| 500 |
body.extend(chunk)
|
| 501 |
try:
|
| 502 |
return _loads_strict(bytes(body))
|
| 503 |
except (TypeError, ValueError, UnicodeError, RecursionError) as exc:
|
| 504 |
+
raise DirectGatewayError(
|
| 505 |
502, "The selected Decision runtime returned invalid JSON"
|
| 506 |
) from exc
|
| 507 |
|
static/app.js
CHANGED
|
@@ -300,9 +300,9 @@ async function run() {
|
|
| 300 |
try {
|
| 301 |
const result = status.backend === 'pull_queue' ? await queuedPrediction(payload, current) : await nativePrediction(payload, current);
|
| 302 |
if(version!==revision || epoch!==modelEpoch || current.signal.aborted) return;
|
| 303 |
-
const expectedModel=status.backend==='
|
| 304 |
if(result.model!==expectedModel || payload.model!==selectedModel) throw Error('The server returned a different model.');
|
| 305 |
-
if(status.backend!=='
|
| 306 |
renderResult(result,payload);
|
| 307 |
} catch(e) {if(e.name!=='AbortError' && version===revision){setModelActivity('error');setRunStatus('Could not complete this request','error');toast(e.message); if(!lastResult)$('results').innerHTML=`<div class="error-message">${esc(e.message)}</div>`;}}
|
| 308 |
finally{if(controller===current)controller=null;updateRunButton();if(epoch===modelEpoch)await checkStatus();}
|
|
|
|
| 300 |
try {
|
| 301 |
const result = status.backend === 'pull_queue' ? await queuedPrediction(payload, current) : await nativePrediction(payload, current);
|
| 302 |
if(version!==revision || epoch!==modelEpoch || current.signal.aborted) return;
|
| 303 |
+
const expectedModel=status.backend==='direct'?modelInfo().repo_id:selectedModel;
|
| 304 |
if(result.model!==expectedModel || payload.model!==selectedModel) throw Error('The server returned a different model.');
|
| 305 |
+
if(status.backend!=='direct'&&result.source!=='live_native') throw Error('This server did not return a live native prediction.');
|
| 306 |
renderResult(result,payload);
|
| 307 |
} catch(e) {if(e.name!=='AbortError' && version===revision){setModelActivity('error');setRunStatus('Could not complete this request','error');toast(e.message); if(!lastResult)$('results').innerHTML=`<div class="error-message">${esc(e.message)}</div>`;}}
|
| 308 |
finally{if(controller===current)controller=null;updateRunButton();if(epoch===modelEpoch)await checkStatus();}
|
tests/studio_contract.test.mjs
CHANGED
|
@@ -73,7 +73,7 @@ test('browser request bytes match the single and batch API routes', () => {
|
|
| 73 |
}
|
| 74 |
});
|
| 75 |
|
| 76 |
-
test('strict
|
| 77 |
assert.equal(validateResult(strict, request, canonical), strict);
|
| 78 |
assert.throws(() => validateResult({ ...strict, model: request.model }, request, canonical));
|
| 79 |
assert.throws(() => validateResult({ ...strict, timing: {} }, request, canonical));
|
|
@@ -83,7 +83,7 @@ test('strict drun single response needs canonical identity and no diagnostics',
|
|
| 83 |
} }, request, canonical));
|
| 84 |
});
|
| 85 |
|
| 86 |
-
test('strict
|
| 87 |
const batchRequest = {
|
| 88 |
...request,
|
| 89 |
states: [
|
|
|
|
| 73 |
}
|
| 74 |
});
|
| 75 |
|
| 76 |
+
test('strict direct single response needs canonical identity and no diagnostics', () => {
|
| 77 |
assert.equal(validateResult(strict, request, canonical), strict);
|
| 78 |
assert.throws(() => validateResult({ ...strict, model: request.model }, request, canonical));
|
| 79 |
assert.throws(() => validateResult({ ...strict, timing: {} }, request, canonical));
|
|
|
|
| 83 |
} }, request, canonical));
|
| 84 |
});
|
| 85 |
|
| 86 |
+
test('strict direct batch response preserves state order and per-row usage', () => {
|
| 87 |
const batchRequest = {
|
| 88 |
...request,
|
| 89 |
states: [
|
tests/{test_drun_gateway.py → test_direct_gateway.py}
RENAMED
|
@@ -1,4 +1,4 @@
|
|
| 1 |
-
"""Direct
|
| 2 |
|
| 3 |
import asyncio
|
| 4 |
import inspect
|
|
@@ -12,7 +12,7 @@ from fastapi.testclient import TestClient
|
|
| 12 |
|
| 13 |
from app import create_app
|
| 14 |
from contract import MODEL
|
| 15 |
-
from
|
| 16 |
from engine import DEFAULT_MANIFEST
|
| 17 |
from model_registry import PROFILES
|
| 18 |
from tetris_arena import Competitor, HTTPDecisionAdapter, LOCAL_MODELS
|
|
@@ -164,13 +164,13 @@ def batch_response(model=CANONICAL):
|
|
| 164 |
}
|
| 165 |
|
| 166 |
|
| 167 |
-
class
|
| 168 |
async def gateway(self, handler, *, control_handler=False, **kwargs):
|
| 169 |
artifact_pins = kwargs.pop("expected_artifacts", {CANONICAL: pins()})
|
| 170 |
client = httpx.AsyncClient(
|
| 171 |
transport=httpx.MockTransport(handler if control_handler else with_control(handler))
|
| 172 |
)
|
| 173 |
-
gateway =
|
| 174 |
{CANONICAL: ORIGIN},
|
| 175 |
expected_artifacts=artifact_pins,
|
| 176 |
client=client,
|
|
@@ -179,6 +179,24 @@ class DrunTransportTests(unittest.IsolatedAsyncioTestCase):
|
|
| 179 |
self.addAsyncCleanup(client.aclose)
|
| 180 |
return gateway
|
| 181 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 182 |
async def test_routes_single_to_exact_model_and_preserves_non_max_confidence(self):
|
| 183 |
seen = []
|
| 184 |
|
|
@@ -258,7 +276,7 @@ class DrunTransportTests(unittest.IsolatedAsyncioTestCase):
|
|
| 258 |
|
| 259 |
client = httpx.AsyncClient(transport=httpx.MockTransport(with_control(handler)))
|
| 260 |
self.addAsyncCleanup(client.aclose)
|
| 261 |
-
gateway =
|
| 262 |
{CANONICAL: ORIGIN, LUX: LUX_ORIGIN},
|
| 263 |
expected_artifacts={CANONICAL: pins(), LUX: pins(LUX)},
|
| 264 |
client=client,
|
|
@@ -279,7 +297,7 @@ class DrunTransportTests(unittest.IsolatedAsyncioTestCase):
|
|
| 279 |
|
| 280 |
client = httpx.AsyncClient(transport=httpx.MockTransport(with_control(handler)))
|
| 281 |
self.addAsyncCleanup(client.aclose)
|
| 282 |
-
gateway =
|
| 283 |
{CANONICAL: bridge_origin},
|
| 284 |
expected_artifacts={CANONICAL: pins()},
|
| 285 |
client=client,
|
|
@@ -313,7 +331,7 @@ class DrunTransportTests(unittest.IsolatedAsyncioTestCase):
|
|
| 313 |
|
| 314 |
gateway = await self.gateway(handler)
|
| 315 |
for expected in ("confidence", "model", "timing"):
|
| 316 |
-
with self.assertRaises(
|
| 317 |
await gateway.evaluate(single_payload(), batch=False)
|
| 318 |
self.assertEqual(raised.exception.code, 502)
|
| 319 |
self.assertNotIn(ORIGIN, raised.exception.message)
|
|
@@ -328,7 +346,7 @@ class DrunTransportTests(unittest.IsolatedAsyncioTestCase):
|
|
| 328 |
)
|
| 329 |
|
| 330 |
gateway = await self.gateway(handler, max_response_bytes=128)
|
| 331 |
-
with self.assertRaises(
|
| 332 |
await gateway.evaluate(single_payload(), batch=False)
|
| 333 |
self.assertEqual(raised.exception.code, 502)
|
| 334 |
self.assertIn("oversized", raised.exception.message)
|
|
@@ -342,7 +360,7 @@ class DrunTransportTests(unittest.IsolatedAsyncioTestCase):
|
|
| 342 |
)
|
| 343 |
|
| 344 |
gateway = await self.gateway(handler)
|
| 345 |
-
with self.assertRaises(
|
| 346 |
await gateway.evaluate(single_payload(), batch=False)
|
| 347 |
error = raised.exception
|
| 348 |
self.assertEqual(error.code, 429)
|
|
@@ -359,7 +377,7 @@ class DrunTransportTests(unittest.IsolatedAsyncioTestCase):
|
|
| 359 |
)
|
| 360 |
|
| 361 |
gateway = await self.gateway(handler)
|
| 362 |
-
with self.assertRaises(
|
| 363 |
await gateway.evaluate(single_payload(), batch=False)
|
| 364 |
self.assertEqual(raised.exception.code, 529)
|
| 365 |
self.assertEqual(raised.exception.retry_after, "1")
|
|
@@ -371,7 +389,7 @@ class DrunTransportTests(unittest.IsolatedAsyncioTestCase):
|
|
| 371 |
return httpx.Response(200, json=single_response())
|
| 372 |
|
| 373 |
gateway = await self.gateway(handler, timeout_seconds=0.1)
|
| 374 |
-
with self.assertRaises(
|
| 375 |
await gateway.evaluate(single_payload(), batch=False)
|
| 376 |
self.assertEqual(raised.exception.code, 504)
|
| 377 |
self.assertNotIn(ORIGIN, raised.exception.message)
|
|
@@ -417,7 +435,7 @@ class DrunTransportTests(unittest.IsolatedAsyncioTestCase):
|
|
| 417 |
)
|
| 418 |
|
| 419 |
gateway = await self.gateway(handler, control_handler=True)
|
| 420 |
-
with self.assertRaises(
|
| 421 |
await getattr(gateway, method_name)(single_payload(), batch=False)
|
| 422 |
self.assertEqual(raised.exception.code, 503)
|
| 423 |
self.assertEqual(calls, ["/v1/systemone"])
|
|
@@ -461,7 +479,7 @@ class DrunTransportTests(unittest.IsolatedAsyncioTestCase):
|
|
| 461 |
control_handler=True,
|
| 462 |
expected_artifacts={CANONICAL: {**pins(), "content_sha256": CONTENT}},
|
| 463 |
)
|
| 464 |
-
with self.assertRaises(
|
| 465 |
await getattr(gateway, method_name)(single_payload(), batch=False)
|
| 466 |
self.assertEqual(raised.exception.code, 503)
|
| 467 |
self.assertEqual(calls, ["/v1/systemone"])
|
|
@@ -541,7 +559,7 @@ class DrunTransportTests(unittest.IsolatedAsyncioTestCase):
|
|
| 541 |
|
| 542 |
def test_direct_mode_requires_pinned_revision(self):
|
| 543 |
with self.assertRaisesRegex(ValueError, "revision"):
|
| 544 |
-
|
| 545 |
{CANONICAL: ORIGIN},
|
| 546 |
expected_artifacts={CANONICAL: {"manifest_sha256": DEFAULT_MANIFEST}},
|
| 547 |
)
|
|
@@ -556,16 +574,16 @@ class DrunTransportTests(unittest.IsolatedAsyncioTestCase):
|
|
| 556 |
)
|
| 557 |
for endpoint in invalid:
|
| 558 |
with self.subTest(endpoint=endpoint), self.assertRaises(ValueError):
|
| 559 |
-
|
| 560 |
{CANONICAL: endpoint},
|
| 561 |
expected_artifacts={CANONICAL: pins()},
|
| 562 |
)
|
| 563 |
|
| 564 |
def test_endpoint_configuration_must_exactly_cover_registry(self):
|
| 565 |
with self.assertRaises(ValueError):
|
| 566 |
-
|
| 567 |
with self.assertRaises(ValueError):
|
| 568 |
-
|
| 569 |
{
|
| 570 |
CANONICAL: ORIGIN,
|
| 571 |
"llm-semantic-router/Decision-1.0-Lux-9B": (
|
|
@@ -575,7 +593,7 @@ class DrunTransportTests(unittest.IsolatedAsyncioTestCase):
|
|
| 575 |
expected_artifacts={CANONICAL: pins()},
|
| 576 |
)
|
| 577 |
with self.assertRaises(ValueError):
|
| 578 |
-
|
| 579 |
{CANONICAL: ORIGIN, LUX: ORIGIN},
|
| 580 |
expected_artifacts={CANONICAL: pins(), LUX: pins(LUX)},
|
| 581 |
)
|
|
@@ -644,12 +662,12 @@ class SlowCloudAdapter:
|
|
| 644 |
await asyncio.sleep(60)
|
| 645 |
|
| 646 |
|
| 647 |
-
class
|
| 648 |
def setUp(self):
|
| 649 |
self.direct = FakeDirectGateway()
|
| 650 |
app = create_app(
|
| 651 |
-
mode="
|
| 652 |
-
|
| 653 |
tetris_manager=FakeTetrisManager(),
|
| 654 |
registry=[
|
| 655 |
{
|
|
@@ -713,11 +731,6 @@ class DrunAppTests(unittest.TestCase):
|
|
| 713 |
self.direct.calls[-1],
|
| 714 |
("evaluate", batch_payload(), True),
|
| 715 |
)
|
| 716 |
-
self.assertEqual(
|
| 717 |
-
self.client.post("/v1/systemone/batch", json=batch_payload()).status_code,
|
| 718 |
-
405,
|
| 719 |
-
)
|
| 720 |
-
|
| 721 |
def test_status_is_compatible_and_never_contains_private_topology(self):
|
| 722 |
response = self.client.get("/api/status", params={"model": MODEL})
|
| 723 |
self.assertEqual(response.status_code, 200)
|
|
@@ -728,7 +741,7 @@ class DrunAppTests(unittest.TestCase):
|
|
| 728 |
self.assertEqual(body["loaded_revision"], REVISION)
|
| 729 |
self.assertEqual(body["loaded_manifest_sha256"], DEFAULT_MANIFEST)
|
| 730 |
self.assertEqual(body["loaded_content_sha256"], CONTENT)
|
| 731 |
-
self.assertEqual(body["backend"], "
|
| 732 |
self.assertTrue(body["loaded"])
|
| 733 |
self.assertTrue(body["context_batch"])
|
| 734 |
serialized = json.dumps(body)
|
|
@@ -779,7 +792,7 @@ class DrunAppTests(unittest.TestCase):
|
|
| 779 |
"revision": LUX_REVISION,
|
| 780 |
},
|
| 781 |
]
|
| 782 |
-
app = create_app(mode="
|
| 783 |
race = {"left": "kai", "right": "lux", "mode": "steps", "max_steps": 1}
|
| 784 |
with TestClient(app) as client:
|
| 785 |
self.assertEqual(client.get("/api/ready").status_code, 503)
|
|
@@ -812,8 +825,8 @@ class DrunAppTests(unittest.TestCase):
|
|
| 812 |
direct = SixModelDirectGateway()
|
| 813 |
direct.probe = disconnected
|
| 814 |
app = create_app(
|
| 815 |
-
mode="
|
| 816 |
-
|
| 817 |
tetris_adapter=SlowCloudAdapter(),
|
| 818 |
registry=[
|
| 819 |
{
|
|
@@ -856,7 +869,7 @@ class DrunAppTests(unittest.TestCase):
|
|
| 856 |
self.assertIsNone(limits["gpu_microbatch"])
|
| 857 |
self.assertEqual(limits["synchronous_wait_seconds"], 17.0)
|
| 858 |
|
| 859 |
-
def
|
| 860 |
with patch.dict(
|
| 861 |
os.environ,
|
| 862 |
{
|
|
@@ -865,8 +878,8 @@ class DrunAppTests(unittest.TestCase):
|
|
| 865 |
},
|
| 866 |
):
|
| 867 |
app = create_app(
|
| 868 |
-
mode="
|
| 869 |
-
|
| 870 |
registry=[
|
| 871 |
{
|
| 872 |
"id": MODEL,
|
|
|
|
| 1 |
+
"""Direct routing stays private and preserves the strict response contract."""
|
| 2 |
|
| 3 |
import asyncio
|
| 4 |
import inspect
|
|
|
|
| 12 |
|
| 13 |
from app import create_app
|
| 14 |
from contract import MODEL
|
| 15 |
+
from direct_gateway import ARTIFACT_RESPONSE_HEADERS, DirectGateway, DirectGatewayError
|
| 16 |
from engine import DEFAULT_MANIFEST
|
| 17 |
from model_registry import PROFILES
|
| 18 |
from tetris_arena import Competitor, HTTPDecisionAdapter, LOCAL_MODELS
|
|
|
|
| 164 |
}
|
| 165 |
|
| 166 |
|
| 167 |
+
class DirectTransportTests(unittest.IsolatedAsyncioTestCase):
|
| 168 |
async def gateway(self, handler, *, control_handler=False, **kwargs):
|
| 169 |
artifact_pins = kwargs.pop("expected_artifacts", {CANONICAL: pins()})
|
| 170 |
client = httpx.AsyncClient(
|
| 171 |
transport=httpx.MockTransport(handler if control_handler else with_control(handler))
|
| 172 |
)
|
| 173 |
+
gateway = DirectGateway(
|
| 174 |
{CANONICAL: ORIGIN},
|
| 175 |
expected_artifacts=artifact_pins,
|
| 176 |
client=client,
|
|
|
|
| 179 |
self.addAsyncCleanup(client.aclose)
|
| 180 |
return gateway
|
| 181 |
|
| 182 |
+
async def test_direct_environment_configures_private_origin_and_timeout(self):
|
| 183 |
+
with patch.dict(
|
| 184 |
+
os.environ,
|
| 185 |
+
{
|
| 186 |
+
"DECISION_DIRECT_ENDPOINTS": json.dumps({CANONICAL: ORIGIN}),
|
| 187 |
+
"DECISION_DIRECT_TIMEOUT_SECONDS": "17",
|
| 188 |
+
},
|
| 189 |
+
clear=True,
|
| 190 |
+
):
|
| 191 |
+
gateway = DirectGateway.from_environment({CANONICAL: pins()})
|
| 192 |
+
self.addAsyncCleanup(gateway.aclose)
|
| 193 |
+
|
| 194 |
+
self.assertEqual(
|
| 195 |
+
gateway.local_systemone_endpoints(),
|
| 196 |
+
{CANONICAL: ORIGIN + "/v1/systemone"},
|
| 197 |
+
)
|
| 198 |
+
self.assertEqual(gateway.timeout_seconds, 17)
|
| 199 |
+
|
| 200 |
async def test_routes_single_to_exact_model_and_preserves_non_max_confidence(self):
|
| 201 |
seen = []
|
| 202 |
|
|
|
|
| 276 |
|
| 277 |
client = httpx.AsyncClient(transport=httpx.MockTransport(with_control(handler)))
|
| 278 |
self.addAsyncCleanup(client.aclose)
|
| 279 |
+
gateway = DirectGateway(
|
| 280 |
{CANONICAL: ORIGIN, LUX: LUX_ORIGIN},
|
| 281 |
expected_artifacts={CANONICAL: pins(), LUX: pins(LUX)},
|
| 282 |
client=client,
|
|
|
|
| 297 |
|
| 298 |
client = httpx.AsyncClient(transport=httpx.MockTransport(with_control(handler)))
|
| 299 |
self.addAsyncCleanup(client.aclose)
|
| 300 |
+
gateway = DirectGateway(
|
| 301 |
{CANONICAL: bridge_origin},
|
| 302 |
expected_artifacts={CANONICAL: pins()},
|
| 303 |
client=client,
|
|
|
|
| 331 |
|
| 332 |
gateway = await self.gateway(handler)
|
| 333 |
for expected in ("confidence", "model", "timing"):
|
| 334 |
+
with self.assertRaises(DirectGatewayError) as raised:
|
| 335 |
await gateway.evaluate(single_payload(), batch=False)
|
| 336 |
self.assertEqual(raised.exception.code, 502)
|
| 337 |
self.assertNotIn(ORIGIN, raised.exception.message)
|
|
|
|
| 346 |
)
|
| 347 |
|
| 348 |
gateway = await self.gateway(handler, max_response_bytes=128)
|
| 349 |
+
with self.assertRaises(DirectGatewayError) as raised:
|
| 350 |
await gateway.evaluate(single_payload(), batch=False)
|
| 351 |
self.assertEqual(raised.exception.code, 502)
|
| 352 |
self.assertIn("oversized", raised.exception.message)
|
|
|
|
| 360 |
)
|
| 361 |
|
| 362 |
gateway = await self.gateway(handler)
|
| 363 |
+
with self.assertRaises(DirectGatewayError) as raised:
|
| 364 |
await gateway.evaluate(single_payload(), batch=False)
|
| 365 |
error = raised.exception
|
| 366 |
self.assertEqual(error.code, 429)
|
|
|
|
| 377 |
)
|
| 378 |
|
| 379 |
gateway = await self.gateway(handler)
|
| 380 |
+
with self.assertRaises(DirectGatewayError) as raised:
|
| 381 |
await gateway.evaluate(single_payload(), batch=False)
|
| 382 |
self.assertEqual(raised.exception.code, 529)
|
| 383 |
self.assertEqual(raised.exception.retry_after, "1")
|
|
|
|
| 389 |
return httpx.Response(200, json=single_response())
|
| 390 |
|
| 391 |
gateway = await self.gateway(handler, timeout_seconds=0.1)
|
| 392 |
+
with self.assertRaises(DirectGatewayError) as raised:
|
| 393 |
await gateway.evaluate(single_payload(), batch=False)
|
| 394 |
self.assertEqual(raised.exception.code, 504)
|
| 395 |
self.assertNotIn(ORIGIN, raised.exception.message)
|
|
|
|
| 435 |
)
|
| 436 |
|
| 437 |
gateway = await self.gateway(handler, control_handler=True)
|
| 438 |
+
with self.assertRaises(DirectGatewayError) as raised:
|
| 439 |
await getattr(gateway, method_name)(single_payload(), batch=False)
|
| 440 |
self.assertEqual(raised.exception.code, 503)
|
| 441 |
self.assertEqual(calls, ["/v1/systemone"])
|
|
|
|
| 479 |
control_handler=True,
|
| 480 |
expected_artifacts={CANONICAL: {**pins(), "content_sha256": CONTENT}},
|
| 481 |
)
|
| 482 |
+
with self.assertRaises(DirectGatewayError) as raised:
|
| 483 |
await getattr(gateway, method_name)(single_payload(), batch=False)
|
| 484 |
self.assertEqual(raised.exception.code, 503)
|
| 485 |
self.assertEqual(calls, ["/v1/systemone"])
|
|
|
|
| 559 |
|
| 560 |
def test_direct_mode_requires_pinned_revision(self):
|
| 561 |
with self.assertRaisesRegex(ValueError, "revision"):
|
| 562 |
+
DirectGateway(
|
| 563 |
{CANONICAL: ORIGIN},
|
| 564 |
expected_artifacts={CANONICAL: {"manifest_sha256": DEFAULT_MANIFEST}},
|
| 565 |
)
|
|
|
|
| 574 |
)
|
| 575 |
for endpoint in invalid:
|
| 576 |
with self.subTest(endpoint=endpoint), self.assertRaises(ValueError):
|
| 577 |
+
DirectGateway(
|
| 578 |
{CANONICAL: endpoint},
|
| 579 |
expected_artifacts={CANONICAL: pins()},
|
| 580 |
)
|
| 581 |
|
| 582 |
def test_endpoint_configuration_must_exactly_cover_registry(self):
|
| 583 |
with self.assertRaises(ValueError):
|
| 584 |
+
DirectGateway({}, expected_artifacts={CANONICAL: pins()})
|
| 585 |
with self.assertRaises(ValueError):
|
| 586 |
+
DirectGateway(
|
| 587 |
{
|
| 588 |
CANONICAL: ORIGIN,
|
| 589 |
"llm-semantic-router/Decision-1.0-Lux-9B": (
|
|
|
|
| 593 |
expected_artifacts={CANONICAL: pins()},
|
| 594 |
)
|
| 595 |
with self.assertRaises(ValueError):
|
| 596 |
+
DirectGateway(
|
| 597 |
{CANONICAL: ORIGIN, LUX: ORIGIN},
|
| 598 |
expected_artifacts={CANONICAL: pins(), LUX: pins(LUX)},
|
| 599 |
)
|
|
|
|
| 662 |
await asyncio.sleep(60)
|
| 663 |
|
| 664 |
|
| 665 |
+
class DirectAppTests(unittest.TestCase):
|
| 666 |
def setUp(self):
|
| 667 |
self.direct = FakeDirectGateway()
|
| 668 |
app = create_app(
|
| 669 |
+
mode="direct",
|
| 670 |
+
direct_gateway=self.direct,
|
| 671 |
tetris_manager=FakeTetrisManager(),
|
| 672 |
registry=[
|
| 673 |
{
|
|
|
|
| 731 |
self.direct.calls[-1],
|
| 732 |
("evaluate", batch_payload(), True),
|
| 733 |
)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 734 |
def test_status_is_compatible_and_never_contains_private_topology(self):
|
| 735 |
response = self.client.get("/api/status", params={"model": MODEL})
|
| 736 |
self.assertEqual(response.status_code, 200)
|
|
|
|
| 741 |
self.assertEqual(body["loaded_revision"], REVISION)
|
| 742 |
self.assertEqual(body["loaded_manifest_sha256"], DEFAULT_MANIFEST)
|
| 743 |
self.assertEqual(body["loaded_content_sha256"], CONTENT)
|
| 744 |
+
self.assertEqual(body["backend"], "direct")
|
| 745 |
self.assertTrue(body["loaded"])
|
| 746 |
self.assertTrue(body["context_batch"])
|
| 747 |
serialized = json.dumps(body)
|
|
|
|
| 792 |
"revision": LUX_REVISION,
|
| 793 |
},
|
| 794 |
]
|
| 795 |
+
app = create_app(mode="direct", direct_gateway=direct, registry=registry)
|
| 796 |
race = {"left": "kai", "right": "lux", "mode": "steps", "max_steps": 1}
|
| 797 |
with TestClient(app) as client:
|
| 798 |
self.assertEqual(client.get("/api/ready").status_code, 503)
|
|
|
|
| 825 |
direct = SixModelDirectGateway()
|
| 826 |
direct.probe = disconnected
|
| 827 |
app = create_app(
|
| 828 |
+
mode="direct",
|
| 829 |
+
direct_gateway=direct,
|
| 830 |
tetris_adapter=SlowCloudAdapter(),
|
| 831 |
registry=[
|
| 832 |
{
|
|
|
|
| 869 |
self.assertIsNone(limits["gpu_microbatch"])
|
| 870 |
self.assertEqual(limits["synchronous_wait_seconds"], 17.0)
|
| 871 |
|
| 872 |
+
def test_tetris_uses_private_direct_route_even_when_legacy_local_url_is_set(self):
|
| 873 |
with patch.dict(
|
| 874 |
os.environ,
|
| 875 |
{
|
|
|
|
| 878 |
},
|
| 879 |
):
|
| 880 |
app = create_app(
|
| 881 |
+
mode="direct",
|
| 882 |
+
direct_gateway=FakeDirectGateway(),
|
| 883 |
registry=[
|
| 884 |
{
|
| 885 |
"id": MODEL,
|
tests/test_public_api_contract.py
CHANGED
|
@@ -99,8 +99,6 @@ class PublicAPIContractTests(unittest.TestCase):
|
|
| 99 |
self.assertEqual(self.client.post("/v1/systemone/batches", json=duplicate).status_code, 422)
|
| 100 |
self.assertEqual(self.client.post("/v1/systemone/batches", json=dict(payload, state="extra")).status_code, 422)
|
| 101 |
self.assertEqual(self.client.post("/v1/systemone/batches", json=dict(payload, model=MODEL)).status_code, 422)
|
| 102 |
-
self.assertEqual(self.client.post("/v1/decision/batches", json=payload).status_code, 405)
|
| 103 |
-
self.assertEqual(self.client.post("/v1/systemone/batch", json=payload).status_code, 405)
|
| 104 |
|
| 105 |
|
| 106 |
if __name__ == "__main__":
|
|
|
|
| 99 |
self.assertEqual(self.client.post("/v1/systemone/batches", json=duplicate).status_code, 422)
|
| 100 |
self.assertEqual(self.client.post("/v1/systemone/batches", json=dict(payload, state="extra")).status_code, 422)
|
| 101 |
self.assertEqual(self.client.post("/v1/systemone/batches", json=dict(payload, model=MODEL)).status_code, 422)
|
|
|
|
|
|
|
| 102 |
|
| 103 |
|
| 104 |
if __name__ == "__main__":
|
tetris_arena.py
CHANGED
|
@@ -23,7 +23,7 @@ from dataclasses import dataclass, field
|
|
| 23 |
from typing import Any
|
| 24 |
from urllib.parse import urlsplit
|
| 25 |
|
| 26 |
-
from
|
| 27 |
from model_registry import MODEL_ORDER, PROFILES
|
| 28 |
|
| 29 |
BOARD_WIDTH = 10
|
|
@@ -469,7 +469,7 @@ class HTTPDecisionAdapter:
|
|
| 469 |
) -> HTTPDecisionAdapter:
|
| 470 |
env = os.environ if env is None else env
|
| 471 |
if direct_gateway is not None and local_endpoints is None:
|
| 472 |
-
raise ValueError("Direct Tetris routing requires the
|
| 473 |
shared_local_key = _first_environment(env, "TETRIS_LOCAL_API_KEY")
|
| 474 |
endpoints: dict[str, Endpoint] = {}
|
| 475 |
if local_endpoints is not None:
|
|
@@ -481,7 +481,7 @@ class HTTPDecisionAdapter:
|
|
| 481 |
raw_url = local_endpoints.get(competitor.request_model)
|
| 482 |
if raw_url:
|
| 483 |
endpoints[competitor.id] = Endpoint(
|
| 484 |
-
_valid_url(raw_url, "Private
|
| 485 |
"",
|
| 486 |
competitor.request_model,
|
| 487 |
)
|
|
@@ -598,7 +598,7 @@ class HTTPDecisionAdapter:
|
|
| 598 |
if self._direct_gateway is not None and competitor.family == "decision":
|
| 599 |
try:
|
| 600 |
document = await self._direct_gateway.evaluate(upstream_payload, batch=False)
|
| 601 |
-
except
|
| 602 |
raise ArenaUpstreamError(
|
| 603 |
"The selected model could not complete this turn."
|
| 604 |
) from exc
|
|
|
|
| 23 |
from typing import Any
|
| 24 |
from urllib.parse import urlsplit
|
| 25 |
|
| 26 |
+
from direct_gateway import DirectGatewayError
|
| 27 |
from model_registry import MODEL_ORDER, PROFILES
|
| 28 |
|
| 29 |
BOARD_WIDTH = 10
|
|
|
|
| 469 |
) -> HTTPDecisionAdapter:
|
| 470 |
env = os.environ if env is None else env
|
| 471 |
if direct_gateway is not None and local_endpoints is None:
|
| 472 |
+
raise ValueError("Direct Tetris routing requires the configured private Decision endpoint map")
|
| 473 |
shared_local_key = _first_environment(env, "TETRIS_LOCAL_API_KEY")
|
| 474 |
endpoints: dict[str, Endpoint] = {}
|
| 475 |
if local_endpoints is not None:
|
|
|
|
| 481 |
raw_url = local_endpoints.get(competitor.request_model)
|
| 482 |
if raw_url:
|
| 483 |
endpoints[competitor.id] = Endpoint(
|
| 484 |
+
_valid_url(raw_url, "Private Decision endpoint"),
|
| 485 |
"",
|
| 486 |
competitor.request_model,
|
| 487 |
)
|
|
|
|
| 598 |
if self._direct_gateway is not None and competitor.family == "decision":
|
| 599 |
try:
|
| 600 |
document = await self._direct_gateway.evaluate(upstream_payload, batch=False)
|
| 601 |
+
except DirectGatewayError as exc:
|
| 602 |
raise ArenaUpstreamError(
|
| 603 |
"The selected model could not complete this turn."
|
| 604 |
) from exc
|