Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
104 changes: 104 additions & 0 deletions packages/vantio-agent-sdk-py/tests/test_observe_bugs.py
Original file line number Diff line number Diff line change
@@ -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()
38 changes: 25 additions & 13 deletions packages/vantio-agent-sdk-py/tests/test_outcome_clarity.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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",
],
)

Expand Down
57 changes: 49 additions & 8 deletions packages/vantio-agent-sdk-py/vantio/_http_observe.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -490,7 +491,7 @@ def _record(
action: str,
mediation: str,
**extra: Any,
) -> None:
) -> dict[str, Any]:
rec = {
"hostname": hostname,
"provider": extra.pop("provider", "other"),
Expand All @@ -506,6 +507,7 @@ def _record(
rec["opticsLabel"] = _human_status(optics)
apply_customer_outcome(rec)
_append(rec)
return rec


def _dispatch_gate(
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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:
Expand Down
30 changes: 30 additions & 0 deletions packages/vantio-cli/bin/ingest-url.cjs
Original file line number Diff line number Diff line change
@@ -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 };
4 changes: 3 additions & 1 deletion packages/vantio-cli/bin/interceptor.cjs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
2 changes: 1 addition & 1 deletion packages/vantio-cli/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -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": {
Expand Down
17 changes: 17 additions & 0 deletions packages/vantio-cli/test/ingest-url.test.js
Original file line number Diff line number Diff line change
@@ -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);
});
Loading