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
7 changes: 6 additions & 1 deletion docs/reference/public-api.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -318,7 +318,12 @@ These modules are public once their plugin package is installed:
- `Proxy`, `on`, `Request`, `ClientResponse`, `Client`, `Handler`,
`TunnelHandle`, `AbridgeError`, `DynamicRoutes`
- `Forward`, `SessionForward` (forward to a host-side sidecar / gateway)
- `Recorder` (record tunnel request/response pairs to JSONL)
- `Recorder` (record tunnel traffic as `abridge.record.v1` JSONL)
- `CaptureLevel`, `RequestFacts`, `ResponseFacts`, `PrefixRelation`,
`PrefixTracker`, `request_facts(...)`, `response_facts(...)`,
`RECORD_SCHEMA_VERSION`, `SESSION_META_SCHEMA_VERSION`
(the record's capture ladder and derivations — see
`agentix/bridge/capture.py`)
- `Sidecar`, `SidecarError`, `Command`
- handler clients under `agentix.bridge.clients` (`AnthropicClient`,
`OpenAIClient`, `AnthropicFromOpenAIClient`, `AnthropicToOpenAI`)
Expand Down
57 changes: 50 additions & 7 deletions plugins/abridge/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -131,22 +131,60 @@ Claude-speaking agent in front of an OpenAI-shaped recording gateway.

### Record the tunnel traffic

`Recorder` wraps any handler client and appends one JSONL line per served
call — `{ts, path, request_id, session_id?, request, response}` — flushed
as it goes, so the file is complete up to the last call even if the host
dies mid-rollout. It exposes the wrapped client's routes and closes it on
teardown, so it drops in transparently:
`Recorder` wraps any handler client and appends one `abridge.record.v1` JSONL
line per served call, flushed as it goes, so the file is complete up to the
last call even if the host dies mid-rollout. It exposes the wrapped client's
routes and closes it on teardown, so it drops in transparently:

```python
from agentix.bridge import Proxy, Recorder
from agentix.bridge import CaptureLevel, Proxy, Recorder

proxy = Proxy(Recorder(client, "runs/rollout-42.jsonl", session_id="rollout-42"))
proxy = Proxy(Recorder(
client, "runs/rollout-42.jsonl",
session_id="rollout-42",
level=CaptureLevel.VERBATIM, # default is METADATA
))
```

How much lands on disk is a level, and full capture is **opt-in**:

| level | on disk |
|---|---|
| `off` | nothing — don't install a Recorder |
| `metadata` (default) | identity, model, sampling, usage, tool **names**, per-message digests, the response block skeleton, and the prefix relation. No conversation text. |
| `verbatim` | all of the above **plus** the complete request body (`system`, `tools` with full schemas, the entire message history) and the complete structured response (`thinking` with its opaque `signature`, `text`, `tool_use`). |

`agentix/bridge/capture.py`'s module docstring is the normative schema
document. The two fields worth knowing about up front:

- **`prefix`** — whether this request's message list *extends* the previous
one. `"stable": false` means **the context was rewritten** (compaction, a
retry rollback, a fresh conversation on the same key), with
`divergence_index` pointing at the first message that differs. A consumer
reads that one boolean instead of parsing a harness's own compaction
boundary records. It is the Anthropic-face analogue of the token gateway's
`prefix_stable`, computed over per-message digests because this layer has
no tokenizer. Interleaved conversations (subagents, helper calls) are kept
in separate lanes keyed by the system prompt, so alternating between them
is not mistaken for a rewrite; a lane collision reports `false`, never a
false `true`.
- **`shape.content_blocks`** — one entry per block the model emitted, so
`{"type": "thinking", "chars": 0, "signature_chars": 210}` is visible even
at the metadata level. Streaming responses are recorded as the structured
`Message` (the bundled clients publish it), not as an SSE blob to re-parse.

The `request_id` in each row is the same id the transport stamps as
`x-request-id` on the upstream hop (bound through a context var), so a
message-level row joins a downstream token recorder's per-turn record;
`session_id`, when given, tags every row with the rollout identity.
`turn_index` is monotonic per file and advances even when a row fails to
serialize, so a dropped row leaves a detectable gap; `aclose()` appends an
`abridge.session.v1` trailer, so a truncated file is distinguishable from a
closed one.

Records never contain credentials — the tunnel carries no HTTP metadata at
all (`Request` is a path plus the decoded JSON body), so no `Authorization`
header or API key can reach one. Files are created `0600`.

## Writing your own handler

Expand Down Expand Up @@ -293,6 +331,9 @@ More serve options:
own session id, i.e. the `session_id` in its token records. The same
`request_id` reaches the upstream as `x-request-id`, so message rows
join the gateway's token records per call as well as per session.
* `--capture-level {off,metadata,verbatim}` (env `ABRIDGE_CAPTURE_LEVEL`,
default `metadata`) — how much of each call `--record-dir` persists.
Turning recording on never turns verbatim capture on; ask for it.

