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
60 changes: 60 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 2 additions & 0 deletions context_intelligence/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -47,6 +48,7 @@

__all__ = [
"AsyncCIClient",
"build_event_payload",
"CIClient",
"resolve_config",
"LOG_SCHEMA",
Expand Down
76 changes: 74 additions & 2 deletions context_intelligence/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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,
Expand Down
41 changes: 41 additions & 0 deletions context_intelligence/events.py
Original file line number Diff line number Diff line change
@@ -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
117 changes: 117 additions & 0 deletions docs/lanes/public-event-ingestion/README.md
Original file line number Diff line number Diff line change
@@ -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 <verified-server-commit> \
--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.
55 changes: 55 additions & 0 deletions docs/lanes/public-event-ingestion/evidence/live.json
Original file line number Diff line number Diff line change
@@ -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
}
Loading
Loading