diff --git a/README.md b/README.md index 4fb2fab6..ecef2cd9 100644 --- a/README.md +++ b/README.md @@ -192,6 +192,66 @@ The actual `destinations` configuration — its sub-keys, routing rules, pattern ## Embedding in an Amplifier application +### Sending application events without mounting a hook + +Install `amplifier-bundle-context-intelligence[client]` to use the public async +client. Applications can build a single-event envelope and submit it through the +same server/auth contract used by the bundle: + +```python +import os +from datetime import datetime, timezone +from uuid import uuid4 + +from context_intelligence import AsyncCIClient +from context_intelligence import build_event_payload + +payload = build_event_payload( + event="application:diagnostic", + workspace="my-project", + data={ + "session_id": "application-session-123", # stable across this session + "timestamp": datetime.now(timezone.utc).isoformat(), + "event_id": str(uuid4()), # generated once per occurrence + "phase": "completed", + }, +) +# Persist this exact envelope in your outbox before dispatching, if needed. +client = AsyncCIClient(os.environ["CI_SERVER_URL"], os.environ["CI_API_KEY"]) +receipt = await client.ingest(payload) +``` + +`build_event_payload` is a pure, detached JSON transform. Its deterministic +`aci-event-v1` key matches the telemetry hook. The caller supplies event/session +identity, timestamps and sanitized data. The current server requires a timestamp; +session-scoped authorization and graph storage also require `data.session_id`. +Include a stable occurrence ID in `data` when two otherwise identical events are +distinct. An optional `working_dir` is metadata outside the v1 key; omit it when +unknown. Neither the helper nor client reads application state or chooses a +destination, privacy policy, local spool, fan-out or retry schedule. + +`ingest` makes one request, using the existing static-key or `auth_strategy` +credential resolution. It does not follow redirects or retry. Failures use +`CIClientError`, including `invalid_payload` for a non-object or non-JSON caller +envelope (before auth/network), and HTTP status/retry-after metadata when supplied. A +timeout or connection loss can have an unknown acceptance outcome: an explicit +retry should reuse the persisted envelope, including its original key. Custom +server-compatible envelopes and idempotency keys are also accepted. + +The returned receipt retains the server's exact status: HTTP 202 `queued` reports +durable queue acceptance, while `duplicate` reports recognition of an existing +key, not a new append. At the currently tested server revision, a duplicate alone +cannot prove durable acceptance after an earlier append failure. Neither proves +graph indexing has finished. Deduplication +scope, retention and crash behavior belong to the server; this API does not +promise exactly-once delivery. See the [live validation and server +limitations](docs/lanes/public-event-ingestion/README.md). + +Synchronous applications can call `asyncio.run(client.ingest(payload))` when they +do not already have a running event loop. There is no synchronous ingestion API. + +### Mounting the telemetry hook + When integrating this hook from Python rather than through the bundle CLI, call `mount()` directly. ```python diff --git a/context_intelligence/__init__.py b/context_intelligence/__init__.py index a5008b8c..a2c3e35b 100644 --- a/context_intelligence/__init__.py +++ b/context_intelligence/__init__.py @@ -18,6 +18,7 @@ from __future__ import annotations from context_intelligence.client import AsyncCIClient, CIClient +from context_intelligence.events import build_event_payload from context_intelligence.config import ( AMPLIFIER_DIR, LOG_SCHEMA, @@ -47,6 +48,7 @@ __all__ = [ "AsyncCIClient", + "build_event_payload", "CIClient", "resolve_config", "LOG_SCHEMA", diff --git a/context_intelligence/client.py b/context_intelligence/client.py index 78c419bc..413e8274 100644 --- a/context_intelligence/client.py +++ b/context_intelligence/client.py @@ -71,7 +71,7 @@ def __init__( ) -> None: super().__init__(message) #: One of "connection_error" | "timeout" | "http_status" | "decode_error" - #: | "auth_error". "auth_error" -- an unusable credential -- is raised + #: | "auth_error" | "invalid_payload". "auth_error" -- an unusable credential -- is raised #: BEFORE the request is attempted, so it is never confused with #: "decode_error" (a bad response body from a server actually reached). self.error_type = error_type @@ -862,7 +862,7 @@ class AsyncCIClient: When ``None``, an ``ApiKeyAuth(api_key)`` is built implicitly (backward compat). timeout: Per-request HTTP timeout (seconds) applied to every ``httpx.AsyncClient`` - constructed by this instance (cypher, fetch_blob, list_blob_keys). Defaults + constructed by this instance (including ingest). Defaults to 30.0, matching the sync helpers' existing ``timeout=30``. See ``ToolConfigResolver.request_timeout`` for how callers resolve this value from config/env. @@ -917,6 +917,78 @@ def _auth_headers(self, url: str) -> dict[str, str]: # Public API # ------------------------------------------------------------------ + async def ingest(self, payload: dict[str, Any]) -> dict[str, Any]: + """Submit an event envelope once, returning the server's acceptance receipt. + + ``build_event_payload`` supplies an optional deterministic v1 envelope; + callers may also supply their own server-compatible envelope and key. + The envelope is detached before resolving auth so concurrent caller edits + cannot change the request while a credential is refreshed. + + HTTP 202 / ``queued`` means durable server acceptance, not graph indexing. + HTTP 202 / ``duplicate`` means the server recognized the idempotency key; + it does not append another event or establish graph completion. The server + owns deduplication scope and retention. This method never retries or follows + redirects. Persist and reuse the same envelope for an explicit retry after + an ambiguous failure. Routing, redaction and retry policy belong to callers. + """ + import asyncio + + url = f"{self._server_url}/events" + # A caller's invalid envelope is distinct from an invalid server receipt. + # Snapshot before auth without allowing NaN or Infinity on the wire. + try: + if not isinstance(payload, dict): + raise TypeError("event payload must be an object") + body = json.loads(json.dumps(payload, allow_nan=False)) + except (TypeError, ValueError) as exc: + raise CIClientError( + "event payload is not a JSON-encodable object", + error_type="invalid_payload", + url=url, + ) from exc + headers = await asyncio.to_thread(self._auth_headers, url) + try: + async with httpx.AsyncClient(timeout=self._timeout, follow_redirects=False) as client: # type: ignore[union-attr] + resp = await client.post(url, json=body, headers=headers) + resp.raise_for_status() + result = resp.json() + if ( + resp.status_code != 202 + or not isinstance(result, dict) + or not isinstance(result.get("status"), str) + or result["status"] not in {"queued", "duplicate"} + or ( + result.get("session_id") is not None + and not isinstance(result["session_id"], str) + ) + ): + raise ValueError("invalid event acceptance receipt") + return result + except httpx.TimeoutException as exc: # type: ignore[union-attr] + raise CIClientError( + "event acceptance timed out", error_type="timeout", url=url + ) from exc + except httpx.HTTPStatusError as exc: # type: ignore[union-attr] + body, code = _error_detail(exc.response) + raise CIClientError( + "event acceptance rejected", + error_type="http_status", + url=url, + status_code=exc.response.status_code, + retry_after=_retry_after_seconds(exc.response.headers), + error_body=body, + error_code=code, + ) from exc + except ValueError as exc: + raise CIClientError( + "invalid event receipt", error_type="decode_error", url=url + ) from exc + except httpx.HTTPError as exc: # type: ignore[union-attr] + raise CIClientError( + "event connection failed", error_type="connection_error", url=url + ) from exc + async def cypher( self, query: str, diff --git a/context_intelligence/events.py b/context_intelligence/events.py new file mode 100644 index 00000000..1096d748 --- /dev/null +++ b/context_intelligence/events.py @@ -0,0 +1,41 @@ +"""Level 1: pure single-event envelopes compatible with the ingestion protocol.""" + +from __future__ import annotations + +import hashlib +import json +from typing import Any + + +def build_event_payload( + event: str, workspace: str, data: dict[str, Any], working_dir: str | None = None +) -> dict[str, Any]: + """Build a detached, JSON-safe envelope with the hook's v1 idempotency key. + + Callers must sanitize content before calling this function. Include a stable + event identity in ``data`` when otherwise-identical occurrences are distinct. + Persist the returned envelope unchanged for retries. ``working_dir`` is + envelope metadata and intentionally excluded from the existing v1 key. + An unknown working directory is omitted, never encoded as a blank string. + """ + if not isinstance(event, str) or not event.strip(): + raise ValueError("event must be a nonblank string") + if not isinstance(workspace, str) or not workspace.strip(): + raise ValueError("workspace must be a nonblank string") + if not isinstance(data, dict): + raise ValueError("data must be an object") + if working_dir is not None and (not isinstance(working_dir, str) or not working_dir.strip()): + raise ValueError("working_dir must be nonblank when supplied") + canonical = json.dumps( + {"event": event, "workspace": workspace, "data": data}, + sort_keys=True, + separators=(",", ":"), + allow_nan=False, + ) + payload = json.loads(canonical) + payload["idempotency_key"] = ( + "aci-event-v1:" + hashlib.sha256(canonical.encode("utf-8")).hexdigest() + ) + if working_dir is not None: + payload["working_dir"] = working_dir + return payload diff --git a/docs/lanes/public-event-ingestion/README.md b/docs/lanes/public-event-ingestion/README.md new file mode 100644 index 00000000..0be9fecf --- /dev/null +++ b/docs/lanes/public-event-ingestion/README.md @@ -0,0 +1,117 @@ +# Public event ingestion: contract and live evidence + +Applications need to send their own event streams without mounting a session +telemetry hook or importing a module's private uploader. This additive change +provides a pure `build_event_payload` helper and `AsyncCIClient.ingest` using the +existing server/auth contract. Callers own routing, consent, redaction, local +retention, fan-out and retry scheduling. No bundle, agent, skill or mode changes. + +## Acceptance contract + +- The helper produces a detached JSON envelope and the hook-compatible + `aci-event-v1` content hash. Distinct occurrences need distinct stable identity + in `data`. JSON containing NaN/Infinity is rejected. Unknown `working_dir` is + omitted; a supplied directory is outside the v1 hash, as in the hook. +- The client snapshots the supplied envelope before refreshing auth, posts once, + does not follow redirects, and returns the server's HTTP 202 receipt. Custom + compatible envelopes and keys are allowed. Invalid non-object/non-JSON payloads + raise `CIClientError(error_type="invalid_payload")` before authentication or HTTP. Existing auth strategies and + `CIClientError` classifications are reused; blocking token refresh runs outside + the event loop. +- `queued` means the server accepted a durable queue append. `duplicate` means + its deduplication cache recognized the key. Neither means graph indexing is + complete. There is no automatic retry or exactly-once guarantee. +- Current server handlers require `data.timestamp`. Session-scoped authorization + and graph ingestion also require a nonblank `data.session_id`; generic events + without a session may be accepted by compatibility mode but not indexed. The + helper deliberately does not invent either field. + +## Real seam run + +[Captured evidence](evidence/live.json) scores five scenarios through the actual +public API, real HTTP/auth and Neo4j. All five passed: first `queued` receipt, +second `duplicate` receipt, exactly one matching graph event, explicit HTTP 401 +for a bad key, and `connection_error` for a stopped endpoint. All event data is +synthetic. There are no credentials, prompts or personal session contents in the +evidence. The tests' receipt doubles were reconciled with these responses. + +Provenance: + +- Library starting commit: `87894048f8b68519f61aea52306d907b04f4bdfb` plus this + change; exact client/helper source SHA-256 values are in the evidence. +- Server: `microsoft/amplifier-context-intelligence` commit + `43973967ff4bc9a02422814a8c00dfce7624a76d`, clean checkout, single authenticated `create_asgi_app` uvicorn factory + process at `http://127.0.0.1:18081`, explicit isolated storage/config. +- Neo4j: official `neo4j:5.26.22-community` image with bundled APOC, + image ID `sha256:24b071534c7cfe9718689041ab9aafea2cd0d88af9ea58b768cb5eed381ab2d0`. + A dedicated container exposed Bolt on 17687 and HTTP on 17474. +- Auth: a dedicated synthetic static identity. Entra refresh is covered by + strategy tests; this run did not authenticate against Azure. + +Reproduce against an explicitly provisioned isolated server, with write/query +access to the `ci-public-api-fixture` workspace: + +```bash +uv run --frozen python scripts/verify-event-ingestion.py \ + --server-url http://127.0.0.1:18081 \ + --token-file /private/path/to/token \ + --server-revision \ + --output /private/path/to/evidence.json +``` + +The harness uses a unique session/event, queries only that session, and prints +pass/fail scores without credentials or arbitrary server error text. The default +network-failure endpoint is `127.0.0.1:1`; override `--unavailable-url` if needed. +It does not provision services, stop them, read user configuration or run models. + +## Server limitation discovered at this seam + +At the server revision above, `context_intelligence_server/main.py:1501–1510` +calls `idempotency_cache.check_and_store` before the durable queue append at +line 1533. If that first append fails, a later request with the same key can +receive `duplicate` even though that attempt never appended. A duplicate receipt +therefore cannot independently establish durability. The client exposes the +receipt without promoting it to `queued` or silently replaying with a new key. +This is a source-level failure-window finding; the live run did not inject a disk +append failure or test server crash/restart durability. It requires a separate +server-side correction and fault-injection test. + +## Validation + +- `uv run --frozen pytest tests -q`: 891 passed. +- Query module, `PYTHONPATH=../.. uv run --frozen pytest -q`: 195 passed. +- `uv run --frozen ruff check .`, `ruff format --check .`, `pyright`: clean. +- `uv build --no-sources`: sdist and wheel built. The wheel's `[client]` extra + installed in a fresh environment; an isolated `python -I` imported both APIs. +- Five real seam scenarios passed, as captured above. +- `scripts/validate-full.sh`: exit 0, `validation_mode: full`, adjudicated + **PASS WITH SUGGESTIONS**. Its mechanical FAIL is the known mode-advertising + false positive; the review also confirmed the standalone README install + convention is intentional. The existing `server-data-ops` agent description + exceeds the cosmetic length recommendation. Those unrelated files are + unchanged. The recipe lacks `pip wheel`; the separate successful `uv build` + and installed-wheel smoke cover package construction. The validator's + generated overview diagram changes are excluded because this change does not + alter bundle structure. + +Test environment notes: bare root `pytest` collects independently packaged module +suites and currently fails collection due to missing module dependencies and +colliding `tests.conftest` names. The root's intended `tests/` suite was selected +explicitly. The query module's lock pins an older library revision without the +existing `CIClientError` class, so `PYTHONPATH=../..` selects this checkout for the +required module suite. Neither unrelated setup issue was changed here. Frozen +uv commands preserve the repository's public-index lock in environments with +an externally configured default package index. + +The current recipe runner injects the CLI interpreter as `AMPLIFIER_PYTHON`, so +the helper's PATH-only dependency environment did not supply hatchling on the +first run (`full_no_build`). For the successful rerun, hatchling was installed +into `/tmp/amplifier-ci-validator-extra` with `uv pip install --target`, and +`PYTHONPATH` pointed there for this command only. No installed CLI environment +was changed. Recipe version: 3.15.0; recipe SHA-256: +`dd1f7b7e87cd384bc8f49d9545cd641a52287f2b9f7860985e99493a9cd6cb5d`. + +Review revision: the Level 1 builder lives in `context_intelligence.events` and +is re-exported from the package root; the reserved uploader namespace is unchanged. +The additive public API increments the package version to 0.2.0. Hash parity +covers Unicode, nested values, nulls, and arrays in addition to the original fixture. diff --git a/docs/lanes/public-event-ingestion/evidence/live.json b/docs/lanes/public-event-ingestion/evidence/live.json new file mode 100644 index 00000000..449c7059 --- /dev/null +++ b/docs/lanes/public-event-ingestion/evidence/live.json @@ -0,0 +1,55 @@ +{ + "executed_at": "2026-09-19T03:08:44.332768+00:00", + "server_url": "http://127.0.0.1:18081", + "server_revision": "43973967ff4bc9a02422814a8c00dfce7624a76d", + "client_source_sha256": { + "context_intelligence/client.py": "62de48c0a34462a378eda4da82a9816db755959efbc67122650963e074a10ab8", + "context_intelligence/events.py": "e513d77839bec9d7e0bface5737c654ce6fabcd39449819d4ca913bf3fc77df4" + }, + "request": { + "data": { + "event_id": "01abd059d8754ce6be13dc590507de8f", + "message": "synthetic public ingestion seam", + "session_id": "public-ingest-d3cf96519e8b43dc8f37711f5f787cf2", + "timestamp": "2026-09-19T03:08:43.905753+00:00" + }, + "event": "host:diagnostic", + "workspace": "ci-public-api-fixture", + "idempotency_key": "aci-event-v1:b65830071eb14d6719ea79d01137ea2e53eb63001f48898fb60e2eae7948a9b3" + }, + "receipts": [ + { + "status": "queued", + "session_id": "public-ingest-d3cf96519e8b43dc8f37711f5f787cf2" + }, + { + "status": "duplicate", + "session_id": "public-ingest-d3cf96519e8b43dc8f37711f5f787cf2" + } + ], + "graph_query": "MATCH (n:Event) WHERE n.session_id = $sid RETURN n.event_name AS event_name, count(n) AS count", + "graph_rows": [ + { + "event_name": "host:diagnostic", + "count": 1 + } + ], + "failures": { + "bad_auth": { + "error_type": "http_status", + "status_code": 401 + }, + "connection_refused": { + "error_type": "connection_error", + "status_code": null + } + }, + "checks": { + "first_acceptance": true, + "duplicate_receipt": true, + "one_graph_event": true, + "auth_fails_loudly": true, + "network_fails_loudly": true + }, + "passed": true +} diff --git a/pyproject.toml b/pyproject.toml index 192133b2..6b574b04 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -1,6 +1,6 @@ [project] name = "amplifier-bundle-context-intelligence" -version = "0.1.3" +version = "0.2.0" requires-python = ">=3.11" dependencies = [ # context_intelligence/auth.py imports azure.identity.DefaultAzureCredential @@ -12,6 +12,9 @@ dependencies = [ "azure-identity>=1.19", ] +[project.optional-dependencies] +client = ["httpx>=0.25"] + [build-system] requires = [ "hatchling", diff --git a/scripts/verify-event-ingestion.py b/scripts/verify-event-ingestion.py new file mode 100644 index 00000000..77154b0a --- /dev/null +++ b/scripts/verify-event-ingestion.py @@ -0,0 +1,130 @@ +#!/usr/bin/env python3 +"""Score the public ingestion seam against an explicitly supplied, isolated server. + +Writes one synthetic event and reads only its unique session. Does not provision a +server, inspect user configuration, or print credentials/response error bodies. +""" + +from __future__ import annotations + +import argparse +import asyncio +import hashlib +import json +from datetime import datetime, timezone +from pathlib import Path +from urllib.parse import urlsplit +from uuid import uuid4 + +from context_intelligence import AsyncCIClient +from context_intelligence.client import CIClientError +from context_intelligence import build_event_payload + + +def server_url(value: str) -> str: + parsed = urlsplit(value) + if ( + parsed.scheme not in {"http", "https"} + or not parsed.hostname + or parsed.username + or parsed.password + or parsed.query + or parsed.fragment + ): + raise argparse.ArgumentTypeError( + "use an HTTP(S) URL without credentials, query or fragment" + ) + return value.rstrip("/") + + +async def score(args: argparse.Namespace) -> dict: + token = args.token_file.read_text().strip() + if not token: + raise ValueError("token file is empty") + sid = "public-ingest-" + uuid4().hex + payload = build_event_payload( + "host:diagnostic", + "ci-public-api-fixture", + { + "session_id": sid, + "timestamp": datetime.now(timezone.utc).isoformat(), + "event_id": uuid4().hex, + "message": "synthetic public ingestion seam", + }, + ) + client = AsyncCIClient(args.server_url, token, timeout=10) + receipts = [await client.ingest(payload), await client.ingest(payload)] + query = ( + "MATCH (n:Event) WHERE n.session_id = $sid " + "RETURN n.event_name AS event_name, count(n) AS count" + ) + rows = [] + for _ in range(80): + rows = await client.cypher(query, workspace=payload["workspace"], params={"sid": sid}) + if rows: + break + await asyncio.sleep(0.25) + failures = {} + for name, failing in ( + ("bad_auth", AsyncCIClient(args.server_url, "synthetic-invalid-token", timeout=2)), + ("connection_refused", AsyncCIClient(args.unavailable_url, token, timeout=2)), + ): + try: + await failing.ingest(payload) + except CIClientError as exc: + failures[name] = {"error_type": exc.error_type, "status_code": exc.status_code} + else: + failures[name] = {"error_type": "unexpected_success"} + source = Path(__file__).resolve().parents[1] + checks = { + "first_acceptance": receipts[0] == {"status": "queued", "session_id": sid}, + "duplicate_receipt": receipts[1] == {"status": "duplicate", "session_id": sid}, + "one_graph_event": rows == [{"event_name": payload["event"], "count": 1}], + "auth_fails_loudly": failures["bad_auth"] + == {"error_type": "http_status", "status_code": 401}, + "network_fails_loudly": failures["connection_refused"] + == {"error_type": "connection_error", "status_code": None}, + } + return { + "executed_at": datetime.now(timezone.utc).isoformat(), + "server_url": args.server_url, + "server_revision": args.server_revision, + "client_source_sha256": { + str(path.relative_to(source)): hashlib.sha256(path.read_bytes()).hexdigest() + for path in ( + source / "context_intelligence/client.py", + source / "context_intelligence/events.py", + ) + }, + "request": payload, + "receipts": receipts, + "graph_query": query, + "graph_rows": rows, + "failures": failures, + "checks": checks, + "passed": all(checks.values()), + } + + +def main() -> int: + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--server-url", type=server_url, required=True) + parser.add_argument("--token-file", type=Path, required=True) + parser.add_argument("--server-revision", required=True, help="operator-verified server commit") + parser.add_argument("--unavailable-url", type=server_url, default="http://127.0.0.1:1") + parser.add_argument("--output", type=Path, required=True) + args = parser.parse_args() + try: + result = asyncio.run(score(args)) + except CIClientError as exc: + print(json.dumps({"passed": False, "error_type": exc.error_type})) + return 1 + args.output.parent.mkdir(parents=True, exist_ok=True) + args.output.write_text(json.dumps(result, indent=2) + "\n") + args.output.chmod(0o600) + print(json.dumps({"passed": result["passed"], "checks": result["checks"]})) + return 0 if result["passed"] else 1 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/tests/test_event_ingestion.py b/tests/test_event_ingestion.py new file mode 100644 index 00000000..7698bdc6 --- /dev/null +++ b/tests/test_event_ingestion.py @@ -0,0 +1,258 @@ +"""Ingestion contracts reconciled with the live-server evidence in docs/lanes/public-event-ingestion.""" + +from __future__ import annotations + +import asyncio +from datetime import datetime, timezone +import importlib.util +import json +from pathlib import Path +import threading +from unittest.mock import patch + +import httpx +import pytest + +from context_intelligence.client import AsyncCIClient, CIClientError +from context_intelligence import build_event_payload + + +def envelope(): + return build_event_payload( + "host:diagnostic", + "fixture-workspace", + { + "session_id": "fixture-session", + "timestamp": "2026-09-18T22:00:00+00:00", + "text": "Résumé", + }, + ) + + +def test_payload_is_detached_canonical_and_compatible_with_hook_v1(): + path = ( + Path(__file__).parents[1] + / "modules/hook-context-intelligence/amplifier_module_hook_context_intelligence/upload.py" + ) + spec = importlib.util.spec_from_file_location("hook_upload", path) + assert spec is not None and spec.loader is not None + module = importlib.util.module_from_spec(spec) + spec.loader.exec_module(module) + data = {"nested": {"value": "Résumé"}, "session_id": "fixture"} + payload = build_event_payload("host:diagnostic", "fixture", data) + assert ( + payload["idempotency_key"] + == module.build_payload("host:diagnostic", "fixture", data)["idempotency_key"] + ) + assert ( + build_event_payload( + "host:diagnostic", "fixture", dict(reversed(list(data.items()))), "/other" + )["idempotency_key"] + == payload["idempotency_key"] + ) + assert "working_dir" not in payload + data["nested"]["value"] = "changed" # type: ignore[index] + assert payload["data"]["nested"]["value"] == "Résumé" + assert json.loads(json.dumps(payload)) == payload + + +@pytest.mark.parametrize( + "event,workspace,data,directory", + [ + ("", "w", {}, None), + ("e", " ", {}, None), + ("e", "w", [], None), + ("e", "w", {}, " "), + ("e", "w", {"value": float("nan")}, None), + ], +) +def test_payload_rejects_unusable_json_or_blank_required_fields(event, workspace, data, directory): + with pytest.raises(ValueError): + build_event_payload(event, workspace, data, directory) + + +def transport(handler): + original = httpx.AsyncClient + return patch( + "context_intelligence.client.httpx.AsyncClient", + side_effect=lambda **kw: original(transport=httpx.MockTransport(handler), **kw), + ) + + +@pytest.mark.parametrize( + "payload", + [ + [], + "not-an-envelope", + {"data": {"unserializable": {1, 2}}}, + {"data": {"value": float("nan")}}, + {"data": {"value": float("inf")}}, + {"data": {"timestamp": datetime(2026, 9, 18, tzinfo=timezone.utc)}}, + ], +) +async def test_unusable_caller_envelope_has_classified_error_before_auth_or_network(payload): + class Strategy: + def headers(self): + pytest.fail("Invalid payload must not request credentials") + + def handle(request): + pytest.fail("Invalid payload must not reach the network") + + with transport(handle), pytest.raises(CIClientError) as caught: + await AsyncCIClient("http://server.invalid", auth_strategy=Strategy()).ingest(payload) + assert caught.value.error_type == "invalid_payload" + assert caught.value.url == "http://server.invalid/events" + + +@pytest.mark.parametrize( + "data", + [ + {}, + {"value": None}, + {"value": [True, False, 1, 1.25]}, + {"nested": {"z": "☀ Résumé", "a": {"value": -5}}}, + { + "event_id": "distinct-occurrence", + "session_id": "root", + "timestamp": "2026-09-18T00:00:00Z", + }, + ], +) +def test_hook_v1_keys_match_across_supported_json_shapes(data): + path = ( + Path(__file__).parents[1] + / "modules/hook-context-intelligence/amplifier_module_hook_context_intelligence/upload.py" + ) + spec = importlib.util.spec_from_file_location("hook_upload_shapes", path) + assert spec is not None and spec.loader is not None + module = importlib.util.module_from_spec(spec) + spec.loader.exec_module(module) + assert ( + build_event_payload("host:diagnostic", "fixture", data)["idempotency_key"] + == module.build_payload("host:diagnostic", "fixture", data, "/hookdir")["idempotency_key"] + ) + + +@pytest.mark.parametrize("status", ["queued", "duplicate"]) +async def test_ingest_preserves_real_server_acceptance_receipts_and_custom_envelope(status): + captured = [] + payload = { + **envelope(), + "idempotency_key": "caller-owned-stable-id", + "extension": {"future": True}, + } + + def handle(request): + captured.append(request) + return httpx.Response(202, json={"status": status, "session_id": "fixture-session"}) + + with transport(handle): + result = await AsyncCIClient("http://server.invalid/", "synthetic-key").ingest(payload) + assert result == {"status": status, "session_id": "fixture-session"} + assert len(captured) == 1 + assert json.loads(captured[0].content) == payload + assert captured[0].headers["authorization"] == "Bearer synthetic-key" + + +@pytest.mark.parametrize( + "code,body", + [ + (200, {"status": "queued"}), + (202, {"status": "complete"}), + (202, {"status": []}), + (202, {"status": "queued", "session_id": 3}), + (202, []), + (202, "malformed"), + ], +) +async def test_unrecognized_receipt_is_not_success(code, body): + def handle(request): + return ( + httpx.Response(code, content=body) + if isinstance(body, str) + else httpx.Response(code, json=body) + ) + + with transport(handle), pytest.raises(CIClientError) as caught: + await AsyncCIClient("http://server.invalid", "synthetic-key").ingest(envelope()) + assert caught.value.error_type == "decode_error" + + +@pytest.mark.parametrize("status", [302, 401, 403, 429, 503]) +async def test_rejection_is_attempted_once_and_redirects_are_not_followed(status): + calls = [] + + def handle(request): + calls.append(request) + return httpx.Response( + status, + json={"detail": "fixture rejection"}, + headers={"Location": "https://elsewhere.invalid", "Retry-After": "7"}, + ) + + with transport(handle), pytest.raises(CIClientError) as caught: + await AsyncCIClient("http://server.invalid", "synthetic-key").ingest(envelope()) + assert len(calls) == 1 and caught.value.status_code == status + assert caught.value.error_type == "http_status" + assert caught.value.retry_after == 7 + + +@pytest.mark.parametrize( + "failure,error_type", [(httpx.ReadTimeout, "timeout"), (httpx.ConnectError, "connection_error")] +) +async def test_transport_failure_does_not_retry(failure, error_type): + calls = [] + + def handle(request): + calls.append(request) + raise failure("fixture failure", request=request) + + with transport(handle), pytest.raises(CIClientError) as caught: + await AsyncCIClient("http://server.invalid", "synthetic-key").ingest(envelope()) + assert caught.value.error_type == error_type and len(calls) == 1 + + +async def test_auth_refresh_does_not_block_loop_or_allow_mutating_inflight_payload(): + started = threading.Event() + release = threading.Event() + captured = [] + + class Strategy: + def headers(self): + started.set() + assert release.wait(3) + return {"Authorization": "Bearer refreshed-fixture"} + + def handle(request): + captured.append(request) + return httpx.Response(202, json={"status": "queued"}) + + payload = envelope() + with transport(handle): + task = asyncio.create_task( + AsyncCIClient("http://server.invalid", auth_strategy=Strategy()).ingest(payload) + ) + try: + async with asyncio.timeout(2): + while not started.is_set(): + await asyncio.sleep(0.001) + payload["data"]["text"] = "changed while resolving auth" + release.set() + await task + finally: + release.set() + assert json.loads(captured[0].content)["data"]["text"] == "Résumé" + assert captured[0].headers["authorization"] == "Bearer refreshed-fixture" + + +async def test_invalid_auth_fails_before_network(): + class Strategy: + def headers(self): + raise ValueError("fixture credential unavailable") + + def handle(request): + pytest.fail("No request for unusable auth") + + with transport(handle), pytest.raises(CIClientError) as caught: + await AsyncCIClient("http://server.invalid", auth_strategy=Strategy()).ingest(envelope()) + assert caught.value.error_type == "auth_error" diff --git a/uv.lock b/uv.lock index 61cbed71..693bd95c 100644 --- a/uv.lock +++ b/uv.lock @@ -4,12 +4,17 @@ requires-python = ">=3.11" [[package]] name = "amplifier-bundle-context-intelligence" -version = "0.1.3" +version = "0.2.0" source = { editable = "." } dependencies = [ { name = "azure-identity" }, ] +[package.optional-dependencies] +client = [ + { name = "httpx" }, +] + [package.dev-dependencies] dev = [ { name = "amplifier-core" }, @@ -23,7 +28,11 @@ dev = [ ] [package.metadata] -requires-dist = [{ name = "azure-identity", specifier = ">=1.19" }] +requires-dist = [ + { name = "azure-identity", specifier = ">=1.19" }, + { name = "httpx", marker = "extra == 'client'", specifier = ">=0.25" }, +] +provides-extras = ["client"] [package.metadata.requires-dev] dev = [