From ec9fe97fc112d0ad81959807a098e26a465ea9e8 Mon Sep 17 00:00:00 2001 From: Vantio Date: Wed, 30 Sep 2026 13:54:53 -0400 Subject: [PATCH] fix(optics): record http.client status and keep the ingest URL path getresponse stores the HTTP status on the pending call and clears that pending state when it throws, so the next response is not attached to the failed request. VANTIO_INGEST_URL keeps its path. A worker thread's call stays on the shield run record. Rebased onto main after enforcement removal. Optics stays observe-only. CANDIDATE_ONLY. Not published. --- .../tests/test_observe_bugs.py | 104 ++++++++++++++++++ .../tests/test_outcome_clarity.py | 38 ++++--- .../vantio/_http_observe.py | 57 ++++++++-- packages/vantio-cli/bin/ingest-url.cjs | 30 +++++ packages/vantio-cli/bin/interceptor.cjs | 4 +- packages/vantio-cli/package.json | 2 +- packages/vantio-cli/test/ingest-url.test.js | 17 +++ 7 files changed, 229 insertions(+), 23 deletions(-) create mode 100644 packages/vantio-agent-sdk-py/tests/test_observe_bugs.py create mode 100644 packages/vantio-cli/bin/ingest-url.cjs create mode 100644 packages/vantio-cli/test/ingest-url.test.js diff --git a/packages/vantio-agent-sdk-py/tests/test_observe_bugs.py b/packages/vantio-agent-sdk-py/tests/test_observe_bugs.py new file mode 100644 index 00000000..bd0af3df --- /dev/null +++ b/packages/vantio-agent-sdk-py/tests/test_observe_bugs.py @@ -0,0 +1,104 @@ +"""Observation bugs: http.client pending, worker-thread records, ingest URL path.""" + +from __future__ import annotations + +import json +import os +import tempfile +import threading +import unittest +import urllib.request +from http.server import BaseHTTPRequestHandler, HTTPServer +from pathlib import Path + +from vantio import shield + + +class _Server(HTTPServer): + def __init__(self) -> None: + super().__init__(("127.0.0.1", 0), _Handler) + + +class _Handler(BaseHTTPRequestHandler): + def do_GET(self) -> None: + if self.path.startswith("/boom"): + self.connection.close() + return + body = b"ok" + self.send_response(201) + self.send_header("Content-Length", str(len(body))) + self.end_headers() + self.wfile.write(body) + + def log_message(self, *_args: object) -> None: + return + + +class ObserveBugTests(unittest.IsolatedAsyncioTestCase): + def setUp(self) -> None: + self._home = tempfile.mkdtemp() + self._saved = {k: os.environ.get(k) for k in ("VANTIO_HOME", "VANTIO_EXTRA_LLM_HOSTS", "VANTIO_API_KEY", "VANTIO_INGEST_URL")} + os.environ["VANTIO_HOME"] = self._home + os.environ["VANTIO_EXTRA_LLM_HOSTS"] = "127.0.0.1" + os.environ.pop("VANTIO_API_KEY", None) + os.environ.pop("VANTIO_INGEST_URL", None) + self._server = _Server() + self._thread = threading.Thread(target=self._server.serve_forever, daemon=True) + self._thread.start() + + def tearDown(self) -> None: + self._server.shutdown() + for key, value in self._saved.items(): + if value is None: + os.environ.pop(key, None) + else: + os.environ[key] = value + + async def test_getresponse_throw_clears_pending_and_does_not_attach_the_next_status(self) -> None: + import http.client + + port = self._server.server_address[1] + async with shield(trace_id="http-pending"): + boom = http.client.HTTPConnection("127.0.0.1", port, timeout=2) + boom.request("GET", "/boom") + with self.assertRaises(Exception): + boom.getresponse() + self.assertIsNone(getattr(boom, "_vantio_pending", None)) + ok = http.client.HTTPConnection("127.0.0.1", port, timeout=2) + ok.request("GET", "/ok") + resp = ok.getresponse() + self.assertEqual(resp.status, 201) + resp.read() + log = json.loads((Path(self._home) / "runs" / "http-pending.json").read_text(encoding="utf-8")) + by_path = {c.get("path"): c for c in log["calls"]} + self.assertEqual(by_path["/ok"]["status"], 201) + boom_call = by_path.get("/boom") + if boom_call is not None: + self.assertNotEqual(boom_call.get("status"), 201) + + async def test_worker_thread_call_reaches_the_shield_record(self) -> None: + port = self._server.server_address[1] + url = f"http://127.0.0.1:{port}/ok" + + async def main() -> None: + async with shield(trace_id="worker-trace"): + box: dict[str, object] = {} + + def work() -> None: + urllib.request.urlopen(url, timeout=2).read() + box["ok"] = True + + thread = threading.Thread(target=work) + thread.start() + thread.join(timeout=3) + urllib.request.urlopen(url, timeout=2).read() + self.assertTrue(box.get("ok")) + + await main() + log = json.loads((Path(self._home) / "runs" / "worker-trace.json").read_text(encoding="utf-8")) + self.assertEqual(log["trace_id"], "worker-trace") + self.assertGreaterEqual(len(log["calls"]), 2) + + +if __name__ == "__main__": + unittest.main() diff --git a/packages/vantio-agent-sdk-py/tests/test_outcome_clarity.py b/packages/vantio-agent-sdk-py/tests/test_outcome_clarity.py index ce6ed8b1..c92e71bd 100644 --- a/packages/vantio-agent-sdk-py/tests/test_outcome_clarity.py +++ b/packages/vantio-agent-sdk-py/tests/test_outcome_clarity.py @@ -175,7 +175,13 @@ class OutcomeMappingTests(unittest.TestCase): def test_http_lines_match_the_status_table(self) -> None: for status, (label, response, category, _token, _ok) in EXPECTED.items(): self.assertEqual(http_outcome_label(status), label) - self.assertEqual(http_response_text(status), response) + if status == 422: + self.assertIn( + http_response_text(status), + ("HTTP 422 Unprocessable Entity", "HTTP 422 Unprocessable Content"), + ) + else: + self.assertEqual(http_response_text(status), response) self.assertNotIn("Application error", label) def test_other_4xx_stays_a_rejection_without_a_new_token(self) -> None: @@ -467,7 +473,13 @@ def _assert_http_call(self, call: dict, status: int, mediation: str) -> None: self.assertEqual(call["opticsLabel"], "Successful") self.assertEqual(call["applicationOutcomeLabel"], label) self.assertEqual(call["applicationLabel"], label) - self.assertEqual(call["providerResponse"], response) + if status == 422: + self.assertIn( + call["providerResponse"], + ("HTTP 422 Unprocessable Entity", "HTTP 422 Unprocessable Content"), + ) + else: + self.assertEqual(call["providerResponse"], response) self.assertEqual(call["providerResponseLabel"], "Upstream response") self.assertEqual(call["upstreamService"], "127.0.0.1") self.assertNotIn("providerName", call) @@ -647,20 +659,20 @@ def boom(*args, **kwargs): self.assertNotIn("error", wrapped) client_calls = [c for c in data["calls"] if c.get("mediation") == "python_http_client"] self.assertEqual(len(client_calls), 1) - unavailable = client_calls[0] - self.assertIs(unavailable["ok"], True) - self.assertNotIn("status", unavailable) - self.assertEqual(unavailable["opticsStatus"], "SUCCESS") - self.assertEqual(unavailable["applicationStatus"], "UNAVAILABLE") - self.assertEqual(unavailable["applicationOutcomeLabel"], "Provider outcome unavailable") - self.assertEqual(unavailable["providerResponse"], "No HTTP response") - self.assertEqual(unavailable["nextActionCategory"], "inspection") + observed = client_calls[0] + self.assertIs(observed["ok"], True) + self.assertEqual(observed["status"], 200) + self.assertEqual(observed["opticsStatus"], "SUCCESS") + self.assertEqual(observed["applicationStatus"], "SUCCESS") + self.assertEqual(observed["applicationOutcomeLabel"], "Successful") + self.assertEqual(observed["providerResponse"], "HTTP 200 OK") + self.assertEqual(observed["nextActionCategory"], "inspection") self.assertEqual( - customer_view_lines(unavailable), + customer_view_lines(observed), [ "Optics status: Successful", - "Observed outcome: Provider outcome unavailable", - "Upstream response: No HTTP response", + "Observed outcome: Successful", + "Upstream response: HTTP 200 OK", ], ) diff --git a/packages/vantio-agent-sdk-py/vantio/_http_observe.py b/packages/vantio-agent-sdk-py/vantio/_http_observe.py index 6b844204..f37c66b5 100644 --- a/packages/vantio-agent-sdk-py/vantio/_http_observe.py +++ b/packages/vantio-agent-sdk-py/vantio/_http_observe.py @@ -115,6 +115,7 @@ _orig_create_connection: Any = None _orig_ssl_connect: Any = None _orig_http_request: Any = None +_orig_http_getresponse: Any = None _orig_http_putrequest: Any = None _orig_urllib3_request: Any = None _orig_pycurl_curl: Any = None @@ -490,7 +491,7 @@ def _record( action: str, mediation: str, **extra: Any, -) -> None: +) -> dict[str, Any]: rec = { "hostname": hostname, "provider": extra.pop("provider", "other"), @@ -506,6 +507,7 @@ def _record( rec["opticsLabel"] = _human_status(optics) apply_customer_outcome(rec) _append(rec) + return rec def _dispatch_gate( @@ -1101,10 +1103,14 @@ def _observe_http_client_request( resp = _http_orig( _orig_http_request, self, method, url, body, headers or {}, encode_chunked=encode_chunked ) - if record_send: + if record_send and getattr(self, "_vantio_pending", None) is None: action = "OBSERVED" - _record(hostname, action, "python_http_client", method=method_s, path=path, ok=True, - duration_ms=int((time.time() - t0) * 1000)) + rec = _record( + hostname, action, "python_http_client", method=method_s, path=path, ok=True, + duration_ms=int((time.time() - t0) * 1000), + ) + _, _, nbytes = _body_to_text(body) + self._vantio_pending = {"rec": rec, "path": path, "t0": t0, "request_bytes": nbytes} return resp except Exception as exc: _record_http_exception( @@ -1134,33 +1140,68 @@ def _observe_http_client_putrequest( kind, _payload, _redactions, record_send = _dispatch_gate( hostname, port, path, None, "python_http_client" ) - if kind != "pass" and record_send: + if kind != "pass" and record_send and getattr(self, "_vantio_pending", None) is None: action = "OBSERVED" - _record(hostname, action, "python_http_client", method=str(method or "GET").upper(), path=path) + rec = _record( + hostname, action, "python_http_client", + method=str(method or "GET").upper(), path=path, + ) + self._vantio_pending = {"rec": rec, "path": path, "t0": time.time(), "request_bytes": None} return _http_orig(_orig_http_putrequest, self, method, url, skip_host, skip_accept_encoding) +def _observe_http_client_getresponse(self: Any, *args: Any, **kwargs: Any) -> Any: + if _http_owns(): + return _orig_http_getresponse(self, *args, **kwargs) + pending = getattr(self, "_vantio_pending", None) + try: + resp = _orig_http_getresponse(self, *args, **kwargs) + except Exception: + self._vantio_pending = None + raise + self._vantio_pending = None + if pending and isinstance(pending.get("rec"), dict): + status = getattr(resp, "status", None) + rec = pending["rec"] + rec["status"] = _normalize_http_status(status) + rec["ok"] = _ok_for_http_status(status) + rec["applicationStatus"] = _application_status(status) + try: + cl = resp.getheader("Content-Length") if hasattr(resp, "getheader") else None + if cl is not None and str(cl).strip() != "": + rec["bytes"] = int(cl) + except (TypeError, ValueError): + pass + apply_customer_outcome(rec) + return resp + + def _install_http_client() -> None: - global _orig_http_request, _orig_http_putrequest + global _orig_http_request, _orig_http_putrequest, _orig_http_getresponse if _orig_http_request is not None: return _orig_http_request = http.client.HTTPConnection.request _orig_http_putrequest = http.client.HTTPConnection.putrequest + _orig_http_getresponse = http.client.HTTPConnection.getresponse http.client.HTTPConnection.request = _observe_http_client_request # type: ignore[assignment] http.client.HTTPConnection.putrequest = _observe_http_client_putrequest # type: ignore[assignment] + http.client.HTTPConnection.getresponse = _observe_http_client_getresponse # type: ignore[assignment] def _uninstall_http_client() -> None: - global _orig_http_request, _orig_http_putrequest + global _orig_http_request, _orig_http_putrequest, _orig_http_getresponse try: if _orig_http_request is not None: http.client.HTTPConnection.request = _orig_http_request if _orig_http_putrequest is not None: http.client.HTTPConnection.putrequest = _orig_http_putrequest + if _orig_http_getresponse is not None: + http.client.HTTPConnection.getresponse = _orig_http_getresponse except Exception: pass _orig_http_request = None _orig_http_putrequest = None + _orig_http_getresponse = None def _observe_urllib3_urlopen(self: Any, method: Any, url: Any, *args: Any, **kwargs: Any) -> Any: diff --git a/packages/vantio-cli/bin/ingest-url.cjs b/packages/vantio-cli/bin/ingest-url.cjs new file mode 100644 index 00000000..9b9a1b0b --- /dev/null +++ b/packages/vantio-cli/bin/ingest-url.cjs @@ -0,0 +1,30 @@ +"use strict"; + +// Node and Python both keep the path on VANTIO_INGEST_URL. The origin alone is not enough. +function parseIngestUrl(raw) { + if (raw == null || String(raw).trim() === "") { + return { ok: true, href: "https://vantio.ai", publicHost: true }; + } + let url; + try { + url = new URL(String(raw).trim()); + } catch { + return { ok: false, href: "", publicHost: false, reason: "VANTIO_INGEST_URL is not a valid URL" }; + } + if (url.protocol !== "http:" && url.protocol !== "https:") { + return { ok: false, href: "", publicHost: false, reason: "VANTIO_INGEST_URL must use http or https" }; + } + if (!url.hostname) { + return { ok: false, href: "", publicHost: false, reason: "VANTIO_INGEST_URL has no host" }; + } + const path = url.pathname && url.pathname !== "/" ? url.pathname.replace(/\/+$/, "") : ""; + const href = `${url.origin}${path}`; + const host = url.hostname.toLowerCase(); + return { + ok: true, + href, + publicHost: host === "vantio.ai" || host === "www.vantio.ai", + }; +} + +module.exports = { parseIngestUrl }; diff --git a/packages/vantio-cli/bin/interceptor.cjs b/packages/vantio-cli/bin/interceptor.cjs index 5f5a4e69..0f24fd70 100755 --- a/packages/vantio-cli/bin/interceptor.cjs +++ b/packages/vantio-cli/bin/interceptor.cjs @@ -45,7 +45,9 @@ const c = { cyan: USE_COLOR ? "\x1b[36m" : "", }; -const INGEST_URL = process.env.VANTIO_INGEST_URL || "https://vantio.ai"; +const { parseIngestUrl } = require("./ingest-url.cjs"); +const _parsedIngest = parseIngestUrl(process.env.VANTIO_INGEST_URL); +const INGEST_URL = _parsedIngest.ok ? _parsedIngest.href : (process.env.VANTIO_INGEST_URL || "https://vantio.ai"); // Keep the path. Do not reduce the URL to its origin. function isPublicCloudHost(raw) { try { diff --git a/packages/vantio-cli/package.json b/packages/vantio-cli/package.json index 61b7e4d5..041dde39 100644 --- a/packages/vantio-cli/package.json +++ b/packages/vantio-cli/package.json @@ -37,7 +37,7 @@ "LICENSE" ], "scripts": { - "lint": "node --check bin/vantio.js && node --check bin/interceptor.cjs && node --check bin/telemetry.cjs && node --check bin/llm-hosts.cjs && node --check bin/optics-cx.cjs", + "lint": "node --check bin/vantio.js && node --check bin/interceptor.cjs && node --check bin/telemetry.cjs && node --check bin/llm-hosts.cjs && node --check bin/optics-cx.cjs && node --check bin/ingest-url.cjs", "test": "node --test" }, "engines": { diff --git a/packages/vantio-cli/test/ingest-url.test.js b/packages/vantio-cli/test/ingest-url.test.js new file mode 100644 index 00000000..cc753c4e --- /dev/null +++ b/packages/vantio-cli/test/ingest-url.test.js @@ -0,0 +1,17 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { parseIngestUrl } from "../bin/ingest-url.cjs"; + +test("Node keeps the VANTIO_INGEST_URL path", () => { + const parsed = parseIngestUrl("http://127.0.0.1:9/custom/base/"); + assert.equal(parsed.ok, true); + assert.equal(parsed.href, "http://127.0.0.1:9/custom/base"); + assert.equal(parsed.href.includes("/custom/base"), true); + assert.notEqual(parsed.href, "http://127.0.0.1:9"); +}); + +test("a missing ingest URL stays on the public origin", () => { + const parsed = parseIngestUrl(""); + assert.equal(parsed.href, "https://vantio.ai"); + assert.equal(parsed.publicHost, true); +});