`GET /_health` reports `translation_spec_sha` — one SHA-256 over the
source of the Anthropic↔OpenAI transform module and both client modules
Expand All @@ -315,6 +356,8 @@ agentix/bridge/
├── proxy.py # Proxy + @on + sandbox tunnel + wire types
├── serve.py # direct mode: @on handlers as a standalone HTTP service
├── forward.py # JSON POST forwarding to a host-side service
├── capture.py # abridge.record.v1: capture levels, digests, prefix relation
├── recorder.py # Recorder: the JSONL sink for that record
├── sidecar.py # local process lifecycle + health supervision
└── clients/ # bundled handler implementations
├── openai.py # OpenAIClient (openai SDK) + PLACEHOLDER_API_KEY
Expand Down
12 changes: 7 additions & 5 deletions plugins/abridge/ROADMAP.md
Original file line number Diff line number Diff line change
Expand Up @@ -71,11 +71,13 @@ clients remain valid while sidecar gateways mature.
`@on(path)` by index. Useful for offline eval reruns, RL buffer
regression tests, CI-friendly assertions without burning tokens.

4. **Capture API.** Today storage is gone from the core (skipped in
the recent cleanup). Add `agentix.bridge.capture` — a small hook
that any handler can call (or a Proxy-level event subscriber) to
record full `(request, response)` pairs. Lightweight; in-memory
list with optional `JsonlSink` / `ParquetSink` overlays.
4. **Capture API.** Shipped as `agentix.bridge.capture` (the
`abridge.record.v1` shape, capture levels, and the prefix relation)
plus the `Recorder` JSONL sink. Still open: an in-memory sink for
tests, a columnar (`ParquetSink`) overlay for large corpora, and
structure recovery for an opaque SSE body relayed by a bare
`Forward` — today such a row keeps the bytes and says so rather
than reconstructing blocks the tunnel never decoded.

## Medium-term — additional bundled clients

Expand Down
22 changes: 21 additions & 1 deletion plugins/abridge/agentix/bridge/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,17 @@

from __future__ import annotations

from .capture import (
RECORD_SCHEMA_VERSION,
SESSION_META_SCHEMA_VERSION,
CaptureLevel,
PrefixRelation,
PrefixTracker,
RequestFacts,
ResponseFacts,
request_facts,
response_facts,
)
from .forward import Forward, SessionForward
from .proxy import (
NAMESPACE,
Expand All @@ -52,21 +63,30 @@
__version__ = "0.5.0"

__all__ = [
"NAMESPACE",
"RECORD_SCHEMA_VERSION",
"SESSION_META_SCHEMA_VERSION",
"AbridgeError",
"CaptureLevel",
"Client",
"ClientResponse",
"Command",
"DynamicRoutes",
"Forward",
"Handler",
"NAMESPACE",
"PrefixRelation",
"PrefixTracker",
"Proxy",
"Recorder",
"Request",
"RequestFacts",
"ResponseFacts",
"SessionForward",
"Sidecar",
"SidecarError",
"TunnelHandle",
"__version__",
"on",
"request_facts",
"response_facts",
]
26 changes: 26 additions & 0 deletions plugins/abridge/agentix/bridge/_request_id.py
Original file line number Diff line number Diff line change
Expand Up @@ -22,19 +22,33 @@
after the lazy session create). The `Recorder` clears it before each handler
call and reads it afterwards into the row's `gateway_session_id`, restoring
the session-level join between caller-side rows and gateway-side records.

`current_response_message` flows the same way and carries the STRUCTURE the
wire loses. An agent that asked for streaming gets `text/event-stream` bytes
back, so a capture layer reading only `ClientResponse` would have to re-parse
an SSE blob to find the assistant's `thinking` / `text` / `tool_use` blocks —
exactly the kind of string-scraping this capture exists to eliminate. The
Anthropic-face clients already hold the completed `Message` dict at that
point (they drain the stream and re-render it), so they publish it here and
the `Recorder` writes the object, not the blob.
"""

from __future__ import annotations

import uuid
from contextvars import ContextVar
from typing import Any

current_request_id: ContextVar[str | None] = ContextVar("abridge_request_id", default=None)

current_upstream_session_id: ContextVar[str | None] = ContextVar(
"abridge_upstream_session_id", default=None
)

current_response_message: ContextVar[dict[str, Any] | None] = ContextVar(
"abridge_response_message", default=None
)


def mint_request_id() -> str:
return uuid.uuid4().hex
Expand All @@ -46,9 +60,21 @@ def get_or_mint_request_id() -> str:
return bound if bound else mint_request_id()


def publish_response_message(message: dict[str, Any]) -> None:
"""Hand the completed provider-native response to an outer capture layer.

A no-op when nothing is capturing — setting a `ContextVar` nobody reads
costs one assignment, so clients call this unconditionally rather than
branching on whether a `Recorder` happens to be installed.
"""
current_response_message.set(message)


__all__ = [
"current_request_id",
"current_response_message",
"current_upstream_session_id",
"get_or_mint_request_id",
"mint_request_id",
"publish_response_message",
]
Loading
Loading