diff --git a/scripts/ops/hermes_lan_generate.sh b/scripts/ops/hermes_lan_generate.sh index 409d1310..f4e6f742 100755 --- a/scripts/ops/hermes_lan_generate.sh +++ b/scripts/ops/hermes_lan_generate.sh @@ -31,7 +31,7 @@ import json, os, sys json.dump({"model": "qwen2.5:14b-instruct-q4_0", "prompt": sys.stdin.read(), "stream": False, "options": {"num_predict": int(os.getenv("HERMES_LAN_MAX_TOKENS", "256"))}}, sys.stdout) ' > "$scratch/request.json" -curl --fail-with-body --silent --show-error --connect-timeout 10 --max-time 310 \ +curl --fail-with-body --silent --show-error --connect-timeout 10 --max-time 1220 \ --noproxy worker.bstein.dev \ --resolve "worker.bstein.dev:443:${HERMES_LAN_ADDRESS:-192.168.22.50}" \ --config "$scratch/curl.conf" \ diff --git a/services/ai-llm/batch-deployment.yaml b/services/ai-llm/batch-deployment.yaml index 2adc648f..20587281 100644 --- a/services/ai-llm/batch-deployment.yaml +++ b/services/ai-llm/batch-deployment.yaml @@ -5,7 +5,7 @@ metadata: name: ollama-batch namespace: ai spec: - replicas: 1 + replicas: 0 revisionHistoryLimit: 2 strategy: type: Recreate diff --git a/services/hermes/LAN_API.md b/services/hermes/LAN_API.md new file mode 100644 index 00000000..7a23424a --- /dev/null +++ b/services/hermes/LAN_API.md @@ -0,0 +1,221 @@ +# Private inference API + +This service accepts stateless inference requests from the Atlas LAN. It does +not run the importer or contact ClickUp. The application supplies its own prompt +and schema and controls selection, validation, caching, and export policy. + +## Connection + +| Setting | Value | +|---|---| +| LAN_HOST | `worker.bstein.dev` | +| LAN_IP | `192.168.22.50` | +| PORT | `443` | +| HEALTH_URL | `https://worker.bstein.dev/local-model/healthz` | +| INFERENCE_URL | `https://worker.bstein.dev/local-model/api/generate` | +| MODEL_OR_ROUTE | `qwen2.5:14b-instruct-q4_0` | +| AUTHENTICATION_HEADER | `Authorization: Bearer ` | +| SECURE_CREDENTIAL_RETRIEVAL | Existing authenticated Vault UI, `kv/atlas/hermes/model-gate-lan-api`, field `token` | +| TLS_OR_CA_REQUIREMENTS | Normal system CA trust; certificate hostname `worker.bstein.dev`; no insecure TLS option | +| INITIAL_CONCURRENCY | One LAN generation; additional concurrent generations return HTTP 429 | +| REQUEST_TIMEOUT | Gateway 1200 s; Traefik response-header timeout 1210 s; client 1220 s | +| INPUT_AND_OUTPUT_LIMITS | Body 131072 bytes; schema 32768 bytes; context 8192 tokens; output 1–2048 tokens | + +Retrieve the token through the existing trusted Vault login, for example the +[secret's Vault UI page](https://vault.bstein.dev/ui/vault/secrets/kv/show/atlas/hermes/model-gate-lan-api). +Copy only this scoped token to the WSL client. No Vault token, Kubernetes access, +or admin credential belongs in the client. The token authorizes this API's +health and inference operations; it grants no dashboard, terminal, or model +management access. It is a static credential: rotate it deliberately in Vault +and change the deployment revision through Git/Flux to reload it. + +Public DNS continues to serve the existing worker site. The commands below +connect explicitly to the LAN IP while retaining the hostname for TLS/SNI. +The permitted source network is `192.168.22.0/24`. VPN routes or WSL host routing +still need verification on the laptop; a private address alone does not prove +that every hop remains on the LAN. + +## WSL commands + +Read the privately retrieved credential without storing it in shell history: + +```bash +read -rsp 'LAN API token: ' LOCAL_INFERENCE_TOKEN; printf '\n' +``` + +Health/connectivity, including runtime and model-digest readiness: + +```bash +curl --noproxy '*' --resolve worker.bstein.dev:443:192.168.22.50 \ + --connect-timeout 10 --max-time 30 --fail-with-body --silent --show-error \ + --config <(printf 'header = "Authorization: Bearer %s"\n' "$LOCAL_INFERENCE_TOKEN") \ + https://worker.bstein.dev/local-model/healthz +``` + +Plain generation: + +```bash +curl --noproxy '*' --resolve worker.bstein.dev:443:192.168.22.50 \ + --connect-timeout 10 --max-time 1220 --fail-with-body --silent --show-error \ + --config <(printf 'header = "Authorization: Bearer %s"\n' "$LOCAL_INFERENCE_TOKEN") \ + --header 'Content-Type: application/json' --data-binary @- \ + https://worker.bstein.dev/local-model/api/generate <<'JSON' +{"model":"qwen2.5:14b-instruct-q4_0","stream":false,"prompt":"Reply with exactly LAN_POC_OK and nothing else.","options":{"num_predict":32,"temperature":0,"seed":0}} +JSON +``` + +Synthetic structured profile: + +```bash +curl --noproxy '*' --resolve worker.bstein.dev:443:192.168.22.50 \ + --connect-timeout 10 --max-time 1220 --fail-with-body --silent --show-error \ + --config <(printf 'header = "Authorization: Bearer %s"\n' "$LOCAL_INFERENCE_TOKEN") \ + --header 'Content-Type: application/json' --data-binary @- \ + https://worker.bstein.dev/local-model/api/generate <<'JSON' +{ + "model": "qwen2.5:14b-instruct-q4_0", + "stream": false, + "prompt": "SYNTHETIC ONLY. Case SYN-LAN-001. Target: calculator HTTP API. Setup: test instance running. Action: POST /add with a=17,b=20. Verify JSON sum=37. Authentication details are unspecified. Return a compact implementation profile using the schema. Preserve the case ID and numeric facts. Do not invent missing details.", + "format": { + "type": "object", + "properties": { + "case_id": {"type": "string"}, + "target": {"type": "string"}, + "setup": {"type": "string"}, + "stimulus": {"type": "string"}, + "expected_sum": {"type": "integer"}, + "missing_information": {"type": "array", "items": {"type": "string"}} + }, + "required": ["case_id", "target", "setup", "stimulus", "expected_sum", "missing_information"], + "additionalProperties": false + }, + "options": {"num_ctx": 8192, "num_predict": 512, "temperature": 0, "top_p": 1, "top_k": 40, "seed": 0} +} +JSON +``` + +These commands bypass environment proxies, verify TLS, do not follow redirects, +and keep the token out of curl's argument list. No files or source records are +uploaded by the examples. Do not use verbose HTTP tracing with real credentials +or case material. `ip route get 192.168.22.50` inside WSL and the Windows host's +route table can help check the remaining laptop/VPN path. + +## Request and response contract + +Only `GET /healthz` and `POST /api/generate` exist under `/local-model`. +Other paths return 404. The POST body is Ollama-native JSON: + +- Required: exact `model`, nonempty string `prompt`, and `stream: false`. +- Optional `format`: `"json"` or an object JSON schema. Remote schema references + are rejected. No schema or referenced document is fetched over the network. +- Optional `options`: `num_ctx` (8192 only), `num_predict` (1–2048, default 256), + `temperature` (0–2, default 0), `top_p` (0–1, default 1), `top_k` (1–100, + default 40), and `seed` (0–2147483647, default 0). +- `messages`, `system`, `context`, custom templates, tools, images, streaming, + keep-alive changes, installation, and model deletion are unavailable. + +The gateway requires `UTF8_bytes(prompt) + num_predict + 1024 <= 8192`. +This deliberately conservative admission rule reserves output and template +space before inference. For example, a 512-token output budget permits a prompt +up to 6656 UTF-8 bytes. It never splits, trims, or silently substitutes input. +Large suites need application-controlled bounded requests and reconciliation. + +The response is an object containing **generated JSON as text in `response`**. +With urllib, parse the HTTP body as JSON, then call +`json.loads(envelope["response"])`. Validate that object against the application's +schema and fidelity criteria. JSON grammar is not a guarantee of semantic truth. +The gateway rejects invalid JSON or an exhausted output budget instead of +returning an incomplete success. It removes Ollama's context token array and +does not expose hidden reasoning or an upstream error body. + +Native metadata includes `model`, `done`, `done_reason`, `created_at`, token counts +`prompt_eval_count`/`eval_count`, and nanosecond durations `total_duration`, +`load_duration`, `prompt_eval_duration`, and `eval_duration` where supplied. +`inference_provenance` records the runtime, exact weight digest, placement, +effective options, protocol version, backend API, and gateway wall seconds. + +Errors are JSON with a static `error` string: 400 invalid request; 401 invalid +credential; 403 non-LAN source; 404 unsupported route; 408 body-read timeout; +413 body/schema/context limit; 422 incomplete or invalid structured output; +429 LAN generation already active; 502 local backend rejected a request; +503 pinned model/runtime or backend capacity unavailable; 504 inference timeout. +There is no automatic retry or alternate provider. Clients choose explicit, +bounded retries for 429/503 after waiting. A disconnected/timed-out caller +should not assume its computation stopped immediately. + +## Runtime and capacity + +The backend is the existing Ollama 0.13.5 on **titan-20, Jetson Xavier with 16 GB +unified memory**, using its CUDA JetPack 5 runtime. The initial model has 14.8B +parameters, Q4_0 weights, and this exact manifest digest: + +`5449194ff8035ccb13a6409a5814de6c8f9c39f555f429e383ae0fb7137001bd` + +The model advertises a 32768-token training context; this deployment and API +support **8192**, with the stricter admission rule above. Both model digest and +runtime version are checked before any prompt is forwarded. A mismatch fails +closed and requires a deliberate configuration update, never a substitution. + +The LAN gateway rejects a second in-flight LAN generation (429), with no LAN +request queue. The shared Ollama server serializes inference and model residency +and can queue behind existing internal consumers. Its existing queue is bounded +by Ollama's configured/default capacity; a queue-full response becomes 503. +Any upstream wait counts toward the 1200-second gateway budget. The backend has +not been moved, restarted, upgraded, or given another GPU for this endpoint. +Normal inference can switch residency between its existing 3B and 14B models. +Other consumers can therefore affect latency; this is not a dedicated capacity +reservation. The separate titan-23 CPU experiment is parked at zero replicas. + +Traefik's existing entrypoints have no shorter read/write/idle timeout. The LAN +Service selects a dedicated ServersTransport with a 1210-second response-header +timeout. Gateway request-body reads are limited to 10 seconds; backend responses +are non-streaming and bounded to 512 KiB. Client timeout is 1220 seconds. + +## Data path and retention + +The LAN handler calls only the fixed in-cluster Ollama `/api/version`, `/api/tags`, +and `/api/generate` operations. It ignores proxy environment variables and +rejects redirects. There is no Switchyard, hosted target, agent execution, +auxiliary model call, or fallback on this path. The gateway NetworkPolicy admits +port 8082 only from Traefik and permits only its required cluster destinations. +The existing shared Ollama pod is not under an egress-deny policy; this endpoint +enforces local execution by pinning its local GGUF manifest and native backend +operation. No claim is made that unrelated consumers of that pod are isolated. + +The gateway stores no prompts, responses, cache, agent memory, or conversation +history. It has no data PVC. Ollama uses a local-path volume for model weights; +native API calls do not create CLI history. Inference uses transient RAM/KV +buffers; this is not a forensic memory-erasure guarantee. No token-array context +is accepted from clients or returned to them. + +Traefik has access logs and tracing disabled. Gateway logs contain only a fixed +route category and status; errors never echo request fields or backend bodies. +Ollama's debug logging is disabled; observed logs contain HTTP/timing metadata. +Fluent Bit explicitly excludes these gateway and inference container logs. +This matters because central OpenSearch uses Longhorn storage whose configured +backup target is external B2 (`s3://atlas-soteria@us-west-004/`). The endpoint does +not send its logs or contents into that pipeline. Existing OpenTelemetry exports +go to in-cluster Data Prepper/OpenSearch; neither this gateway nor Ollama enables +request tracing. The model volume is local-path, outside Longhorn backups. + +## Deployment and rollback + +Changes are Git/Flux managed: the dedicated MetalLB LAN address, Traefik LAN +LoadBalancer, allowlist/prefix ingress, scoped Vault token, restricted gateway +listener/NetworkPolicy, native request validation, model pins, timeout transport, +and Fluent Bit exclusions. Existing model/GPU placement is unchanged. + +To disable the endpoint, remove `model-gate-lan-ingress.yaml` from +`services/hermes/kustomization.yaml`, commit and push the reviewed change, then +reconcile `hermes` with source. Flux pruning removes the route and its middleware +and transport. Revert that commit to restore it. This does not restart Ollama. +Keep log exclusions while any real-data inference remains enabled. To roll back +code, revert the relevant deployment commit(s) through Git, retain the log +exclusions, and reconcile `hermes`; do not use manual kubectl edits. + +The focused tests use synthetic inputs and mocked upstreams. They cover auth, +LAN-source validation, strict routes, schema forwarding, context limits, exact +model/runtime checks, concurrent rejection, redacted logs/errors, timeout, +redirect rejection, and failure without fallback. They make no inference calls. +Live synthetic evidence and the laptop-test boundary are recorded separately +after deployment. No roster, importer application, or FA01 pilot has been tested. diff --git a/services/hermes/NOTES.md b/services/hermes/NOTES.md index 8a277078..29e167d8 100644 --- a/services/hermes/NOTES.md +++ b/services/hermes/NOTES.md @@ -8,24 +8,21 @@ client IPs; the shared public LoadBalancer masks them. The service and middlewar allow only `192.168.22.0/24`, and the gateway also requires a bearer token from Vault at `kv/atlas/hermes/model-gate-lan-api`, field `token`. -Use `scripts/ops/hermes_lan_generate.sh < prompt.txt` from the LAN with a logged-in -Vault CLI, or set `HERMES_LAN_TOKEN_FILE` to a private file containing that token. -The helper connects directly to the LAN IP while verifying the existing -`worker.bstein.dev` TLS certificate. It needs no DNS override. Public DNS continues -to serve the worker dashboard; it does not select the LAN gateway automatically. +The endpoint contract and literal WSL curl commands are in +[LAN_API.md](LAN_API.md). The laptop needs only curl or Python urllib and the +scoped bearer credential. Public DNS still selects the worker dashboard; use +`--resolve worker.bstein.dev:443:192.168.22.50` for the private connection. -`GET /healthz` and `POST /api/generate` are the only gateway operations. Generate -accepts `model: qwen2.5:14b-instruct-q4_0`, a string `prompt`, `stream: false`, and -optional `num_predict`, `temperature`, `top_p`, and `seed` inside `options`. -Requests are capped at 128 KiB and 2048 output tokens. Concurrent LAN generations -receive 429 and can retry. Requests use only the local Ollama service; no hosted -fallback, model management, tools, sessions, image, or voice APIs are exposed. +Only authenticated `GET /healthz` and `POST /api/generate` are exposed under +`/local-model`. The initial model is pinned by tag and weight digest to the +existing titan-20 Jetson runtime. JSON schemas, bounded sampling controls, +8,192-token context admission, and 20-minute inference requests are supported. +The LAN API has no Switchyard, cloud fallback, session, tool, or management path. +Experimental CPU batch routes are disabled and their deployment is parked. -A missing/wrong token returns 401, an unavailable credential or local model -returns 503, and a non-LAN source returns 403. Flux bootstraps the Vault roles -and creates the token only if absent. After deliberate token rotation, roll the -model-gate through its Flux-tracked deployment revision to reload the injected -credential. +Flux bootstraps the Vault roles and creates the token only if absent. After +rotation, change the model-gate deployment revision through Git/Flux to reload +the init-injected credential. The client token grants no Vault or cluster access. This is the mental model and demonstration script for the operator instance at `triage.bstein.dev`. Read it once, then prove each section in the live UI. The diff --git a/services/hermes/kustomization.yaml b/services/hermes/kustomization.yaml index 7d0c7f0c..00c1ddad 100644 --- a/services/hermes/kustomization.yaml +++ b/services/hermes/kustomization.yaml @@ -55,6 +55,7 @@ configMapGenerator: - name: hermes-batch-api files: - batch_api.py=scripts/batch_api.py + - lan_generate.py=scripts/lan_generate.py options: disableNameSuffixHash: true - name: hermes-chat-oauth-templates diff --git a/services/hermes/model-gate-configmap.yaml b/services/hermes/model-gate-configmap.yaml index e9e02337..b3d4330f 100644 --- a/services/hermes/model-gate-configmap.yaml +++ b/services/hermes/model-gate-configmap.yaml @@ -18,7 +18,7 @@ data: import ssl import threading import time - import batch_api + import lan_generate from urllib.error import HTTPError, URLError from urllib.request import Request, urlopen @@ -33,10 +33,9 @@ data: IMAGE_NAMESPACE = os.environ.get("IMAGE_NAMESPACE", "hermes") IMAGE_DEPLOYMENT = os.environ.get("IMAGE_DEPLOYMENT", "hermes-local-image") LAN_API_TOKEN_FILE = Path(os.environ.get("LAN_API_TOKEN_FILE", "/vault/secrets/lan-api-token")) - LAN_MODEL = "qwen2.5:14b-instruct-q4_0" + LAN_MODEL = lan_generate.MODEL LAN_NETWORK = ipaddress.ip_network("192.168.22.0/24") - LAN_MAX_BODY = 131072 - _lan_inference = threading.BoundedSemaphore(1) + LAN_MAX_BODY = lan_generate.MAX_BODY KUBE_HOST = os.environ.get("KUBERNETES_SERVICE_HOST", "kubernetes.default.svc") KUBE_PORT = os.environ.get("KUBERNETES_SERVICE_PORT_HTTPS", "443") KUBE_TOKEN = Path("/var/run/secrets/kubernetes.io/serviceaccount/token") @@ -118,43 +117,6 @@ data: return False, last_error - def _normalize_lan_generate(body: bytes | None) -> bytes: - """Accept only the one supported stateless local generation operation.""" - - if not body: - raise ValueError("request body is required") - try: - payload = json.loads(body) - except (TypeError, ValueError, json.JSONDecodeError) as exc: - raise ValueError("invalid JSON") from exc - if not isinstance(payload, dict): - raise ValueError("JSON object is required") - if payload.get("model") != LAN_MODEL: - raise ValueError(f"model must be {LAN_MODEL}") - if payload.get("stream") is not False: - raise ValueError("stream must be false") - if not isinstance(payload.get("prompt"), str): - raise ValueError("prompt must be a string") - # Only stateless text and bounded sampling settings cross this boundary. - invalid = sorted(set(payload) - {"model", "prompt", "stream", "options"}) - if invalid: - raise ValueError(f"unsupported fields: {', '.join(invalid)}") - options = payload.get("options", {}) - if not isinstance(options, dict) or set(options) - {"num_predict", "temperature", "top_p", "seed"}: - raise ValueError("unsupported options") - count = options.get("num_predict", 256) - if type(count) is not int or not 1 <= count <= 2048: - raise ValueError("num_predict must be between 1 and 2048") - for key, upper in (("temperature", 2), ("top_p", 1)): - value = options.get(key, 1) - if type(value) not in (int, float) or not 0 <= value <= upper: - raise ValueError(f"invalid {key}") - if "seed" in options and (type(options["seed"]) is not int or not -1 <= options["seed"] <= 2147483647): - raise ValueError("invalid seed") - payload["options"] = {**options, "num_predict": count} - return json.dumps(payload, separators=(",", ":")).encode("utf-8") - - def _is_lan_client(client_ip: str) -> bool: """Validate the address appended by the trusted ingress proxy.""" try: @@ -204,18 +166,15 @@ data: def do_GET(self) -> None: """Expose authenticated readiness without model or session metadata.""" - if self.path not in ("/healthz", "/api/batch/models"): + if self.path != "/healthz": self._json(404, {"error": "not found"}) return if self._authorized(): - if self.path == "/api/batch/models": - self._json(*batch_api.models_response()) - return - self._json(200, {"status": "ok", "model": LAN_MODEL}) + self._json(*lan_generate.health()) def do_POST(self) -> None: """Validate and serialize bounded inference against local Ollama only.""" - if self.path not in ("/api/generate", "/api/batch/generate"): + if self.path != "/api/generate": self._json(404, {"error": "not found"}) return if not self._authorized(): @@ -224,38 +183,33 @@ data: if self.headers.get("Transfer-Encoding"): raise ValueError("Transfer-Encoding is unsupported") length = int(self.headers.get("Content-Length", "0")) - maximum = batch_api.MAX_BODY if self.path == "/api/batch/generate" else LAN_MAX_BODY + maximum = LAN_MAX_BODY if not 0 < length <= maximum: self._json(413, {"error": f"body must be between 1 and {maximum} bytes"}) return body = self.rfile.read(length) - if self.path == "/api/batch/generate": - self._json(*batch_api.generate(body)) - return - body = _normalize_lan_generate(body) - except ValueError as exc: - self._json(400, {"error": str(exc)}) + if len(body) != length: + raise ValueError("incomplete body") + except ValueError: + self._json(400, {"error": "invalid request framing"}) return - if not _lan_inference.acquire(blocking=False): - self._json(429, {"error": "local model is busy; retry later"}) + except TimeoutError: + self._json(408, {"error": "request body timed out"}) return - request = Request(f"{UPSTREAM_URL}/api/generate", data=body, headers={"Content-Type": "application/json"}, method="POST") - try: - try: - with urlopen(request, timeout=300) as response: - payload = json.load(response) - self._json(response.status, payload) - except HTTPError: - self._json(502, {"error": "local Ollama rejected the request"}) - except (TimeoutError, URLError, ValueError): - # This endpoint has one upstream only; never route externally. - self._json(503, {"error": "local Ollama unavailable"}) - finally: - _lan_inference.release() + self._json(*lan_generate.generate(body)) + + def send_error(self, code, message=None, explain=None) -> None: + """Do not reflect malformed methods, request targets, or parser input.""" + self._json(code, {"error": "unsupported or malformed HTTP request"}) + + def log_request(self, code="-", size="-") -> None: + """Only fixed route names and status codes enter container logs.""" + route = {"/healthz": "health", "/api/generate": "generate"}.get(getattr(self, "path", ""), "rejected") + print(f"model-gate-lan route={route} status={code}", flush=True) def log_message(self, format_string: str, *args) -> None: - # Deliberately log only transport metadata. Never log prompts or responses. - print(f"model-gate-lan {self.address_string()} {format_string % args}", flush=True) + # Parser errors may contain attacker-controlled text; omit it entirely. + return class Handler(BaseHTTPRequestHandler): diff --git a/services/hermes/model-gate-deployment.yaml b/services/hermes/model-gate-deployment.yaml index 34db7e3c..d3db3998 100644 --- a/services/hermes/model-gate-deployment.yaml +++ b/services/hermes/model-gate-deployment.yaml @@ -15,7 +15,7 @@ spec: template: metadata: annotations: - ai.bstein.dev/config-rev: "20260928-lan-batch-api-v2" + ai.bstein.dev/config-rev: "20260928-lan-native-v3" vault.hashicorp.com/agent-inject: "true" vault.hashicorp.com/agent-pre-populate-only: "true" vault.hashicorp.com/agent-init-first: "true" @@ -159,6 +159,8 @@ kind: Service metadata: name: hermes-model-gate-lan-api namespace: hermes + annotations: + traefik.ingress.kubernetes.io/service.serverstransport: hermes-hermes-model-gate-lan@kubernetescrd labels: app: hermes-model-gate spec: diff --git a/services/hermes/model-gate-lan-ingress.yaml b/services/hermes/model-gate-lan-ingress.yaml index 6c2e8af1..336b96e5 100644 --- a/services/hermes/model-gate-lan-ingress.yaml +++ b/services/hermes/model-gate-lan-ingress.yaml @@ -1,5 +1,16 @@ # services/hermes/model-gate-lan-ingress.yaml apiVersion: traefik.io/v1alpha1 +kind: ServersTransport +metadata: + name: hermes-model-gate-lan + namespace: hermes +spec: + forwardingTimeouts: + dialTimeout: 10s + responseHeaderTimeout: 1210s + idleConnTimeout: 90s +--- +apiVersion: traefik.io/v1alpha1 kind: Middleware metadata: name: hermes-model-gate-lan-allowlist diff --git a/services/hermes/scripts/lan_generate.py b/services/hermes/scripts/lan_generate.py new file mode 100644 index 00000000..19d341a9 --- /dev/null +++ b/services/hermes/scripts/lan_generate.py @@ -0,0 +1,173 @@ +#!/usr/bin/env python3 +"""Stateless, pinned Jetson inference without routing, proxies, or persistence.""" + +import json +import threading +import time +from urllib.error import HTTPError, URLError +from urllib.request import HTTPRedirectHandler, ProxyHandler, Request, build_opener + +UPSTREAM = "http://ollama.ai.svc.cluster.local:11434" +MODEL = "qwen2.5:14b-instruct-q4_0" +DIGEST = "5449194ff8035ccb13a6409a5814de6c8f9c39f555f429e383ae0fb7137001bd" +RUNTIME = "0.13.5" +MAX_BODY = 131072 +CONTEXT = 8192 +TIMEOUT = 1200 +_inference = threading.BoundedSemaphore(1) + + +class NoRedirect(HTTPRedirectHandler): + """Reject redirects before another host can receive a request.""" + + def redirect_request(self, req, fp, code, msg, headers, newurl): + return None + + +_http = build_opener(ProxyHandler({}), NoRedirect()) + + +class InputTooLarge(ValueError): + """Distinguish context admission failures from malformed requests.""" + + +def _finite_json(raw): + """Reject nonstandard constants without including input in diagnostics.""" + def invalid_constant(_value): + raise ValueError("nonfinite JSON value") + + return json.loads(raw, parse_constant=invalid_constant) + + +def _local_schema(value): + """Allow bounded schemas without external references or dynamic resolution.""" + if isinstance(value, dict): + for key, child in value.items(): + if key in ("$ref", "$dynamicRef") and ( + not isinstance(child, str) or not child.startswith("#") + ): + raise ValueError("external schema references are unsupported") + _local_schema(child) + elif isinstance(value, list): + for child in value: + _local_schema(child) + + +def normalize(body): + """Return a bounded native generate payload; never trim source material.""" + if not body or len(body) > MAX_BODY: + raise InputTooLarge("request body exceeds limit") + payload = _finite_json(body) + if not isinstance(payload, dict) or set(payload) - {"model", "prompt", "stream", "format", "options"}: + raise ValueError("unsupported fields") + if payload.get("model") != MODEL or payload.get("stream") is not False: + raise ValueError("exact model and stream=false are required") + prompt = payload.get("prompt") + if not isinstance(prompt, str) or not prompt.strip(): + raise ValueError("nonempty prompt is required") + if "format" in payload: + shape = payload["format"] + if shape != "json": + if not isinstance(shape, dict) or shape.get("type") != "object": + raise ValueError("format must be json or an object JSON schema") + if len(json.dumps(shape).encode()) > 32768: + raise InputTooLarge("schema exceeds limit") + _local_schema(shape) + options = payload.get("options", {}) + allowed = {"num_ctx", "num_predict", "temperature", "top_p", "top_k", "seed"} + if not isinstance(options, dict) or set(options) - allowed: + raise ValueError("unsupported options") + options = {"num_ctx": CONTEXT, "num_predict": 256, "temperature": 0, + "top_p": 1, "top_k": 40, "seed": 0, **options} + for key, lower, upper in (("num_ctx", CONTEXT, CONTEXT), ("num_predict", 1, 2048), + ("top_k", 1, 100), ("seed", 0, 2147483647)): + if type(options[key]) is not int or not lower <= options[key] <= upper: + raise ValueError("unsupported integer option") + for key, upper in (("temperature", 2), ("top_p", 1)): + if type(options[key]) not in (int, float) or not 0 <= options[key] <= upper: + raise ValueError("unsupported sampling option") + # UTF-8 bytes upper-bound this model's byte-level tokenizer. Reserve the + # full output budget plus 1024 template tokens before calling Ollama. + if len(prompt.encode("utf-8")) + options["num_predict"] + 1024 > CONTEXT: + raise InputTooLarge("prompt exceeds conservative context budget") + return {**payload, "options": options} + + +def _request(path, payload=None, timeout=10): + """Call only the fixed in-cluster backend and bound the response body.""" + data = None if payload is None else json.dumps(payload, allow_nan=False).encode() + request = Request(UPSTREAM + path, data=data, headers={"Content-Type": "application/json"}) + with _http.open(request, timeout=timeout) as response: + raw = response.read(4 * MAX_BODY + 1) + if len(raw) > 4 * MAX_BODY: + raise ValueError("upstream response exceeds limit") + return _finite_json(raw) + + +def verify_model(): + """Check the runtime and local manifest before sending any prompt text.""" + if _request("/api/version").get("version") != RUNTIME: + raise ValueError("runtime changed") + matches = [item for item in _request("/api/tags")["models"] if item.get("name") == MODEL] + if len(matches) != 1 or matches[0].get("digest", "").removeprefix("sha256:") != DIGEST: + raise ValueError("pinned model unavailable") + return {"model": MODEL, "model_digest": DIGEST, "runtime": RUNTIME, + "placement": "titan-20/Jetson-Xavier-16GB", "context_tokens": CONTEXT, + "max_output_tokens": 2048, "timeout_seconds": TIMEOUT, + "concurrency": 1, "fallback": None, "protocol_version": 1} + + +def health(): + """Report authenticated model readiness, not merely process liveness.""" + try: + return 200, {"status": "ok", **verify_model()} + except (OSError, URLError, ValueError, KeyError, TypeError, AttributeError): + return 503, {"error": "pinned local model unavailable; no fallback"} + + +def generate(body): + """Reject busy requests and return native output with reproducible metadata.""" + try: + payload = normalize(body) + except InputTooLarge: + return 413, {"error": "request exceeds body, schema, or context budget"} + except (ValueError, TypeError, RecursionError): + return 400, {"error": "invalid inference request; consult endpoint contract"} + if not _inference.acquire(blocking=False): + return 429, {"error": "local inference busy; retry explicitly"} + started = time.monotonic() + try: + provenance = verify_model() + result = _request("/api/generate", payload, timeout=max(1, TIMEOUT - (time.monotonic() - started))) + if result.get("model") != MODEL or result.get("done") is not True: + raise ValueError("model mismatch or incomplete response") + if result.get("done_reason") == "length": + return 422, {"error": "output token budget exhausted; result is incomplete"} + if not isinstance(result.get("response"), str): + raise ValueError("missing output") + if "format" in payload: + try: + _finite_json(result["response"]) + except ValueError: + return 422, {"error": "model returned invalid JSON"} + # Project the response so context token IDs, debug fields, or upstream + # error text can never become an accidental history or logging channel. + fields = ("model", "created_at", "response", "done", "done_reason", + "total_duration", "load_duration", "prompt_eval_count", + "prompt_eval_duration", "eval_count", "eval_duration") + output = {key: result[key] for key in fields if key in result} + output["inference_provenance"] = { + **provenance, "options": payload["options"], "backend_api": "/api/generate", + "wall_seconds": round(time.monotonic() - started, 3), + } + return 200, output + except TimeoutError: + return 504, {"error": "local inference timed out; no fallback"} + except HTTPError as exc: + if exc.code in (429, 503): + return 503, {"error": "local inference capacity unavailable; no fallback"} + return 502, {"error": "local inference failed; no fallback"} + except (OSError, URLError, ValueError, KeyError, TypeError, AttributeError, RecursionError): + return 503, {"error": "pinned local inference unavailable; no fallback"} + finally: + _inference.release() diff --git a/services/logging/fluent-bit-helmrelease.yaml b/services/logging/fluent-bit-helmrelease.yaml index 34b755ff..1da3af10 100644 --- a/services/logging/fluent-bit-helmrelease.yaml +++ b/services/logging/fluent-bit-helmrelease.yaml @@ -74,7 +74,7 @@ spec: Name tail Tag kube.* Path /var/log/containers/*.log - Exclude_Path /var/log/containers/*_POD_*.log + Exclude_Path /var/log/containers/*_POD_*.log,/var/log/containers/ollama-*_ai_ollama-*.log,/var/log/containers/hermes-model-gate-*_hermes_model-gate-*.log Parser cri Mem_Buf_Limit 50MB Skip_Long_Lines On diff --git a/testing/tests/test_hermes_model_gate.py b/testing/tests/test_hermes_model_gate.py index ced7415c..6ce56ddc 100644 --- a/testing/tests/test_hermes_model_gate.py +++ b/testing/tests/test_hermes_model_gate.py @@ -6,10 +6,11 @@ import json from pathlib import Path import yaml -from services.hermes.scripts import batch_api +from services.hermes.scripts import batch_api, lan_generate import sys sys.modules.setdefault("batch_api", batch_api) +sys.modules.setdefault("lan_generate", lan_generate) HERMES = Path(__file__).parents[2] / "services/hermes" diff --git a/testing/tests/test_hermes_model_gate_lan.py b/testing/tests/test_hermes_model_gate_lan.py index 005f5bde..6644fd38 100644 --- a/testing/tests/test_hermes_model_gate_lan.py +++ b/testing/tests/test_hermes_model_gate_lan.py @@ -9,14 +9,14 @@ from urllib.error import URLError import pytest import yaml -from services.hermes.scripts import batch_api +from services.hermes.scripts import lan_generate import sys -sys.modules.setdefault("batch_api", batch_api) +sys.modules.setdefault("lan_generate", lan_generate) @pytest.fixture -def gateway(tmp_path): +def gateway(tmp_path, monkeypatch): """Run the shipped handler with an isolated token and intercepted upstream.""" path = Path(__file__).parents[2] / "services/hermes/model-gate-configmap.yaml" namespace = {"__name__": "test_gateway"} @@ -26,13 +26,17 @@ def gateway(tmp_path): namespace["LAN_API_TOKEN_FILE"] = token calls = [] - def upstream(request, timeout): - calls.append(request) - response = io.BytesIO(b'{"response":"LAN_OK","done":true}') - response.status = 200 - return response + def upstream(path, payload=None, timeout=10): + calls.append((path, payload, timeout)) + if path == "/api/version": + return {"version": lan_generate.RUNTIME} + if path == "/api/tags": + return {"models": [{"name": lan_generate.MODEL, "digest": lan_generate.DIGEST}]} + return {"model": lan_generate.MODEL, "response": "LAN_OK", "done": True, + "context": [1, 2, 3], "thinking": "never forward this"} - namespace["urlopen"] = upstream + monkeypatch.setattr(lan_generate, "_request", upstream) + monkeypatch.setattr(lan_generate, "_inference", threading.BoundedSemaphore(1)) server = namespace["ThreadingHTTPServer"](("127.0.0.1", 0), namespace["LanHandler"]) thread = threading.Thread(target=server.serve_forever, daemon=True) thread.start() @@ -72,16 +76,20 @@ def test_health_requires_a_readable_credential(gateway): assert request()[0] == 200 namespace["LAN_API_TOKEN_FILE"].unlink() assert request()[0] == 503 - assert not calls + assert len(calls) == 2 def test_generation_uses_only_local_upstream_without_forwarding_token(gateway): namespace, calls, request = gateway body = json.dumps({"model": namespace["LAN_MODEL"], "prompt": "Say LAN_OK", "stream": False}) - assert request("POST", "/api/generate", body) == (200, {"response": "LAN_OK", "done": True}) - assert calls[0].full_url == "http://ollama.ai.svc.cluster.local:11434/api/generate" - assert "Authorization" not in calls[0].headers - assert json.loads(calls[0].data)["options"]["num_predict"] == 256 + status, result = request("POST", "/api/generate", body) + assert status == 200 and result["response"] == "LAN_OK" + assert "context" not in result and "thinking" not in result + assert result["inference_provenance"]["model_digest"] == lan_generate.DIGEST + assert [call[0] for call in calls] == ["/api/version", "/api/tags", "/api/generate"] + assert calls[-1][1]["options"]["num_predict"] == 256 + assert calls[-1][1]["options"]["num_ctx"] == 8192 + assert 1190 < calls[-1][2] <= 1200 @pytest.mark.parametrize("changes", [ @@ -104,20 +112,20 @@ def test_invalid_request_lengths_are_rejected_before_read(gateway, length, statu assert not calls -def test_outage_and_concurrent_request_fail_without_fallback(gateway): +def test_outage_and_concurrent_request_fail_without_fallback(gateway, monkeypatch): namespace, calls, request = gateway body = json.dumps({"model": namespace["LAN_MODEL"], "prompt": "hello", "stream": False}) - namespace["_lan_inference"].acquire() + lan_generate._inference.acquire() assert request("POST", "/api/generate", body)[0] == 429 - namespace["_lan_inference"].release() + lan_generate._inference.release() def unavailable(*args, **kwargs): raise URLError("local service unavailable") - namespace["urlopen"] = unavailable + monkeypatch.setattr(lan_generate, "_request", unavailable) assert request("POST", "/api/generate", body)[0] == 503 - assert namespace["_lan_inference"].acquire(timeout=1) - namespace["_lan_inference"].release() + assert lan_generate._inference.acquire(timeout=1) + lan_generate._inference.release() assert not calls @@ -128,12 +136,53 @@ def test_model_management_and_agent_routes_are_not_exposed(gateway, path): assert not calls -def test_batch_routes_share_lan_authentication(gateway, monkeypatch): +@pytest.mark.parametrize("path", ["/api/batch/models", "/api/batch/generate", "/terminal", "/api/tags"]) +def test_only_initial_pilot_operations_are_exposed(gateway, path): _, calls, request = gateway - monkeypatch.setattr(batch_api, "models_response", lambda: (200, {"models": []})) - monkeypatch.setattr(batch_api, "generate", lambda body: (200, {"batch": True})) - assert request(path="/api/batch/models", Authorization="")[0] == 401 - assert request("POST", "/api/batch/generate", "{}", Authorization="")[0] == 401 - assert request(path="/api/batch/models") == (200, {"models": []}) - assert request("POST", "/api/batch/generate", "{}") == (200, {"batch": True}) + assert request(path=path)[0] == 404 + assert request("POST", path, "{}")[0] == 404 + assert not calls + + +def test_logs_and_errors_never_echo_request_material(gateway, capsys): + namespace, calls, request = gateway + marker = "PRIVATE_CANARY_TEST_ONLY" + assert request(path="/unknown?prompt=" + marker)[0] == 404 + body = json.dumps({"model": namespace["LAN_MODEL"], "prompt": marker, + "stream": False, marker: "unknown field"}) + status, error = request("POST", "/api/generate", body) + assert status == 400 and marker not in json.dumps(error) + assert marker not in capsys.readouterr().out + assert not calls + + +def test_context_admission_does_not_truncate_or_call_backend(gateway): + namespace, calls, request = gateway + body = json.dumps({"model": namespace["LAN_MODEL"], "prompt": "x" * 8000, "stream": False}) + assert request("POST", "/api/generate", body)[0] == 413 + assert not calls + + +def test_structured_output_and_exact_model_pins(gateway, monkeypatch): + _, calls, request = gateway + original = lan_generate._request + schema = {"type": "object", "properties": {"status": {"type": "string"}}, "required": ["status"]} + body = json.dumps({"model": lan_generate.MODEL, "prompt": "Return JSON status ok", + "stream": False, "format": schema}) + + def structured(path, payload=None, timeout=10): + result = original(path, payload, timeout) + if path == "/api/generate": + result["response"] = '{"status":"ok"}' + return result + + monkeypatch.setattr(lan_generate, "_request", structured) + status, result = request("POST", "/api/generate", body) + assert status == 200 and json.loads(result["response"]) == {"status": "ok"} + assert calls[-1][1]["format"] == schema + calls.clear() + monkeypatch.setattr(lan_generate, "RUNTIME", "mismatch") + # The fixture follows RUNTIME, so force an unexpected version independently. + monkeypatch.setattr(lan_generate, "_request", lambda *a, **k: {"version": "different"}) + assert request("POST", "/api/generate", body)[0] == 503 assert not calls diff --git a/testing/tests/test_lan_generate.py b/testing/tests/test_lan_generate.py new file mode 100644 index 00000000..3844b449 --- /dev/null +++ b/testing/tests/test_lan_generate.py @@ -0,0 +1,85 @@ +"""Local-only transport and failure tests; no real model or credentials used.""" + +import io +import json +from urllib.error import HTTPError +from urllib.request import ProxyHandler + +import pytest + +from services.hermes.scripts import lan_generate as api + + +def body(**changes): + """Build a synthetic native request, allowing invalid-field variants.""" + return json.dumps({"model": api.MODEL, "stream": False, "prompt": "Synthetic only", **changes}).encode() + + +def test_transport_ignores_proxies_and_rejects_redirects(monkeypatch): + """Poisoned proxy variables and redirects cannot choose another backend.""" + monkeypatch.setenv("HTTP_PROXY", "http://unapproved.invalid:8080") + assert not any(isinstance(h, ProxyHandler) and h.proxies for h in api._http.handlers) + redirect = next(h for h in api._http.handlers if isinstance(h, api.NoRedirect)) + assert redirect.redirect_request(None, None, 307, "redirect", {}, "https://unapproved.invalid") is None + captured = [] + + class Opener: + """Capture the single actual HTTP boundary and emulate an outage.""" + + def open(self, request, timeout): + captured.append((request, timeout)) + raise HTTPError(request.full_url, 307, "PRIVATE_ERROR_TEXT", {}, io.BytesIO(b"private")) + + monkeypatch.setattr(api, "_http", Opener()) + status, result = api.generate(body()) + assert status == 502 and "PRIVATE_ERROR_TEXT" not in json.dumps(result) + assert len(captured) == 1 + assert captured[0][0].full_url == api.UPSTREAM + "/api/version" + assert "Authorization" not in captured[0][0].headers + + +@pytest.mark.parametrize("changes", [ + {"options": {"num_ctx": 32768}}, {"options": {"seed": True}}, + {"options": {"num_predict": True}}, {"format": {"type": "object", "$ref": "https://unapproved.invalid"}}, + {"format": {"type": "array"}}, {"format": "xml"}, {"prompt": ""}, + {"system": "new system prompt"}, {"options": {"temperature": float("inf")}}, +]) +def test_invalid_input_fails_before_metadata_or_inference(monkeypatch, changes): + monkeypatch.setattr(api, "_request", lambda *a, **k: pytest.fail("backend must not be called")) + assert api.generate(body(**changes))[0] == 400 + + +@pytest.mark.parametrize("mutation,expected", [ + ({"model": "different"}, 503), ({"done": False}, 503), + ({"done_reason": "length"}, 422), ({"response": "not JSON"}, 422), +]) +def test_incomplete_or_substituted_output_is_not_success(monkeypatch, mutation, expected): + monkeypatch.setattr(api, "verify_model", lambda: {"model_digest": api.DIGEST}) + monkeypatch.setattr(api, "_request", lambda *a, **k: { + "model": api.MODEL, "done": True, "response": '{"status":"ok"}', **mutation, + }) + assert api.generate(body(format="json"))[0] == expected + + +def test_changed_digest_prevents_prompt_submission(monkeypatch): + calls = [] + + def upstream(path, *args, **kwargs): + calls.append(path) + return ({"version": api.RUNTIME} if path == "/api/version" else + {"models": [{"name": api.MODEL, "digest": "changed"}]}) + + monkeypatch.setattr(api, "_request", upstream) + assert api.generate(body())[0] == 503 + assert calls == ["/api/version", "/api/tags"] + + +def test_timeout_releases_capacity_and_redacts_error(monkeypatch): + def timeout(*args, **kwargs): + raise TimeoutError("PRIVATE_PROMPT") + + monkeypatch.setattr(api, "_request", timeout) + status, result = api.generate(body()) + assert status == 504 and "PRIVATE_PROMPT" not in json.dumps(result) + assert api._inference.acquire(blocking=False) + api._inference.release()