diff --git a/cmd/gomodel/docs/docs.go b/cmd/gomodel/docs/docs.go index dacf4d64d..90f72c5dc 100644 --- a/cmd/gomodel/docs/docs.go +++ b/cmd/gomodel/docs/docs.go @@ -7573,6 +7573,144 @@ const docTemplate = `{ ] } }, + "/v1/systemone/permute": { + "post": { + "description": "A Kev server diagnostic: the request is a System One request, and n_perm (1 to 64, default 6) sets how many option orders run. Only jev providers pointing at a Kev server serve it.", + "consumes": [ + "application/json" + ], + "produces": [ + "application/json" + ], + "tags": [ + "systemone" + ], + "summary": "Run one Choice question with several option orders (Kev)", + "parameters": [ + { + "description": "System One request with one Choice question", + "name": "request", + "in": "body", + "required": true, + "schema": { + "type": "object" + } + } + ], + "responses": { + "200": { + "description": "Kev's answer, in the provider's shape", + "schema": { + "type": "object" + } + }, + "400": { + "description": "Bad Request", + "schema": { + "$ref": "#/definitions/core.OpenAIErrorEnvelope" + } + }, + "401": { + "description": "Unauthorized", + "schema": { + "$ref": "#/definitions/core.OpenAIErrorEnvelope" + } + }, + "404": { + "description": "Not Found", + "schema": { + "$ref": "#/definitions/core.OpenAIErrorEnvelope" + } + }, + "429": { + "description": "Too Many Requests", + "schema": { + "$ref": "#/definitions/core.OpenAIErrorEnvelope" + } + }, + "502": { + "description": "Bad Gateway", + "schema": { + "$ref": "#/definitions/core.OpenAIErrorEnvelope" + } + } + }, + "security": [ + { + "BearerAuth": [] + } + ] + } + }, + "/v1/systemone/separate": { + "post": { + "description": "A Kev server diagnostic that answers each question separately. Only jev providers pointing at a Kev server serve it.", + "consumes": [ + "application/json" + ], + "produces": [ + "application/json" + ], + "tags": [ + "systemone" + ], + "summary": "Run each System One question in its own forward pass (Kev)", + "parameters": [ + { + "description": "System One request: model, state, and questions", + "name": "request", + "in": "body", + "required": true, + "schema": { + "type": "object" + } + } + ], + "responses": { + "200": { + "description": "Kev's answer, in the provider's shape", + "schema": { + "type": "object" + } + }, + "400": { + "description": "Bad Request", + "schema": { + "$ref": "#/definitions/core.OpenAIErrorEnvelope" + } + }, + "401": { + "description": "Unauthorized", + "schema": { + "$ref": "#/definitions/core.OpenAIErrorEnvelope" + } + }, + "404": { + "description": "Not Found", + "schema": { + "$ref": "#/definitions/core.OpenAIErrorEnvelope" + } + }, + "429": { + "description": "Too Many Requests", + "schema": { + "$ref": "#/definitions/core.OpenAIErrorEnvelope" + } + }, + "502": { + "description": "Bad Gateway", + "schema": { + "$ref": "#/definitions/core.OpenAIErrorEnvelope" + } + } + }, + "security": [ + { + "BearerAuth": [] + } + ] + } + }, "/v1/usage": { "get": { "description": "Returns recorded usage, budget statuses, and rate limit statuses for the caller's effective user path (the path bound to the managed API key, or the user-path header for master-key callers).", diff --git a/config/config.example.yaml b/config/config.example.yaml index d1213070e..1a177aaa6 100644 --- a/config/config.example.yaml +++ b/config/config.example.yaml @@ -598,6 +598,7 @@ providers: # server speaks the same API without authentication: set base_url # (e.g. "http://localhost:8009") and omit api_key. Name it "kev" to see # that name in logs and usage; no separate provider type is needed. + # Pinned versions such as "jev-1.13.0" route here without being listed. # Jev is priced per input token and is not in the upstream model catalog; # declare its pricing here to have the gateway cost System One requests. # models: diff --git a/docs/advanced/api-endpoints.mdx b/docs/advanced/api-endpoints.mdx index f39cc87da..ff9b452fc 100644 --- a/docs/advanced/api-endpoints.mdx +++ b/docs/advanced/api-endpoints.mdx @@ -12,8 +12,8 @@ documented separately in [Admin Endpoints](/advanced/admin-endpoints). For request and response details, see the dedicated guides: [Responses API](/advanced/responses-api), [Conversations API](/advanced/conversations-api), [Anthropic Messages API](/advanced/anthropic-messages-api), -[Audio API](/advanced/audio-api), [Images API](/advanced/images-api), and -[Usage API](/advanced/usage-api). +[Audio API](/advanced/audio-api), [Images API](/advanced/images-api), +[System One API](/advanced/systemone-api), and [Usage API](/advanced/usage-api). ## OpenAI-Compatible API @@ -111,6 +111,17 @@ metered. | `/v1/messages` | POST | Anthropic Messages API through translated model routing (streaming supported) | | `/v1/messages/count_tokens` | POST | Heuristic Anthropic Messages input token estimate | +## System One API + +Available when a `jev` or `openrouter` provider is configured; see +[System One API](/advanced/systemone-api). + +| Endpoint | Method | Description | +| ---------------------------- | ------ | ------------------------------------------------------------------------ | +| `/v1/systemone` | POST | Evaluate a state against typed questions (Jev, Kev), forwarded natively | +| `/v1/systemone/permute` | POST | Kev only: run one Choice question with several option orders | +| `/v1/systemone/separate` | POST | Kev only: run each question in its own forward pass | + ## Gateway Extensions | Endpoint | Method | Description | diff --git a/docs/advanced/systemone-api.mdx b/docs/advanced/systemone-api.mdx new file mode 100644 index 000000000..8e765fe86 --- /dev/null +++ b/docs/advanced/systemone-api.mdx @@ -0,0 +1,179 @@ +--- +title: "System One API" +description: "Send TypeSafe System One decision requests (Jev, Kev) through GoModel, forwarded natively with virtual models, guardrails, caching, failover, audit, and usage." +icon: "scale" +keywords: ["System One", "systemone", "Jev", "Kev", "TypeSafe", "decision model", "noul", "choice", "score", "OpenRouter"] +--- + +`POST /v1/systemone` serves TypeSafe's System One API: a request carries a +`state` (the text or record to evaluate) and a map of typed questions, and the +answer is a calibrated probability per question. It is a decision API, not a +text generator, so GoModel forwards it **natively** and never translates it to +or from chat. + +The endpoint is available once a [`jev` provider](/providers/jev) (hosted Jev +or a self-hosted Kev server) or an `openrouter` provider is configured. Without +one, it answers `404`. + +## Request and answer + +```bash +curl -s http://localhost:8080/v1/systemone \ + -H "Authorization: Bearer $GOMODEL_MASTER_KEY" \ + -H "Content-Type: application/json" \ + -d '{ + "model": "jev-latest", + "state": "Shoes arrived two weeks late and in the wrong size. Also I see two charges on my card.", + "questions": { + "department": {"type": "choice", "instructions": "Which team should handle this?", + "criteria": {"returns": "Exchanges, refunds, wrong items", + "shipping": "Delivery status, delays", + "billing": "Charges, invoices"}}, + "escalate": {"type": "noul", "instructions": "Does this need urgent human attention?"}, + "frustration": {"type": "score", "instructions": "How frustrated is the customer?", + "criteria": ["Calm", "Frustrated", "Very angry"]} + } + }' +``` + +The answer is the provider's own, relayed unchanged: + +```json +{ + "model": "jev-1.13.0", + "answers": { + "department": {"type": "choice", "choice": "returns", "confidence": 0.21, + "probabilities": {"returns": 0.47, "shipping": 0.28, "billing": 0.25}}, + "escalate": {"type": "noul", "noul": 0.93}, + "frustration": {"type": "score", "score": 1.44, "confidence": 0.78, + "legend": {"0": "Calm", "1": "Frustrated", "2": "Very angry"}, + "probabilities": {"0": 0.00, "1": 0.56, "2": 0.44}} + }, + "usage": {"input_tokens": 101, "output_tokens": 161} +} +``` + +The TypeSafe SDKs send `POST {base_url}/v1/systemone`, so point them at the +gateway root (`base_url="http://localhost:8080"`) with your GoModel key; see +[Jev / Kev](/providers/jev#using-the-typesafe-sdks). + +## Routes + +| Route | What it does | +| --- | --- | +| `POST /v1/systemone` | Evaluate a state against a map of questions | +| `POST /v1/systemone/permute` | Kev only: run one Choice question with several option orders (`n_perm`, 1 to 64, default 6) | +| `POST /v1/systemone/separate` | Kev only: run each question in its own forward pass | + +The Kev routes behave like `/v1/systemone`. They are refused for OpenRouter, +which answers only the evaluation route; a hosted TypeSafe `jev` provider +returns its own `404` for them. + +## What the gateway does + +1. Resolves `model` like any other endpoint: a bare name, a provider-qualified + name (`jev/jev-latest`), or a [virtual model](/features/virtual-models), + then applies the caller's [model allowlist](/features/users), + [rate limits](/features/rate-limits), and [budgets](/features/budgets). +2. Runs the workflow's prompt [guardrails](/advanced/guardrails) over `state`. +3. Serves an identical earlier request from the [response cache](/features/cache). +4. Forwards the body with only `model` (the resolved name) and `state` (if a + guardrail edited it) changed. Questions, criteria, and every other field + reach the provider byte for byte. +5. Relays the answer unchanged and records it in the audit log (request type + **System One**) and in usage. + +## Models + +| Provider | Model names | +| --- | --- | +| `jev` (hosted) | `jev/jev-latest`, `jev/jev-preview`, and any versioned ID such as `jev/jev-1.13.0` | +| `jev` (Kev server) | `kev/kev-latest` and the checkpoint's aliases, for a provider named `kev` | +| `openrouter` | `openrouter/typesafe/jev-1.13`, `openrouter/~typesafe/jev-latest`, and OpenRouter's other decision models, such as `openrouter/jaredpalmer/kev-4b` | + +System One models are listed in `GET /v1/models` as utility models with no +generation mode. + +TypeSafe lists only its aliases but accepts any versioned ID, so a pinned +version works without being declared: GoModel routes a model it does not list +to a `jev` provider when the name says which one (`jev/jev-1.13.0`), or, for a +bare name, when exactly one `jev` provider is configured. A virtual model can +pin a version the same way. + +OpenRouter accepts `jev-latest` itself, but GoModel routes on its catalog IDs. +To keep a plain `jev-latest` (the TypeSafe SDKs' default) working through +OpenRouter, add a virtual model: + +```yaml +virtual_models: + - source: jev-latest + target: openrouter/~typesafe/jev-latest +``` + +## Caching + +With the [response cache](/features/cache) enabled, an identical request (same +route, resolved model, guardrails, and body after guardrail edits) is answered +from the exact cache (`X-Cache: HIT (exact)`) and recorded in usage as a cache +hit. The semantic cache never serves System One: a state that is merely +similar is not the same decision. Send `Cache-Control: no-cache` to skip the +cache for one request. + +## Failover + +A virtual model with the `failover` strategy moves a request to its next +target when the current one fails with an availability error (`429` or `5xx`, +including TypeSafe's `529`, by default; see [Failover](/features/failover)). +Every target receives the request in its own System One form. A target without +the API, such as a chat model, is skipped without using a failover attempt, +and client errors such as a malformed question (`422`) are returned without +failover: + +```yaml +virtual_models: + - source: decider + strategy: failover + targets: + - { model: kev/kev-latest } # local Kev first + - { model: openrouter/typesafe/jev-1.13 } # hosted Jev when Kev is down +``` + +The audit log shows each attempt, usage is recorded under the target that +answered, and a failover answer is not cached. + +## Guardrails + +Guardrails see `state` as a single user message: a string state as its text, +any other JSON value as its encoded JSON, which must still be valid JSON after +an edit. That is what anonymizing and blocking guardrails need; for example, a +`string_replace` rule that masks card numbers applies to `state` before it +leaves the gateway. The questions are your application's fixed schema and are +not exposed. + +Edits a decision request has no place for, such as a system prompt injected by +a guardrail that also covers chat models, are dropped. The gateway logs one +warning per kind of dropped edit, then logs repeats at debug level. A +guardrail that would answer the request itself blocks it instead, since System +One callers expect typed answers, not text. + +## Errors and misuse + +The endpoint never translates, and it says so when a request cannot work: + +| Situation | Result | +| --- | --- | +| No `jev` or `openrouter` provider configured | `404` | +| `model` missing | `400` | +| Model on a provider without System One, or a chat, embedding, or other generation model (including through a virtual model) | `400 invalid_request_error` explaining why; the gateway logs a warning | +| A System One model sent to `/v1/chat/completions`, `/v1/responses`, or `/v1/embeddings` | `400 invalid_request_error` pointing at `/v1/systemone` | +| Upstream error, such as a malformed question | The provider's status, with its message | + +## Audit, usage, and cost + +Each call is an audit entry under its route, with the requested and resolved +model, provider, request and response bodies, guardrail outcomes, and failover +attempts; filter the audit log by the **System One** request type. Usage +records the answer's `input_tokens` and `output_tokens` under the model that +answered. OpenRouter reports its own `usage.cost`, which is recorded as the +request's cost; for hosted Jev, declare pricing on the provider (see +[Jev / Kev](/providers/jev#models-access-control-and-cost)). diff --git a/docs/docs.json b/docs/docs.json index f01997512..511eb42e5 100644 --- a/docs/docs.json +++ b/docs/docs.json @@ -143,6 +143,7 @@ "advanced/responses-compatibility", "advanced/conversations-api", "advanced/anthropic-messages-api", + "advanced/systemone-api", "advanced/extra-content", "advanced/audio-api", "advanced/images-api", diff --git a/docs/features/cache.mdx b/docs/features/cache.mdx index 75bd9f4d3..512773685 100644 --- a/docs/features/cache.mdx +++ b/docs/features/cache.mdx @@ -14,6 +14,7 @@ requests on: - `/v1/responses` - `/v1/messages` - `/v1/embeddings` +- `/v1/systemone` (and Kev's `/permute` and `/separate`) Streaming and non-streaming variants of the same request are cached independently: a streaming miss stores the raw SSE bytes and a streaming hit @@ -37,7 +38,8 @@ X-Cache: HIT (semantic) represent the exact text it was requested for, so replaying the vector of a merely similar input would be a wrong answer rather than an equivalent one. Embeddings requests are also never streamed, so only the JSON response is - cached. + cached. The same holds for [System One](/advanced/systemone-api#caching) + decisions: a similar state is not the same decision. ## Enable the exact cache diff --git a/docs/openapi.json b/docs/openapi.json index efa2de15f..a6cb55a95 100644 --- a/docs/openapi.json +++ b/docs/openapi.json @@ -11200,6 +11200,92 @@ "systemone" ], "summary": "Evaluate a System One decision request (Jev / Kev)", + "requestBody": { + "$ref": "#/components/requestBodies/Request2" + }, + "responses": { + "200": { + "description": "System One answers, in the provider's shape", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + }, + "400": { + "description": "Bad Request", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/core.OpenAIErrorEnvelope" + } + } + } + }, + "401": { + "description": "Unauthorized", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/core.OpenAIErrorEnvelope" + } + } + } + }, + "404": { + "description": "Not Found", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/core.OpenAIErrorEnvelope" + } + } + } + }, + "429": { + "description": "Too Many Requests", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/core.OpenAIErrorEnvelope" + } + } + } + }, + "502": { + "description": "Bad Gateway", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/core.OpenAIErrorEnvelope" + } + } + } + } + }, + "security": [ + { + "BearerAuth": [] + } + ], + "x-mint": { + "metadata": { + "sidebarTitle": "/v1/systemone", + "title": "Evaluate a System One decision request (Jev / Kev)", + "description": "GoModel API reference for POST /v1/systemone: Evaluate a System One decision request (Jev / Kev)." + } + } + } + }, + "/v1/systemone/permute": { + "post": { + "description": "A Kev server diagnostic: the request is a System One request, and n_perm (1 to 64, default 6) sets how many option orders run. Only jev providers pointing at a Kev server serve it.", + "tags": [ + "systemone" + ], + "summary": "Run one Choice question with several option orders (Kev)", "requestBody": { "content": { "application/json": { @@ -11208,12 +11294,12 @@ } } }, - "description": "System One request: model, state, and questions", + "description": "System One request with one Choice question", "required": true }, "responses": { "200": { - "description": "System One answers, in the provider's shape", + "description": "Kev's answer, in the provider's shape", "content": { "application/json": { "schema": { @@ -11280,9 +11366,95 @@ ], "x-mint": { "metadata": { - "sidebarTitle": "/v1/systemone", - "title": "Evaluate a System One decision request (Jev / Kev)", - "description": "GoModel API reference for POST /v1/systemone: Evaluate a System One decision request (Jev / Kev)." + "sidebarTitle": "/v1/systemone/permute", + "title": "Run one Choice question with several option orders (Kev)", + "description": "GoModel API reference for POST /v1/systemone/permute: Run one Choice question with several option orders (Kev)." + } + } + } + }, + "/v1/systemone/separate": { + "post": { + "description": "A Kev server diagnostic that answers each question separately. Only jev providers pointing at a Kev server serve it.", + "tags": [ + "systemone" + ], + "summary": "Run each System One question in its own forward pass (Kev)", + "requestBody": { + "$ref": "#/components/requestBodies/Request2" + }, + "responses": { + "200": { + "description": "Kev's answer, in the provider's shape", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + }, + "400": { + "description": "Bad Request", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/core.OpenAIErrorEnvelope" + } + } + } + }, + "401": { + "description": "Unauthorized", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/core.OpenAIErrorEnvelope" + } + } + } + }, + "404": { + "description": "Not Found", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/core.OpenAIErrorEnvelope" + } + } + } + }, + "429": { + "description": "Too Many Requests", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/core.OpenAIErrorEnvelope" + } + } + } + }, + "502": { + "description": "Bad Gateway", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/core.OpenAIErrorEnvelope" + } + } + } + } + }, + "security": [ + { + "BearerAuth": [] + } + ], + "x-mint": { + "metadata": { + "sidebarTitle": "/v1/systemone/separate", + "title": "Run each System One question in its own forward pass (Kev)", + "description": "GoModel API reference for POST /v1/systemone/separate: Run each System One question in its own forward pass (Kev)." } } } @@ -11452,6 +11624,17 @@ }, "description": "Anthropic Messages request", "required": true + }, + "Request2": { + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + }, + "description": "System One request: model, state, and questions", + "required": true } }, "securitySchemes": { diff --git a/docs/providers/jev.mdx b/docs/providers/jev.mdx index ca23eabce..c983f5485 100644 --- a/docs/providers/jev.mdx +++ b/docs/providers/jev.mdx @@ -132,75 +132,23 @@ providers, name the model with its provider (`kev/kev-latest`) or a The SDKs' model listing expects TypeSafe's shape, while the gateway's `/v1/models` is OpenAI-shaped; list upstream models at `/p/jev/v1/models`. -## The native endpoint +## System One API -`POST /v1/systemone` is a gateway endpoint, not a raw proxy. For each request -GoModel: +`POST /v1/systemone` is a gateway endpoint, not a raw proxy: it applies +virtual models, guardrails on `state`, the response cache, failover, audit, +and usage, and it forwards the request natively without translating it. Kev's +`/v1/systemone/permute` and `/v1/systemone/separate` work the same way. See +[System One API](/advanced/systemone-api) for the full behavior, including +OpenRouter, which serves Jev natively too. -1. Resolves `model` like any other endpoint: a bare name, a provider-qualified - name (`jev/jev-latest`), or a virtual model, then applies the caller's - model allowlist, rate limits, and budgets. -2. Runs the workflow's prompt [guardrails](/advanced/guardrails) over `state`. -3. Forwards the body with only `model` (to the resolved name) and `state` (if - a guardrail edited it) changed. Questions, criteria, and every other field - reach the provider byte for byte. -4. Relays the answer unchanged, and records it in the audit log (as a - **System One** request) and in usage. - -The endpoint never translates. A model that cannot answer System One fails -with `400 invalid_request_error` explaining why, and the gateway logs a -warning: one on a provider without the API, or one the catalog lists as a -chat, embedding, or other generation model, such as a virtual model pointing -at a chat model. Without a `jev` or `openrouter` provider, the route answers -`404`. - -Response caching and failover do not apply to this endpoint yet. - -### Through OpenRouter - -[OpenRouter serves Jev natively](https://openrouter.ai/docs/guides/community/jev) -at the same path, so an OpenRouter key alone is enough: - -```bash -OPENROUTER_API_KEY=sk-or-... -``` - -OpenRouter's decision models appear in `GET /v1/models` as utility models, -priced from OpenRouter's listing: `openrouter/typesafe/jev-1.13`, -`openrouter/~typesafe/jev-latest` (tracks the newest Jev), and other decision -models such as Kev 4B (`openrouter/jaredpalmer/kev-4b`). Name them that way in -`model`. OpenRouter accepts `jev-latest` itself, but GoModel routes on its -catalog IDs, so to keep the TypeSafe SDK's plain `jev-latest` working, add a -[virtual model](/features/virtual-models) `jev-latest` that targets -`openrouter/~typesafe/jev-latest`. Chat models on the same provider are -rejected, since they have no System One API. - -The answer carries OpenRouter's `id`, `provider`, and `usage.cost`, and GoModel -records that reported cost with the request's usage. With a `jev` provider -configured as well, one virtual model can front a local Kev server and -OpenRouter's Jev together. - -### Guardrails - -Guardrails see `state` as a single user message: a string state as its text, -any other JSON value as its encoded JSON (which must still be valid JSON after -an edit). This is what anonymizing or blocking guardrails need. The questions -are your application's fixed schema and are not exposed. - -Guardrail edits a decision request has no place for, such as a system prompt -injected by a workflow that also covers chat models, are dropped with a -warning in the logs rather than failing the request. A guardrail that would -answer the request itself blocks it instead, since System One callers expect -typed answers, not text. - -## Native routes +Every other upstream route is reachable through +[passthrough](/features/passthrough-api), without virtual models, guardrails, +or caching: | Route | What it does | | --- | --- | -| `POST /p/jev/v1/systemone` | Evaluate a state against a map of questions, without virtual models or guardrails | | `GET /p/jev/v1/models` | The names the `model` field accepts, in the upstream's own shape | -| `POST /p/jev/v1/systemone/permute` | Kev only: run one Choice question with several option orders | -| `POST /p/jev/v1/systemone/separate` | Kev only: run each question in its own forward pass | +| `POST /p/jev/v1/systemone` | The evaluation route, forwarded as sent | Upstream errors keep their status code, with the provider's body carried in the gateway error message: a malformed question comes back as TypeSafe's `422` @@ -219,10 +167,8 @@ Every System One request names its model, so both `/v1/systemone` and the passthrough surface apply the caller's [model allowlist](/features/users) to it like any other request. -`/v1/systemone` routes only to models in the catalog. To pin a version the -upstream does not list, such as `jev-1.13.0`, declare it under the provider's -`models` (as in the pricing example below) and set -`CONFIGURED_PROVIDER_MODELS_MODE=merge`; passthrough accepts any name. +A pinned version such as `jev-1.13.0` works without being declared; see +[Models](/advanced/systemone-api#models) for how unlisted names are routed. The response's `usage.input_tokens` and `usage.output_tokens` are recorded, so System One calls appear in the usage API and dashboard under the model that diff --git a/internal/core/endpoint_operations.go b/internal/core/endpoint_operations.go index 8bb2b1395..26ac74342 100644 --- a/internal/core/endpoint_operations.go +++ b/internal/core/endpoint_operations.go @@ -29,7 +29,7 @@ var operationPaths = map[Operation]OperationPaths{ "/v1/realtime/translations", "/v1/realtime/translations/calls", "/v1/realtime/translations/client_secrets", }}, OperationMCP: {Prefixes: []string{"/mcp"}}, - OperationSystemOne: {Exact: []string{"/v1/systemone"}}, + OperationSystemOne: {Exact: []string{"/v1/systemone", "/v1/systemone/permute", "/v1/systemone/separate"}}, OperationProviderPassthrough: {Prefixes: []string{"/p"}}, } diff --git a/internal/core/endpoints.go b/internal/core/endpoints.go index 386d3c726..0e2edc494 100644 --- a/internal/core/endpoints.go +++ b/internal/core/endpoints.go @@ -170,10 +170,11 @@ func describeEndpointPath(path string) EndpointDescriptor { Dialect: "openai_compat", Operation: OperationImageEdits, } - case path == "/v1/systemone": - // TypeSafe's System One decision API (Jev, Kev). It has no canonical - // translation: the body is forwarded to a System One provider - // unchanged, apart from the routed model and guardrail edits to state. + case path == "/v1/systemone" || path == "/v1/systemone/permute" || path == "/v1/systemone/separate": + // TypeSafe's System One decision API (Jev, Kev) and the diagnostic + // variants Kev servers add. It has no canonical translation: the body + // is forwarded to a System One provider unchanged, apart from the + // routed model and guardrail edits to state. return EndpointDescriptor{ ModelInteraction: true, IngressManaged: true, diff --git a/internal/core/endpoints_test.go b/internal/core/endpoints_test.go index b881965de..3fd2b0046 100644 --- a/internal/core/endpoints_test.go +++ b/internal/core/endpoints_test.go @@ -43,7 +43,9 @@ func TestDescribeEndpointPath(t *testing.T) { {path: "/mcp", managed: false, dialect: "mcp", operation: OperationMCP, bodyMode: BodyModeNone, interaction: true}, {path: "/mcp/linear", managed: false, dialect: "mcp", operation: OperationMCP, bodyMode: BodyModeNone, interaction: true}, {path: "/v1/systemone", managed: true, dialect: "systemone", operation: OperationSystemOne, bodyMode: BodyModeJSON, interaction: true}, - {path: "/v1/systemone/permute", managed: false, dialect: "", operation: "", bodyMode: BodyModeNone, interaction: false}, + {path: "/v1/systemone/permute", managed: true, dialect: "systemone", operation: OperationSystemOne, bodyMode: BodyModeJSON, interaction: true}, + {path: "/v1/systemone/separate", managed: true, dialect: "systemone", operation: OperationSystemOne, bodyMode: BodyModeJSON, interaction: true}, + {path: "/v1/systemone/other", managed: false, dialect: "", operation: "", bodyMode: BodyModeNone, interaction: false}, {path: "/p/openai/responses", managed: true, dialect: "provider_passthrough", operation: OperationProviderPassthrough, bodyMode: BodyModeOpaque, interaction: true}, {path: "/v1/models", managed: false, dialect: "", operation: "", bodyMode: BodyModeNone, interaction: false}, } diff --git a/internal/gateway/failover.go b/internal/gateway/failover.go index 370db8477..f20951283 100644 --- a/internal/gateway/failover.go +++ b/internal/gateway/failover.go @@ -42,6 +42,7 @@ func tryFailoverResponse[T any]( workflow *core.Workflow, model, provider string, primaryErr error, + eligible func(selector core.ModelSelector, providerType string) bool, call func(selector core.ModelSelector, providerType, providerName string) (T, string, error), ) (T, ExecutionMeta, error) { var zero T @@ -75,6 +76,16 @@ func tryFailoverResponse[T any]( qualified := selector.QualifiedModel() providerType := o.ProviderTypeForSelector(selector, ProviderTypeFromWorkflow(workflow)) providerName := ResolvedProviderName(o.provider, selector, ProviderNameFromWorkflow(workflow)) + // A target that cannot serve the request is skipped before it counts + // against the attempt cap, so it never crowds out a later valid one. + if eligible != nil && !eligible(selector, providerType) { + slog.Info("skipping failover target that cannot serve the request", + "request_id", requestID, + "to", qualified, + "provider_type", providerType, + ) + continue + } if o.routeGate != nil && !o.routeGate.RouteAvailable(providerName, qualified) { slog.Info("skipping rate-limited failover target", "request_id", requestID, @@ -121,13 +132,14 @@ func executeWithFailoverResponse[T any]( workflow *core.Workflow, model, provider string, primary func() (T, string, string, error), + eligible func(selector core.ModelSelector, providerType string) bool, failoverFn func(selector core.ModelSelector, providerType, providerName string) (T, string, error), ) (T, ExecutionMeta, error) { resp, resolvedProviderType, resolvedProviderName, err := primary() if err == nil { return resp, ExecutionMeta{ProviderType: resolvedProviderType, ProviderName: resolvedProviderName}, nil } - return tryFailoverResponse(ctx, o, workflow, model, provider, err, failoverFn) + return tryFailoverResponse(ctx, o, workflow, model, provider, err, eligible, failoverFn) } func executeTranslatedWithFailover[Req any, Resp any]( @@ -158,6 +170,7 @@ func executeTranslatedWithFailover[Req any, Resp any]( } return resp, ResponseProviderType(ProviderTypeFromWorkflow(workflow), responseProvider), ProviderNameFromWorkflow(workflow), nil }, + nil, func(selector core.ModelSelector, providerType, providerName string) (Resp, string, error) { // A failover target gets a different request body, so it must not // reuse the client's idempotency key. @@ -255,3 +268,57 @@ func firstNonEmptyString(values ...string) string { } return "" } + +// PassthroughCall sends one native request to selector's provider. Provider +// error statuses must come back as errors so the failover policy can judge +// them. +type PassthroughCall func(ctx context.Context, selector core.ModelSelector, providerType, providerName string) (*core.PassthroughResponse, error) + +// ExecutePassthroughWithFailover runs a native, untranslated request against +// the workflow's resolved route and then, while the failover policy allows, +// against its failover targets. It is the native-endpoint counterpart of the +// translated failover path: attempts are recorded the same way, but every +// target receives the client's own dialect, so eligible must reject a +// failover target that cannot serve it; rejected targets are skipped without +// counting against the attempt cap. The selector that answered is returned +// with the response. +func (o *InferenceOrchestrator) ExecutePassthroughWithFailover(ctx context.Context, workflow *core.Workflow, eligible func(selector core.ModelSelector, providerType string) bool, call PassthroughCall) (*core.PassthroughResponse, core.ModelSelector, ExecutionMeta, error) { + primary := core.ModelSelector{} + if workflow != nil && workflow.Resolution != nil { + primary = workflow.Resolution.ResolvedSelector + } + type answer struct { + resp *core.PassthroughResponse + selector core.ModelSelector + } + result, meta, err := executeWithFailoverResponse(ctx, o, workflow, primary.Model, primary.Provider, + func() (answer, string, string, error) { + started := time.Now() + providerType, providerName := ProviderTypeFromWorkflow(workflow), ProviderNameFromWorkflow(workflow) + qualified := primary.QualifiedModel() + // A rate-saturated primary route must not reach the provider; its + // stored 429 becomes the primary failure that starts the sweep. + if saturated := core.PrimaryRouteSaturated(ctx); saturated != nil { + recordProviderAttempt(ctx, providerAttemptFromResult(AttemptKindPrimary, providerType, providerName, qualified, started, saturated)) + return answer{}, "", "", saturated + } + resp, err := call(ctx, primary, providerType, providerName) + recordProviderAttempt(ctx, providerAttemptFromResult(AttemptKindPrimary, providerType, providerName, qualified, started, err)) + if err != nil { + return answer{}, "", "", err + } + return answer{resp: resp, selector: primary}, providerType, providerName, nil + }, + eligible, + func(selector core.ModelSelector, providerType, providerName string) (answer, string, error) { + // A failover target gets a different body, so it must not reuse + // the client's idempotency key. + resp, err := call(core.WithIdempotencyKey(ctx, ""), selector, providerType, providerName) + if err != nil { + return answer{}, "", err + } + return answer{resp: resp, selector: selector}, providerType, nil + }, + ) + return result.resp, result.selector, meta, err +} diff --git a/internal/gateway/failover_policy_test.go b/internal/gateway/failover_policy_test.go index fdf4f9176..cdb96d676 100644 --- a/internal/gateway/failover_policy_test.go +++ b/internal/gateway/failover_policy_test.go @@ -92,7 +92,7 @@ func TestTryFailoverResponseHonorsMaxAttempts(t *testing.T) { return "", "", core.NewProviderError("openai", http.StatusBadGateway, selector.Model+" down", nil) } - _, meta, err := tryFailoverResponse(context.Background(), o, workflow, "openai/gpt-4o", "openai", primaryErr, call) + _, meta, err := tryFailoverResponse(context.Background(), o, workflow, "openai/gpt-4o", "openai", primaryErr, nil, call) require.False(t, meta.UsedFailover) require.Error(t, err) @@ -104,6 +104,25 @@ func TestTryFailoverResponseHonorsMaxAttempts(t *testing.T) { } } +// Targets that cannot serve the request (a chat model in a System One chain) +// are skipped before a call, so they do not consume attempts either. +func TestTryFailoverResponseIneligibleTargetsDoNotConsumeAttempts(t *testing.T) { + o, workflow := threeTargetFixture(&FailoverPolicy{MaxAttempts: 1}) + primaryErr := core.NewProviderError("openai", http.StatusBadGateway, "primary down", nil) + eligible := func(selector core.ModelSelector, _ string) bool { return selector.Model != "a" } + var calls []string + call := func(selector core.ModelSelector, _, _ string) (string, string, error) { + calls = append(calls, selector.QualifiedModel()) + return "ok", "openai", nil + } + + _, meta, err := tryFailoverResponse(context.Background(), o, workflow, "openai/gpt-4o", "openai", primaryErr, eligible, call) + + require.NoError(t, err) + require.True(t, meta.UsedFailover) + require.Equal(t, []string{"openai/b"}, calls) +} + // Targets skipped before a call (rate-limited routes) do not consume attempts. func TestTryFailoverResponseMaxAttemptsCountsCallsOnly(t *testing.T) { o, workflow := threeTargetFixture(&FailoverPolicy{MaxAttempts: 1}) @@ -115,7 +134,7 @@ func TestTryFailoverResponseMaxAttemptsCountsCallsOnly(t *testing.T) { return "ok", "openai", nil } - _, meta, err := tryFailoverResponse(context.Background(), o, workflow, "openai/gpt-4o", "openai", primaryErr, call) + _, meta, err := tryFailoverResponse(context.Background(), o, workflow, "openai/gpt-4o", "openai", primaryErr, nil, call) require.True(t, meta.UsedFailover) require.NoError(t, err) @@ -150,7 +169,7 @@ func TestTryFailoverResponseSkipsWhenPolicyDoesNotMatch(t *testing.T) { return "ok", "openai", nil } - _, meta, err := tryFailoverResponse(context.Background(), o, workflow, "openai/gpt-4o", "openai", primaryErr, call) + _, meta, err := tryFailoverResponse(context.Background(), o, workflow, "openai/gpt-4o", "openai", primaryErr, nil, call) require.False(t, called) require.False(t, meta.UsedFailover) diff --git a/internal/gateway/failover_test.go b/internal/gateway/failover_test.go index db345517e..f6ac6e541 100644 --- a/internal/gateway/failover_test.go +++ b/internal/gateway/failover_test.go @@ -48,7 +48,7 @@ func TestTryFailoverResponseSkipsWhenContextCanceled(t *testing.T) { return "", "", core.NewProviderError("openai", http.StatusBadGateway, "unexpected failover call", nil) } - _, meta, err := tryFailoverResponse(ctx, o, workflow, "openai/gpt-4o", "openai", primaryErr, call) + _, meta, err := tryFailoverResponse(ctx, o, workflow, "openai/gpt-4o", "openai", primaryErr, nil, call) require.False(t, called) require.False(t, meta.UsedFailover) @@ -66,7 +66,7 @@ func TestTryFailoverResponseAttemptsWhenContextLive(t *testing.T) { return "ok", "openai", nil } - resp, meta, err := tryFailoverResponse(context.Background(), o, workflow, "openai/gpt-4o", "openai", primaryErr, call) + resp, meta, err := tryFailoverResponse(context.Background(), o, workflow, "openai/gpt-4o", "openai", primaryErr, nil, call) require.True(t, called) require.True(t, meta.UsedFailover) @@ -100,7 +100,7 @@ func TestTryFailoverResponseSkipsRateLimitedTargets(t *testing.T) { return "ok", "anthropic", nil } - resp, meta, err := tryFailoverResponse(context.Background(), o, workflow, "openai/gpt-4o", "openai", primaryErr, call) + resp, meta, err := tryFailoverResponse(context.Background(), o, workflow, "openai/gpt-4o", "openai", primaryErr, nil, call) require.Len(t, attempted, 1) require.Equal(t, "anthropic/claude", attempted[0]) diff --git a/internal/guardrails/workflow_executor.go b/internal/guardrails/workflow_executor.go index 82783ec31..01920a1bc 100644 --- a/internal/guardrails/workflow_executor.go +++ b/internal/guardrails/workflow_executor.go @@ -3,6 +3,8 @@ package guardrails import ( "context" "log/slog" + "strings" + "sync" "github.com/enterpilot/gomodel/internal/core" "github.com/enterpilot/gomodel/internal/plugins" @@ -49,9 +51,20 @@ func (p *WorkflowRequestPatcher) PatchSystemOneRequest(ctx context.Context, req return processGuarded(ctx, p.chain(ctx), req, "System One", exchange.FromSystemOneRequest, applySystemOneEdits) } +// systemOneDropsWarned records the kinds of dropped guardrail edits already +// logged at warning level. A guardrail scoped to every model edits every +// System One request the same way, so the misconfiguration is reported once +// per kind; repeats are logged at debug level. +var systemOneDropsWarned sync.Map + func applySystemOneEdits(req *core.SystemOneRequest, prompt *pluginapi.Prompt) (*core.SystemOneRequest, error) { if dropped := exchange.SystemOneUncarriedEdits(prompt); len(dropped) > 0 { - slog.Warn("guardrail edits a System One request cannot carry were dropped; only the state is guarded", "model", req.Model, "dropped", dropped) + const message = "guardrail edits a System One request cannot carry were dropped; only the state is guarded" + if _, repeated := systemOneDropsWarned.LoadOrStore(strings.Join(dropped, ","), struct{}{}); repeated { + slog.Debug(message, "model", req.Model, "dropped", dropped) + } else { + slog.Warn(message+" (repeats are logged at debug level)", "model", req.Model, "dropped", dropped) + } } return exchange.ApplyToSystemOneRequest(req, prompt) } diff --git a/internal/plugins/exchange/systemone_request.go b/internal/plugins/exchange/systemone_request.go index 73ad3ad33..eaef9b685 100644 --- a/internal/plugins/exchange/systemone_request.go +++ b/internal/plugins/exchange/systemone_request.go @@ -86,21 +86,31 @@ func ApplyToSystemOneRequest(original *core.SystemOneRequest, p *pluginapi.Promp return &result, nil } -// SystemOneUncarriedEdits describes the prompt edits ApplyToSystemOneRequest -// does not apply, in a stable order, or nil when every edit was carried. +// SystemOneUncarriedEdits describes the kinds of prompt edits +// ApplyToSystemOneRequest does not apply ("inserted message", parameter +// "temperature"), deduplicated and sorted, or nil when every edit was +// carried. Message IDs are left out: they differ per request and name +// nothing an operator can act on. func SystemOneUncarriedEdits(p *pluginapi.Prompt) []string { if p == nil { return nil } changes := p.Changes() - var uncarried []string + seen := map[string]struct{}{} for id, kind := range changes.Messages { if id != SystemOneStateMessageID { - uncarried = append(uncarried, fmt.Sprintf("%s message %q", kind, id)) + seen[string(kind)+" message"] = struct{}{} } } for name := range changes.Params { - uncarried = append(uncarried, fmt.Sprintf("parameter %q", name)) + seen[fmt.Sprintf("parameter %q", name)] = struct{}{} + } + if len(seen) == 0 { + return nil + } + uncarried := make([]string, 0, len(seen)) + for description := range seen { + uncarried = append(uncarried, description) } sort.Strings(uncarried) return uncarried diff --git a/internal/plugins/exchange/systemone_request_test.go b/internal/plugins/exchange/systemone_request_test.go index 62d5ca08b..4cf0d4078 100644 --- a/internal/plugins/exchange/systemone_request_test.go +++ b/internal/plugins/exchange/systemone_request_test.go @@ -121,7 +121,9 @@ func TestSystemOneUncarriedEdits(t *testing.T) { require.NoError(t, p.SetText(SystemOneStateMessageID, 0, "[PERSON]")) assert.Nil(t, SystemOneUncarriedEdits(p), "a state edit is carried") - id := p.Insert(0, pluginapi.TextMessage(pluginapi.RoleSystem, "be safe")) + p.Insert(0, pluginapi.TextMessage(pluginapi.RoleSystem, "be safe")) + p.Append(pluginapi.TextMessage(pluginapi.RoleSystem, "be brief")) p.SetParam("temperature", 0.1) - assert.Equal(t, []string{`inserted message "` + id + `"`, `parameter "temperature"`}, SystemOneUncarriedEdits(p)) + assert.Equal(t, []string{"inserted message", `parameter "temperature"`}, SystemOneUncarriedEdits(p), + "kinds are listed once, without per-request message IDs") } diff --git a/internal/providers/registry_normalization_test.go b/internal/providers/registry_normalization_test.go index 0af8eb32b..3d29b31d5 100644 --- a/internal/providers/registry_normalization_test.go +++ b/internal/providers/registry_normalization_test.go @@ -365,3 +365,20 @@ func TestRouterLookupModel(t *testing.T) { _, ok = (&Router{}).LookupModel("openrouter/~typesafe/jev-latest") assert.False(t, ok, "a lookup without single-model access describes nothing") } + +// ProviderNamesForType lists every configured instance of one type, so a +// caller can tell a single jev provider from several. +func TestRouterProviderNamesForType(t *testing.T) { + registry := newTestRegistryWithModels( + registryModelEntry{provider: &mockProvider{name: "kev"}, providerName: "kev", providerType: "jev", modelID: "kev-latest"}, + registryModelEntry{provider: &mockProvider{name: "jev"}, providerName: "jev", providerType: "jev", modelID: "jev-latest"}, + registryModelEntry{provider: &mockProvider{name: "openrouter"}, providerName: "openrouter", providerType: "openrouter", modelID: "typesafe/jev-1.13"}, + ) + router, err := NewRouter(registry) + require.NoError(t, err) + + assert.Equal(t, []string{"jev", "kev"}, router.ProviderNamesForType("jev")) + assert.Equal(t, []string{"openrouter"}, router.ProviderNamesForType("openrouter")) + assert.Empty(t, router.ProviderNamesForType("anthropic")) + assert.Empty(t, router.ProviderNamesForType("")) +} diff --git a/internal/providers/router_models.go b/internal/providers/router_models.go index b49d87f70..80779bb49 100644 --- a/internal/providers/router_models.go +++ b/internal/providers/router_models.go @@ -196,3 +196,20 @@ func (r *Router) LookupModel(model string) (*core.Model, bool) { cloned := info.Model return &cloned, true } + +// ProviderNamesForType lists the configured provider instance names of one +// type, sorted, or nil when the lookup cannot enumerate its providers. +func (r *Router) ProviderNamesForType(providerType string) []string { + providerType = strings.TrimSpace(providerType) + if providerType == "" || r.caps.nameLister == nil { + return nil + } + var names []string + for _, name := range r.caps.nameLister.ProviderNames() { + if r.GetProviderTypeForName(name) == providerType { + names = append(names, name) + } + } + sort.Strings(names) + return names +} diff --git a/internal/server/http.go b/internal/server/http.go index f6e4dab94..3815dff75 100644 --- a/internal/server/http.go +++ b/internal/server/http.go @@ -499,6 +499,8 @@ func New(provider core.RoutableProvider, cfg *Config) *Server { // System One decisions (Jev / Kev). The handler answers 404 until a jev // or openrouter provider is configured. e.POST("/v1/systemone", handler.SystemOne) + e.POST("/v1/systemone/permute", handler.SystemOnePermute) + e.POST("/v1/systemone/separate", handler.SystemOneSeparate) if cfg == nil || cfg.RealtimeEnabled { e.GET("/v1/realtime", handler.Realtime) e.POST("/v1/realtime/calls", handler.RealtimeCalls) diff --git a/internal/server/model_validation.go b/internal/server/model_validation.go index 52b54be42..dbcd57b09 100644 --- a/internal/server/model_validation.go +++ b/internal/server/model_validation.go @@ -103,12 +103,14 @@ func deriveWorkflowWithPolicy( } return workflow, nil - case core.OperationChatCompletions, core.OperationResponses, core.OperationEmbeddings, core.OperationSystemOne: - if desc.Operation == core.OperationSystemOne && !systemOneAvailable(provider) { - // The handler answers 404; resolving the model first would - // report a model error for an endpoint that is not there. - return nil, nil - } + case core.OperationSystemOne: + // The System One handler resolves the model itself: only it knows + // whether the endpoint is available (answering 404 before any model + // error) and when an unlisted pinned version may still route to a + // jev provider. + return nil, nil + + case core.OperationChatCompletions, core.OperationResponses, core.OperationEmbeddings: workflow.Mode = core.ExecutionModeTranslated if desc.BodyMode != core.BodyModeJSON { // Responses lifecycle routes (GET/DELETE /v1/responses/{id}, @@ -131,6 +133,9 @@ func deriveWorkflowWithPolicy( } return workflow, nil } + if systemOneOnlyModel(provider, resolution) { + return nil, systemOneOnlyModelError(desc.Operation, resolution) + } return translatedWorkflow(c.Request().Context(), requestID, desc, resolution, policyResolver) default: diff --git a/internal/server/systemone_dispatch.go b/internal/server/systemone_dispatch.go new file mode 100644 index 000000000..e9ad381a2 --- /dev/null +++ b/internal/server/systemone_dispatch.go @@ -0,0 +1,158 @@ +package server + +import ( + "bytes" + "context" + "errors" + "io" + "net/http" + "strings" + + "github.com/labstack/echo/v5" + + "github.com/enterpilot/gomodel/internal/auditlog" + "github.com/enterpilot/gomodel/internal/core" + "github.com/enterpilot/gomodel/internal/gateway" + "github.com/enterpilot/gomodel/internal/responsecache" +) + +// dispatchSystemOneWithCache serves a guarded System One body from the +// response cache when the workflow allows it, and forwards it otherwise. +// Only the exact layer applies: the key covers the route, the resolved model, +// the guardrail chain, and the forwarded body, and the semantic layer never +// serves these paths, since a similar state is not the same decision. +func (s *translatedInferenceService) dispatchSystemOneWithCache(c *echo.Context, route systemOneRoute, workflow *core.Workflow, body []byte) error { + dispatch := func() error { return s.dispatchSystemOne(c, route, workflow, body) } + if s.responseCache == nil || !workflow.CacheEnabled() { + return dispatch() + } + c.SetRequest(c.Request().WithContext(s.inference().WithCacheRequestContext(c.Request().Context(), workflow))) + err := s.responseCache.HandleRequest(c, body, dispatch) + if replayErr, ok := errors.AsType[*responsecache.ReplayError](err); ok { + recordCachedStreamError(c, replayErr.Err) + return nil + } + return err +} + +// dispatchSystemOne forwards the body to the resolved route, and to the +// workflow's failover targets while the failover policy allows, then relays +// the answer unchanged with audit and usage accounting. +func (s *translatedInferenceService) dispatchSystemOne(c *echo.Context, route systemOneRoute, workflow *core.Workflow, body []byte) error { + passthroughProvider, ok := s.provider.(core.RoutablePassthrough) + if !ok { + return handleError(c, core.NewInvalidRequestError("provider passthrough is not supported by the current provider router", nil)) + } + // Record each provider attempt so the audit entry shows a failed primary + // and the failover that answered, as it does for chat. + c.SetRequest(c.Request().WithContext(gateway.WithAttemptRecorder(c.Request().Context()))) + s.observeLiveProviderAttempts(c, workflow) + + failovers := len(s.inference().FailoverSelectors(workflow)) + adm, err := enforceAdmission(c, s.rateLimiter, s.budgetChecker, rateLimitRouteFromWorkflow(workflow).withFailovers(failovers)) + if err != nil { + return handleError(c, err) + } + defer adm.release() + ctx := adm.dispatchContext(c.Request().Context()) + + headers := buildPassthroughHeaders(ctx, c.Request().Header) + // The client's Idempotency-Key reaches the primary through the request + // context, which failover attempts clear; forwarded as an explicit header + // it would also mark every failover target's different body. + headers.Del(core.IdempotencyKeyHeader) + eligible := func(selector core.ModelSelector, providerType string) bool { + return s.systemOneUnsupportedReason(route, selector, providerType) == "" + } + resp, executed, meta, err := s.inference().ExecutePassthroughWithFailover(ctx, workflow, eligible, + func(ctx context.Context, selector core.ModelSelector, providerType, providerName string) (*core.PassthroughResponse, error) { + return s.sendSystemOne(ctx, passthroughProvider, route, selector, providerType, providerName, headers, body) + }) + enrichAuditEntryWithProviderAttempts(c) + if err != nil { + return handleError(c, err) + } + + if meta.UsedFailover { + markRequestFailoverUsed(c) + auditlog.EnrichEntryWithFailover(c, meta.FailoverModel) + workflow = executedSystemOneWorkflow(workflow, executed, meta) + storeWorkflow(c, workflow) + } + auditlog.EnrichEntryWithWorkflow(c, workflow) + auditlog.EnrichEntryWithResolvedRoute(c, executed.QualifiedModel(), meta.ProviderType, meta.ProviderName) + info := &core.PassthroughRouteInfo{ + Provider: meta.ProviderType, + ProviderName: meta.ProviderName, + NormalizedEndpoint: route.endpoint, + SemanticOperation: route.operation, + AuditPath: route.path, + Model: executed.Model, + } + return proxyPassthroughResponse(c, s.logger, s.usageLogger, s.pricingResolver, meta.ProviderType, meta.ProviderName, route.endpoint, info, resp) +} + +// maxSystemOneErrorBodyBytes caps how much of an upstream error body is read +// to build the gateway error, so a misbehaving upstream cannot make the +// gateway buffer an unbounded body. +const maxSystemOneErrorBodyBytes = 64 << 10 + +// sendSystemOne sends the body to one target under its own model name. Only +// targets that serve the route reach it: the handler checks the primary and +// the failover sweep skips ineligible targets. An upstream error status comes +// back as an error the failover policy can judge. +func (s *translatedInferenceService) sendSystemOne( + ctx context.Context, + passthroughProvider core.RoutablePassthrough, + route systemOneRoute, + selector core.ModelSelector, + providerType, providerName string, + headers http.Header, + body []byte, +) (*core.PassthroughResponse, error) { + forwarded, err := rewriteMessagesModel(body, selector.Model) + if err != nil { + return nil, core.NewInvalidRequestError("invalid request body: "+err.Error(), err) + } + resp, err := passthroughProvider.Passthrough(ctx, providerType, &core.PassthroughRequest{ + Method: http.MethodPost, + Endpoint: route.endpoint, + Operation: route.operation, + Model: selector.Model, + Body: io.NopCloser(bytes.NewReader(forwarded)), + Headers: headers.Clone(), + ProviderName: providerName, + }) + if err != nil { + return nil, err + } + if resp == nil || resp.Body == nil { + return nil, core.NewProviderError(providerType, http.StatusBadGateway, "provider returned empty passthrough response", nil) + } + if resp.StatusCode < http.StatusBadRequest { + return resp, nil + } + defer func() { _ = resp.Body.Close() }() + errorBody, err := io.ReadAll(io.LimitReader(resp.Body, maxSystemOneErrorBodyBytes)) + if err != nil { + return nil, core.NewProviderError(providerType, http.StatusBadGateway, "failed to read provider error response", err) + } + return nil, core.ParseProviderError(providerType, resp.StatusCode, errorBody, nil) +} + +// executedSystemOneWorkflow returns a copy of workflow routed to the failover +// target that answered, so usage is priced and audited under the model that +// did the work rather than the primary that failed. +func executedSystemOneWorkflow(workflow *core.Workflow, executed core.ModelSelector, meta gateway.ExecutionMeta) *core.Workflow { + if workflow == nil || workflow.Resolution == nil { + return workflow + } + resolution := *workflow.Resolution + resolution.ResolvedSelector = executed + resolution.ProviderType = strings.TrimSpace(meta.ProviderType) + resolution.ProviderName = strings.TrimSpace(meta.ProviderName) + cloned := *workflow + cloned.ProviderType = resolution.ProviderType + cloned.Resolution = &resolution + return &cloned +} diff --git a/internal/server/systemone_dispatch_test.go b/internal/server/systemone_dispatch_test.go new file mode 100644 index 000000000..106dc5466 --- /dev/null +++ b/internal/server/systemone_dispatch_test.go @@ -0,0 +1,334 @@ +package server + +import ( + "context" + "io" + "net/http" + "slices" + "strings" + "testing" + "time" + + "github.com/goccy/go-json" + "github.com/labstack/echo/v5" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/enterpilot/gomodel/internal/auditlog" + "github.com/enterpilot/gomodel/internal/cache" + "github.com/enterpilot/gomodel/internal/core" + "github.com/enterpilot/gomodel/internal/echotest" + "github.com/enterpilot/gomodel/internal/gateway" + "github.com/enterpilot/gomodel/internal/responsecache" + "github.com/enterpilot/gomodel/internal/usage" +) + +// scriptedSystemOneProvider answers each passthrough by the model the body +// forwards, and records every call as " ". +type scriptedSystemOneProvider struct { + *mockProvider + // statuses maps a forwarded model to the status it answers with; + // 200 by default. + statuses map[string]int + catalog map[string]core.Model + calls []string + // idempotencyKeys records the explicit Idempotency-Key of each call. + idempotencyKeys []string +} + +// newScriptedSystemOneProvider configures the given "/" +// selectors, keyed to their provider types. +func newScriptedSystemOneProvider(models map[string]string) *scriptedSystemOneProvider { + mock := &mockProvider{providerTypes: map[string]string{}, providerNames: map[string]string{}} + for qualified, providerType := range models { + providerName, model, _ := strings.Cut(qualified, "/") + mock.supportedModels = append(mock.supportedModels, model) + mock.providerTypes[qualified] = providerType + mock.providerNames[qualified] = providerName + } + return &scriptedSystemOneProvider{mockProvider: mock, statuses: map[string]int{}} +} + +func (p *scriptedSystemOneProvider) Passthrough(_ context.Context, _ string, req *core.PassthroughRequest) (*core.PassthroughResponse, error) { + raw, err := io.ReadAll(req.Body) + if err != nil { + return nil, err + } + var sent struct { + Model string `json:"model"` + } + if err := json.Unmarshal(raw, &sent); err != nil { + return nil, err + } + p.calls = append(p.calls, req.ProviderName+" "+req.Endpoint+" "+sent.Model) + p.idempotencyKeys = append(p.idempotencyKeys, req.Headers.Get(core.IdempotencyKeyHeader)) + + status := p.statuses[sent.Model] + body := `{"error":{"message":"overloaded"}}` + if status == 0 { + status = http.StatusOK + body = `{"model":"` + sent.Model + `-answered","answers":{},"usage":{"input_tokens":10,"output_tokens":1}}` + } + return &core.PassthroughResponse{ + StatusCode: status, + Headers: map[string][]string{"Content-Type": {"application/json"}}, + Body: io.NopCloser(strings.NewReader(body)), + }, nil +} + +func (p *scriptedSystemOneProvider) LookupModel(model string) (*core.Model, bool) { + found, ok := p.catalog[model] + return &found, ok +} + +func (p *scriptedSystemOneProvider) ProviderNamesForType(providerType string) []string { + var names []string + for qualified, candidate := range p.providerTypes { + if name := p.providerNames[qualified]; candidate == providerType && !slices.Contains(names, name) { + names = append(names, name) + } + } + slices.Sort(names) + return names +} + +func systemOneRequest(model string) string { + return `{"model":"` + model + `","state":"I was charged twice.","questions":{"refund":{"type":"noul","instructions":"Refund?"}}}` +} + +// An identical request is answered from the exact cache without reaching the +// provider. +func TestSystemOne_ServesRepeatsFromTheExactCache(t *testing.T) { + provider := newScriptedSystemOneProvider(map[string]string{"kev/kev-latest": "jev"}) + store := cache.NewMapStore() + defer store.Close() + mw := responsecache.NewResponseCacheMiddlewareWithStore(store, time.Hour) + handler := NewHandler(provider, nil, nil, nil) + handler.responseCache = mw + + c, first := echotest.Post(t, "/v1/systemone", systemOneRequest("kev-latest")) + require.NoError(t, handler.SystemOne(c)) + require.Equal(t, http.StatusOK, first.Code, first.Body.String()) + // The cache write is asynchronous; drain it before the repeat. + require.NoError(t, mw.Close()) + + c, second := echotest.Post(t, "/v1/systemone", systemOneRequest("kev-latest")) + require.NoError(t, handler.SystemOne(c)) + require.Equal(t, http.StatusOK, second.Code, second.Body.String()) + + assert.Equal(t, "HIT (exact)", second.Header().Get("X-Cache")) + assert.JSONEq(t, first.Body.String(), second.Body.String()) + assert.Len(t, provider.calls, 1, "the repeat must not reach the provider") +} + +// A different state is a different decision: it misses the cache. +func TestSystemOne_CacheKeyCoversTheState(t *testing.T) { + provider := newScriptedSystemOneProvider(map[string]string{"kev/kev-latest": "jev"}) + store := cache.NewMapStore() + defer store.Close() + mw := responsecache.NewResponseCacheMiddlewareWithStore(store, time.Hour) + handler := NewHandler(provider, nil, nil, nil) + handler.responseCache = mw + + c, rec := echotest.Post(t, "/v1/systemone", systemOneRequest("kev-latest")) + require.NoError(t, handler.SystemOne(c)) + require.Equal(t, http.StatusOK, rec.Code, rec.Body.String()) + require.NoError(t, mw.Close()) + + c, rec = echotest.Post(t, "/v1/systemone", strings.Replace(systemOneRequest("kev-latest"), "twice", "once", 1)) + require.NoError(t, handler.SystemOne(c)) + require.Equal(t, http.StatusOK, rec.Code, rec.Body.String()) + + assert.Empty(t, rec.Header().Get("X-Cache")) + assert.Len(t, provider.calls, 2) +} + +// When the primary fails with an availability error, the request moves to +// the virtual model's next target in its own dialect. A target without the +// System One API is skipped rather than called, and usage and audit carry the +// model that answered. +func TestSystemOne_FailsOverToTheNextSystemOneTarget(t *testing.T) { + provider := newScriptedSystemOneProvider(map[string]string{ + "kev/kev-latest": "jev", + "openai/gpt-5-mini": "openai", + "openrouter/typesafe/jev-1.13": "openrouter", + }) + provider.statuses["kev-latest"] = http.StatusServiceUnavailable + usageLogger := &collectingUsageLogger{config: usage.Config{Enabled: true}} + handler := newHandler(provider, nil, usageLogger, nil, nil, nil, failoverResolverStub{selectors: []core.ModelSelector{ + {Provider: "openai", Model: "gpt-5-mini"}, + {Provider: "openrouter", Model: "typesafe/jev-1.13"}, + }}, nil) + + entry := &auditlog.LogEntry{Data: &auditlog.LogData{}} + c, rec := echotest.Post(t, "/v1/systemone", systemOneRequest("kev/kev-latest"), echotest.WithValue(string(auditlog.LogEntryKey), entry)) + require.NoError(t, handler.SystemOne(c)) + require.Equal(t, http.StatusOK, rec.Code, rec.Body.String()) + + assert.Equal(t, []string{"kev systemone kev-latest", "openrouter systemone typesafe/jev-1.13"}, provider.calls, + "the chat model must be skipped, not called") + require.NotNil(t, entry.Data.Failover) + assert.Equal(t, "openrouter/typesafe/jev-1.13", entry.Data.Failover.TargetModel) + assert.Equal(t, "openrouter/typesafe/jev-1.13", entry.ResolvedModel) + assert.Equal(t, "openrouter", entry.Provider) + require.Len(t, entry.Data.Attempts, 2, "a skipped target is not an attempt") + assert.Equal(t, http.StatusServiceUnavailable, entry.Data.Attempts[0].StatusCode) + assert.True(t, entry.Data.Attempts[1].Success) + + require.Len(t, usageLogger.entries, 1) + assert.Equal(t, "openrouter", usageLogger.entries[0].Provider) + assert.Equal(t, "typesafe/jev-1.13-answered", usageLogger.entries[0].Model) +} + +// A skipped target does not count against max_attempts, so a chat model ahead +// of a valid target in the chain cannot use up the only failover attempt. The +// client's Idempotency-Key is not forwarded as a header: it reaches the +// primary through the request context, and a failover target's different +// body must not carry it. +func TestSystemOne_FailoverSkipsTargetsWithoutUsingAttempts(t *testing.T) { + provider := newScriptedSystemOneProvider(map[string]string{ + "kev/kev-latest": "jev", + "openai/gpt-5-mini": "openai", + "openrouter/typesafe/jev-1.13": "openrouter", + }) + provider.statuses["kev-latest"] = http.StatusServiceUnavailable + handler := newHandler(provider, nil, nil, nil, nil, nil, failoverResolverStub{selectors: []core.ModelSelector{ + {Provider: "openai", Model: "gpt-5-mini"}, + {Provider: "openrouter", Model: "typesafe/jev-1.13"}, + }}, nil) + handler.failoverPolicy = &gateway.FailoverPolicy{MaxAttempts: 1} + + c, rec := echotest.Post(t, "/v1/systemone", systemOneRequest("kev/kev-latest"), echotest.WithHeader(core.IdempotencyKeyHeader, "client-key-1")) + require.NoError(t, handler.SystemOne(c)) + require.Equal(t, http.StatusOK, rec.Code, rec.Body.String()) + + assert.Equal(t, []string{"kev systemone kev-latest", "openrouter systemone typesafe/jev-1.13"}, provider.calls) + assert.Equal(t, []string{"", ""}, provider.idempotencyKeys, "the key must not travel as an explicit header") +} + +// A client error such as a malformed question is not an availability +// problem: it is returned as the upstream reported it, without failover. +func TestSystemOne_DoesNotFailOverOnClientErrors(t *testing.T) { + provider := newScriptedSystemOneProvider(map[string]string{ + "kev/kev-latest": "jev", + "openrouter/typesafe/jev-1.13": "openrouter", + }) + provider.statuses["kev-latest"] = http.StatusUnprocessableEntity + handler := newHandler(provider, nil, nil, nil, nil, nil, failoverResolverStub{selectors: []core.ModelSelector{ + {Provider: "openrouter", Model: "typesafe/jev-1.13"}, + }}, nil) + + c, rec := echotest.Post(t, "/v1/systemone", systemOneRequest("kev/kev-latest")) + require.NoError(t, handler.SystemOne(c)) + + assert.Equal(t, http.StatusUnprocessableEntity, rec.Code, rec.Body.String()) + assert.Equal(t, []string{"kev systemone kev-latest"}, provider.calls) +} + +// Kev's diagnostic routes are served natively on jev providers and refused +// on OpenRouter, which serves only the evaluation route. +func TestSystemOne_KevDiagnosticRoutes(t *testing.T) { + provider := newScriptedSystemOneProvider(map[string]string{ + "kev/kev-latest": "jev", + "openrouter/typesafe/jev-1.13": "openrouter", + }) + handler := NewHandler(provider, nil, nil, nil) + + for path, serve := range map[string]func(*echo.Context) error{ + "/v1/systemone/permute": handler.SystemOnePermute, + "/v1/systemone/separate": handler.SystemOneSeparate, + } { + t.Run(path, func(t *testing.T) { + provider.calls = nil + c, rec := echotest.Post(t, path, systemOneRequest("kev/kev-latest")) + require.NoError(t, serve(c)) + require.Equal(t, http.StatusOK, rec.Code, rec.Body.String()) + assert.Equal(t, []string{"kev " + strings.TrimPrefix(path, "/v1/") + " kev-latest"}, provider.calls) + + c, rec = echotest.Post(t, path, systemOneRequest("openrouter/typesafe/jev-1.13")) + require.NoError(t, serve(c)) + assert.Equal(t, http.StatusBadRequest, rec.Code) + assert.Contains(t, rec.Body.String(), "this route is served by Kev servers") + }) + } +} + +// TypeSafe lists only its aliases but accepts any versioned ID, so a pinned +// version routes to a jev provider without being declared: named with its +// provider, bare when one jev provider is configured, or through a virtual +// model. A bare name stays unrouted when several jev providers could own it. +func TestSystemOne_RoutesUnlistedPinnedVersions(t *testing.T) { + aliases := systemOneAliasResolver{"pinned": {Provider: "jev", Model: "jev-1.13.0"}} + tests := []struct { + name string + models map[string]string + model string + wantCall string + wantCode int + }{ + {name: "provider-qualified", models: map[string]string{"jev/jev-latest": "jev"}, model: "jev/jev-1.13.0", wantCall: "jev systemone jev-1.13.0"}, + {name: "bare with one jev provider", models: map[string]string{"jev/jev-latest": "jev", "openrouter/typesafe/jev-1.13": "openrouter"}, model: "jev-1.13.0", wantCall: "jev systemone jev-1.13.0"}, + {name: "virtual model", models: map[string]string{"jev/jev-latest": "jev"}, model: "pinned", wantCall: "jev systemone jev-1.13.0"}, + {name: "self-hosted provider name", models: map[string]string{"kev/kev-latest": "jev"}, model: "kev/kev-4b-2026-09", wantCall: "kev systemone kev-4b-2026-09"}, + {name: "bare with two jev providers", models: map[string]string{"jev/jev-latest": "jev", "kev/kev-latest": "jev"}, model: "jev-1.13.0", wantCode: http.StatusNotFound}, + {name: "unknown on OpenRouter", models: map[string]string{"openrouter/typesafe/jev-1.13": "openrouter"}, model: "openrouter/typesafe/jev-9", wantCode: http.StatusNotFound}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + provider := newScriptedSystemOneProvider(tt.models) + handler := newHandler(provider, nil, nil, nil, aliases, nil, nil, nil) + + c, rec := echotest.Post(t, "/v1/systemone", systemOneRequest(tt.model)) + require.NoError(t, handler.SystemOne(c)) + + if tt.wantCode != 0 { + assert.Equal(t, tt.wantCode, rec.Code, rec.Body.String()) + assert.Empty(t, provider.calls) + return + } + require.Equal(t, http.StatusOK, rec.Code, rec.Body.String()) + assert.Equal(t, []string{tt.wantCall}, provider.calls) + }) + } +} + +// A System One model sent to an OpenAI-compatible route is refused by the +// gateway with a pointer at /v1/systemone, from the catalog alone, so the +// caller never sees the upstream's advice to use the upstream's own endpoint. +// Chat models on the same provider are unaffected. +func TestSystemOneModels_OnOpenAIRoutesPointAtSystemOne(t *testing.T) { + provider := newScriptedSystemOneProvider(map[string]string{ + "openrouter/~typesafe/jev-latest": "openrouter", + "openrouter/openai/gpt-4o-mini": "openrouter", + }) + provider.catalog = map[string]core.Model{ + "openrouter/~typesafe/jev-latest": {ID: "~typesafe/jev-latest", Metadata: &core.ModelMetadata{ + Categories: []core.ModelCategory{core.CategoryUtility}, + }}, + "openrouter/openai/gpt-4o-mini": {ID: "openai/gpt-4o-mini", Metadata: &core.ModelMetadata{ + Modes: []string{"chat"}, Categories: []core.ModelCategory{core.CategoryTextGeneration}, + }}, + } + provider.response = &core.ChatResponse{ID: "c1", Object: "chat.completion", Model: "openai/gpt-4o-mini", + Choices: []core.Choice{{Message: core.ResponseMessage{Role: "assistant", Content: "ok"}, FinishReason: "stop"}}} + srv := New(provider, &Config{ModelResolver: systemOneAliasResolver{"jev-latest": {Provider: "openrouter", Model: "~typesafe/jev-latest"}}}) + + requests := map[string]string{ + "/v1/chat/completions": `{"model":"jev-latest","messages":[{"role":"user","content":"hi"}]}`, + "/v1/responses": `{"model":"openrouter/~typesafe/jev-latest","input":"hi"}`, + "/v1/embeddings": `{"model":"openrouter/~typesafe/jev-latest","input":"hi"}`, + "/v1/messages": `{"model":"jev-latest","max_tokens":8,"messages":[{"role":"user","content":"hi"}]}`, + } + for path, body := range requests { + t.Run(path, func(t *testing.T) { + rec := postJSON(t, srv, path, body) + assert.Equal(t, http.StatusBadRequest, rec.Code, rec.Body.String()) + assert.Contains(t, rec.Body.String(), "System One decision model") + assert.Contains(t, rec.Body.String(), "POST /v1/systemone") + }) + } + assert.Zero(t, provider.chatCompletionCalls) + + rec := postJSON(t, srv, "/v1/chat/completions", `{"model":"openrouter/openai/gpt-4o-mini","messages":[{"role":"user","content":"hi"}]}`) + assert.Equal(t, http.StatusOK, rec.Code, rec.Body.String()) +} diff --git a/internal/server/systemone_handler.go b/internal/server/systemone_handler.go index b297e7fcf..6c43c262e 100644 --- a/internal/server/systemone_handler.go +++ b/internal/server/systemone_handler.go @@ -2,8 +2,8 @@ package server import ( "bytes" + "context" "fmt" - "io" "log/slog" "net/http" "slices" @@ -12,22 +12,37 @@ import ( "github.com/goccy/go-json" "github.com/labstack/echo/v5" - "github.com/enterpilot/gomodel/internal/auditlog" "github.com/enterpilot/gomodel/internal/core" "github.com/enterpilot/gomodel/internal/gateway" "github.com/enterpilot/gomodel/internal/plugins" ) -const ( - systemOnePath = "/v1/systemone" - systemOneEndpoint = "systemone" +// systemOneRoute is one System One route the gateway serves natively. +type systemOneRoute struct { + // path is the gateway route; endpoint is the same route as a provider's + // passthrough spells it, without the /v1 prefix. + path string + endpoint string + // operation names the call in provider metrics and logs. + operation string + // kevOnly marks the diagnostic routes only Kev servers implement; the + // hosted API and OpenRouter serve the evaluation route alone. + kevOnly bool +} + +var ( + systemOneEvaluate = systemOneRoute{path: "/v1/systemone", endpoint: "systemone", operation: "systemone"} + systemOnePermute = systemOneRoute{path: "/v1/systemone/permute", endpoint: "systemone/permute", operation: "systemone_permute", kevOnly: true} + systemOneSeparate = systemOneRoute{path: "/v1/systemone/separate", endpoint: "systemone/separate", operation: "systemone_separate", kevOnly: true} ) +const jevProviderType = "jev" + // systemOneProviderTypes are the provider types that serve the System One API // natively: jev (TypeSafe's hosted Jev and self-hosted Kev servers) and // OpenRouter, which serves Jev and Kev at the same path with the same request // and answer shapes. Configuring either makes /v1/systemone available. -var systemOneProviderTypes = []string{"jev", "openrouter"} +var systemOneProviderTypes = []string{jevProviderType, "openrouter"} // SystemOne handles POST /v1/systemone. // @@ -52,13 +67,53 @@ var systemOneProviderTypes = []string{"jev", "openrouter"} // @Failure 502 {object} core.OpenAIErrorEnvelope // @Router /v1/systemone [post] func (h *Handler) SystemOne(c *echo.Context) error { - return h.translatedInference().SystemOne(c) + return h.translatedInference().serveSystemOne(c, systemOneEvaluate) } -// SystemOne resolves, guards, and forwards one System One request. -func (s *translatedInferenceService) SystemOne(c *echo.Context) error { +// SystemOnePermute handles POST /v1/systemone/permute. +// +// @Summary Run one Choice question with several option orders (Kev) +// @Description A Kev server diagnostic: the request is a System One request, and n_perm (1 to 64, default 6) sets how many option orders run. Only jev providers pointing at a Kev server serve it. +// @Tags systemone +// @Accept json +// @Produce json +// @Security BearerAuth +// @Param request body object true "System One request with one Choice question" +// @Success 200 {object} object "Kev's answer, in the provider's shape" +// @Failure 400 {object} core.OpenAIErrorEnvelope +// @Failure 401 {object} core.OpenAIErrorEnvelope +// @Failure 404 {object} core.OpenAIErrorEnvelope +// @Failure 429 {object} core.OpenAIErrorEnvelope +// @Failure 502 {object} core.OpenAIErrorEnvelope +// @Router /v1/systemone/permute [post] +func (h *Handler) SystemOnePermute(c *echo.Context) error { + return h.translatedInference().serveSystemOne(c, systemOnePermute) +} + +// SystemOneSeparate handles POST /v1/systemone/separate. +// +// @Summary Run each System One question in its own forward pass (Kev) +// @Description A Kev server diagnostic that answers each question separately. Only jev providers pointing at a Kev server serve it. +// @Tags systemone +// @Accept json +// @Produce json +// @Security BearerAuth +// @Param request body object true "System One request: model, state, and questions" +// @Success 200 {object} object "Kev's answer, in the provider's shape" +// @Failure 400 {object} core.OpenAIErrorEnvelope +// @Failure 401 {object} core.OpenAIErrorEnvelope +// @Failure 404 {object} core.OpenAIErrorEnvelope +// @Failure 429 {object} core.OpenAIErrorEnvelope +// @Failure 502 {object} core.OpenAIErrorEnvelope +// @Router /v1/systemone/separate [post] +func (h *Handler) SystemOneSeparate(c *echo.Context) error { + return h.translatedInference().serveSystemOne(c, systemOneSeparate) +} + +// serveSystemOne resolves, guards, and forwards one System One request. +func (s *translatedInferenceService) serveSystemOne(c *echo.Context, route systemOneRoute) error { if !systemOneAvailable(s.provider) { - return handleError(c, core.NewNotFoundError("POST "+systemOnePath+" is available only when a jev or openrouter provider is configured")) + return handleError(c, core.NewNotFoundError("POST "+route.path+" is available only when a jev or openrouter provider is configured")) } body, err := requestBodyBytes(c) if err != nil { @@ -76,26 +131,21 @@ func (s *translatedInferenceService) SystemOne(c *echo.Context) error { if err != nil { return handleError(c, err) } - ctx := c.Request().Context() resolution := workflow.Resolution if s.modelAuthorizer != nil { - if err := s.modelAuthorizer.ValidateModelAccess(ctx, resolution.ResolvedSelector); err != nil { + if err := s.modelAuthorizer.ValidateModelAccess(c.Request().Context(), resolution.ResolvedSelector); err != nil { return handleError(c, err) } } - if reason := s.systemOneUnsupportedReason(resolution); reason != "" { - return handleError(c, systemOneUnsupportedModelError(c, resolution, reason)) + if reason := s.systemOneUnsupportedReason(route, resolution.ResolvedSelector, resolution.ProviderType); reason != "" { + return handleError(c, systemOneUnsupportedModelError(c, route, resolution, reason)) } body, err = s.guardSystemOneState(c, workflow, &req, body) if err != nil { return handleError(c, err) } - model := resolution.ResolvedSelector.Model - if body, err = rewriteMessagesModel(body, model); err != nil { - return handleError(c, core.NewInvalidRequestError("invalid request body: "+err.Error(), err)) - } - return s.dispatchSystemOne(c, workflow, model, body) + return s.dispatchSystemOneWithCache(c, route, workflow, body) } // systemOneAvailable reports whether a provider that serves System One is @@ -114,17 +164,23 @@ func systemOneAvailable(provider core.RoutableProvider) bool { return false } -// systemOneWorkflow returns the request's workflow with its model resolved. -// The workflow middleware resolves it from the body; a request that reached -// the handler without one (the body was not parsed there) is resolved here, -// so virtual models and workflow policy apply either way. +// systemOneWorkflow resolves the request's model (virtual models first, then +// the catalog) and builds its workflow. A model the catalog does not list may +// still be a pinned version a jev provider accepts; see unlistedJevResolution. func (s *translatedInferenceService) systemOneWorkflow(c *echo.Context, model string) (*core.Workflow, error) { if workflow := core.GetWorkflow(c.Request().Context()); workflow != nil && workflow.Resolution != nil { return workflow, nil } - resolution, err := resolveAndStoreRequestModelResolution(c, s.provider, s.modelResolver, nil, model, "") - if err != nil { - return nil, err + requested := core.NewRequestedModelSelector(model, "") + resolution, ok := s.unlistedJevResolution(c.Request().Context(), requested) + if ok { + enrichAuditEntryWithRequestedModel(c, requested) + } else { + var err error + resolution, err = resolveAndStoreRequestModelResolution(c, s.provider, s.modelResolver, nil, model, "") + if err != nil { + return nil, err + } } workflow, err := translatedWorkflowForRequest(c, resolution, s.workflowPolicyResolver) if err != nil { @@ -134,41 +190,125 @@ func (s *translatedInferenceService) systemOneWorkflow(c *echo.Context, model st return workflow, nil } +// providerNamesByType lists configured provider instances of one type; the +// provider router implements it. +type providerNamesByType interface { + ProviderNamesForType(providerType string) []string +} + +// unlistedJevResolution routes a model the catalog does not list to a jev +// provider. TypeSafe lists only its aliases, yet accepts every versioned ID +// (jev-1.13.0) in the model field, so a pinned version must work without +// being declared first. The provider is the one the model names +// (jev/jev-1.13.0, kev/...), or for a bare name the only jev provider +// configured; virtual models apply first, so one can pin a version too. It +// is checked before the regular resolution, which would refresh the +// provider's model list on every such request; the upstream reports a name +// it rejects. +func (s *translatedInferenceService) unlistedJevResolution(ctx context.Context, requested core.RequestedModelSelector) (*core.RequestModelResolution, bool) { + selector, aliasApplied, err := gateway.ResolveExecutionSelector(ctx, s.provider, s.modelResolver, requested) + if err != nil || selector.Model == "" || s.provider.Supports(selector.QualifiedModel()) { + return nil, false + } + providerName := "" + switch named, _ := s.provider.(core.ProviderNameTypeResolver); { + case selector.Provider == "": + if lister, ok := s.provider.(providerNamesByType); ok { + if names := lister.ProviderNamesForType(jevProviderType); len(names) == 1 { + providerName = names[0] + } + } + case named != nil && named.GetProviderTypeForName(selector.Provider) == jevProviderType: + providerName = selector.Provider + case selector.Provider == jevProviderType: + if byType, ok := s.provider.(core.ProviderTypeNameResolver); ok { + providerName = byType.GetProviderNameForType(jevProviderType) + } + } + if strings.TrimSpace(providerName) == "" { + return nil, false + } + return &core.RequestModelResolution{ + Requested: requested, + ResolvedSelector: core.ModelSelector{Provider: providerName, Model: selector.Model}, + ProviderType: jevProviderType, + ProviderName: providerName, + AliasApplied: aliasApplied, + }, true +} + // modelCatalog describes single catalog models; the provider router // implements it. type modelCatalog interface { LookupModel(model string) (*core.Model, bool) } -// systemOneUnsupportedReason explains why the resolved model cannot answer a -// System One request, or returns "" when it can. The provider must serve the -// API, and since OpenRouter also serves chat models, the model must not be -// catalogued with a generation mode. A model the catalog does not describe is -// given the benefit of the doubt: the upstream reports it if it is wrong. -func (s *translatedInferenceService) systemOneUnsupportedReason(resolution *core.RequestModelResolution) string { - providerType := strings.TrimSpace(resolution.ProviderType) +// systemOneUnsupportedReason explains why a model cannot answer a request on +// route, or returns "" when it can. The provider must serve the API, a Kev +// diagnostic route needs a jev provider, and since OpenRouter also serves +// chat models, the model must not be catalogued with a generation mode. A +// model the catalog does not describe is given the benefit of the doubt: the +// upstream reports it if it is wrong. +func (s *translatedInferenceService) systemOneUnsupportedReason(route systemOneRoute, selector core.ModelSelector, providerType string) string { + providerType = strings.TrimSpace(providerType) if !slices.Contains(systemOneProviderTypes, providerType) { - return fmt.Sprintf("is served by a %s provider, which has no System One API", providerType) + return fmt.Sprintf("is served by provider type %s, which has no System One API", providerType) + } + if route.kevOnly && providerType != jevProviderType { + return fmt.Sprintf("is served by provider type %s, which answers only %s; this route is served by Kev servers", providerType, systemOneEvaluate.path) } catalog, ok := s.provider.(modelCatalog) if !ok { return "" } - model, ok := catalog.LookupModel(resolution.ResolvedQualifiedModel()) + model, ok := catalog.LookupModel(selector.QualifiedModel()) if !ok || model == nil || model.Metadata == nil || len(model.Metadata.Modes) == 0 { return "" } return fmt.Sprintf("is a %s model, not a System One model", strings.Join(model.Metadata.Modes, "/")) } +// systemOneOnlyModel reports whether the resolved model answers only System +// One requests: served by a System One provider and catalogued as a utility +// model with no generation mode, as jev models and OpenRouter's decision +// models are. The catalog is consulted rather than the provider, so the check +// holds from startup, before any provider has listed its models again. +func systemOneOnlyModel(provider core.RoutableProvider, resolution *core.RequestModelResolution) bool { + if resolution == nil || !slices.Contains(systemOneProviderTypes, strings.TrimSpace(resolution.ProviderType)) { + return false + } + catalog, ok := provider.(modelCatalog) + if !ok { + return false + } + model, ok := catalog.LookupModel(resolution.ResolvedQualifiedModel()) + if !ok || model == nil || model.Metadata == nil || len(model.Metadata.Modes) > 0 { + return false + } + return slices.Contains(model.Metadata.Categories, core.CategoryUtility) +} + +// systemOneOnlyModelError points a chat, Responses, or embeddings request for +// a System One model at /v1/systemone. Without it the caller would see the +// upstream's own advice, which names the upstream's endpoint, not the +// gateway's. +func systemOneOnlyModelError(operation core.Operation, resolution *core.RequestModelResolution) error { + surface := strings.ReplaceAll(string(operation), "_", " ") + return core.NewInvalidRequestError(fmt.Sprintf( + "model %q is a System One decision model and does not support %s; it answers decision requests, which GoModel does not translate: send them to POST %s", + resolution.RequestedQualifiedModel(), surface, systemOneEvaluate.path, + ), nil).WithParam("model") +} + // systemOneUnsupportedModelError explains a request whose model cannot answer // System One. It is also logged: a virtual model that sends System One // traffic to a chat model is an operator mistake the caller cannot fix. -func systemOneUnsupportedModelError(c *echo.Context, resolution *core.RequestModelResolution, reason string) error { +func systemOneUnsupportedModelError(c *echo.Context, route systemOneRoute, resolution *core.RequestModelResolution, reason string) error { requested := resolution.RequestedQualifiedModel() resolved := resolution.ResolvedQualifiedModel() slog.Warn("System One request routed to a model without the System One API", "request_id", requestIDFromContextOrHeader(c.Request()), + "path", route.path, "requested_model", requested, "resolved_model", resolved, "provider_type", resolution.ProviderType, @@ -180,7 +320,7 @@ func systemOneUnsupportedModelError(c *echo.Context, resolution *core.RequestMod } return core.NewInvalidRequestError(fmt.Sprintf( "model %s %s; %s forwards requests natively and does not translate them to other APIs, so use a System One model such as a jev model or OpenRouter's typesafe/jev-1.13", - target, reason, systemOnePath, + target, reason, route.path, ), nil).WithParam("model") } @@ -213,48 +353,3 @@ func (s *translatedInferenceService) guardSystemOneState(c *echo.Context, workfl } return rewritten, nil } - -// dispatchSystemOne forwards the body to the resolved provider and relays its -// answer unchanged, with admission, audit, and usage accounting. -func (s *translatedInferenceService) dispatchSystemOne(c *echo.Context, workflow *core.Workflow, model string, body []byte) error { - passthroughProvider, ok := s.provider.(core.RoutablePassthrough) - if !ok { - return handleError(c, core.NewInvalidRequestError("provider passthrough is not supported by the current provider router", nil)) - } - s.observeLiveProviderAttempts(c, workflow) - - adm, err := enforceAdmission(c, s.rateLimiter, s.budgetChecker, rateLimitRouteFromWorkflow(workflow)) - if err != nil { - return handleError(c, err) - } - defer adm.release() - ctx := adm.dispatchContext(c.Request().Context()) - - resolution := workflow.Resolution - providerType := strings.TrimSpace(resolution.ProviderType) - providerName := strings.TrimSpace(resolution.ProviderName) - resp, err := passthroughProvider.Passthrough(ctx, providerType, &core.PassthroughRequest{ - Method: http.MethodPost, - Endpoint: systemOneEndpoint, - Operation: "systemone", - Model: model, - Body: io.NopCloser(bytes.NewReader(body)), - Headers: buildPassthroughHeaders(ctx, c.Request().Header), - ProviderName: providerName, - }) - if err != nil { - return handleError(c, err) - } - - auditlog.EnrichEntryWithWorkflow(c, workflow) - auditlog.EnrichEntryWithResolvedRoute(c, resolution.ResolvedQualifiedModel(), providerType, providerName) - info := &core.PassthroughRouteInfo{ - Provider: providerType, - ProviderName: providerName, - NormalizedEndpoint: systemOneEndpoint, - SemanticOperation: "systemone", - AuditPath: systemOnePath, - Model: model, - } - return proxyPassthroughResponse(c, s.logger, s.usageLogger, s.pricingResolver, providerType, providerName, systemOneEndpoint, info, resp) -} diff --git a/internal/usage/extractor.go b/internal/usage/extractor.go index 46ed07b3d..250b74461 100644 --- a/internal/usage/extractor.go +++ b/internal/usage/extractor.go @@ -281,6 +281,9 @@ func ExtractFromCachedResponseBody( if entry == nil { entry = extractFromCachedSSEBody(body, requestID, model, provider, endpoint, pricing...) } + if entry == nil { + entry = extractFromCachedJSONBody(body, requestID, model, provider, endpoint, pricing...) + } if entry == nil { entry = &UsageEntry{ @@ -346,6 +349,31 @@ func extractFromCachedSSEBody( return observer.cachedEntry } +// extractFromCachedJSONBody reads usage from a cached JSON body of an +// endpoint without a typed response, such as a native /v1/systemone answer, +// the same way a live passthrough response is read. +func extractFromCachedJSONBody( + body []byte, + requestID, model, provider, endpoint string, + pricing ...*core.ModelPricing, +) *UsageEntry { + var payload map[string]any + if err := json.Unmarshal(body, &payload); err != nil || payload == nil { + return nil + } + observer := &StreamUsageObserver{ + model: strings.TrimSpace(model), + provider: strings.TrimSpace(provider), + requestID: strings.TrimSpace(requestID), + endpoint: endpoint, + } + if len(pricing) > 0 && pricing[0] != nil { + observer.pricingResolver = staticPricingResolver{pricing: pricing[0]} + } + observer.OnJSONEvent(payload) + return observer.cachedEntry +} + func normalizeCachedResponseEndpoint(endpoint string) string { normalized := strings.TrimSpace(endpoint) if normalized == "" { diff --git a/internal/usage/extractor_test.go b/internal/usage/extractor_test.go index 86bfed117..01cc313a8 100644 --- a/internal/usage/extractor_test.go +++ b/internal/usage/extractor_test.go @@ -493,6 +493,20 @@ func TestExtractFromCachedResponseBody(t *testing.T) { require.Equal(t, 10, entry.TotalTokens) }) + // A native endpoint without a typed response (System One) is read like a + // live passthrough answer, so a cache hit keeps its token counts. + t.Run("reads usage from an untyped JSON body", func(t *testing.T) { + body := []byte(`{"model":"jev-1.13.0","answers":{"refund":{"type":"noul","noul":0.98}},"usage":{"input_tokens":275,"output_tokens":20}}`) + + entry := ExtractFromCachedResponseBody(body, "req-systemone", "jev-latest", "jev", "/v1/systemone", CacheTypeExact) + require.NotNil(t, entry) + require.Equal(t, CacheTypeExact, entry.CacheType) + require.Equal(t, "/v1/systemone", entry.Endpoint) + require.Equal(t, "jev", entry.Provider) + require.Equal(t, 275, entry.InputTokens) + require.Equal(t, 20, entry.OutputTokens) + }) + t.Run("falls back to synthetic entry when body cannot be parsed", func(t *testing.T) { entry := ExtractFromCachedResponseBody([]byte("{"), "req-cache-fallback", "gpt-4o", "openai", "/v1/chat/completions", CacheTypeExact) require.NotNil(t, entry) diff --git a/web/dashboard/src/pages/audit-logs/audit-operations.js b/web/dashboard/src/pages/audit-logs/audit-operations.js index 7bbd43321..67da43c21 100644 --- a/web/dashboard/src/pages/audit-logs/audit-operations.js +++ b/web/dashboard/src/pages/audit-logs/audit-operations.js @@ -37,6 +37,8 @@ const EXACT_PATHS = { "/v1/images/generations": "images", "/v1/images/edits": "images", "/v1/systemone": "systemone", + "/v1/systemone/permute": "systemone", + "/v1/systemone/separate": "systemone", "/v1/realtime": "realtime", "/v1/realtime/calls": "realtime", "/v1/realtime/client_secrets": "realtime", diff --git a/web/dashboard/tests/audit-operations.test.js b/web/dashboard/tests/audit-operations.test.js index e94d99faf..2c2b7d885 100644 --- a/web/dashboard/tests/audit-operations.test.js +++ b/web/dashboard/tests/audit-operations.test.js @@ -21,7 +21,9 @@ test("auditTypeForPath mirrors the gateway endpoint classification", () => { ["/v1/audio/transcriptions?x=1", "audio"], ["/v1/images/edits/", "images"], ["/v1/systemone", "systemone"], - ["/v1/systemone/permute", ""], + ["/v1/systemone/permute", "systemone"], + ["/v1/systemone/separate/", "systemone"], + ["/v1/systemone/other", ""], ["/v1/realtime/translations/calls", "realtime"], ["/mcp", "mcp"], ["/mcp/github", "mcp"],