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 0000000..bd0af3d --- /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 ce6ed8b..c92e71b 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 6b84420..f37c66b 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 0000000..9b9a1b0 --- /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 5f5a4e6..0f24fd7 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 61b7e4d..041dde3 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 0000000..cc753c4 --- /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); +});