Metadata-Version: 2.4
Name: runlatent
Version: 0.1.0
Summary: Latent SDK: one call at startup posts every finished OpenAI / Anthropic / Google GenAI call to Latent's scoring API, fire-and-forget.
License: Proprietary
Project-URL: Homepage, https://runlatent.ai
Project-URL: Documentation, https://github.com/runlatentai/core/blob/main/runlatent/README.md
Project-URL: Source, https://github.com/runlatentai/core
Classifier: Programming Language :: Python :: 3
Classifier: Operating System :: OS Independent
Requires-Python: >=3.10
Description-Content-Type: text/markdown

# runlatent — ingest finished LLM calls into Latent's scoring API

```
pip install runlatent
```

A Python SDK that captures every finished call your app makes through the OpenAI, Anthropic
or Google GenAI Python SDK and posts the (prompt, answer) pair to Latent's scoring API, which
is the sidecar's `POST /score` contract (`latent_sidecar/README.md`). The app's call is never
slower by a network round trip, never sees a different return value and never gets an
exception from here. Step 1 of the self-serve path: no gateway, no log tail, one line at startup.

```
your app ──openai / anthropic / google-genai SDK──▶ the provider
   │ (after the SDK returned, fire-and-forget, a daemon thread)
   └──POST <LATENT_API_URL>/score {request_id, messages, output, model, finish_reason}──▶ Latent
```

## Install

`pip install runlatent` (the `runlatent` distribution on PyPI, built from `runlatent/` in this
repository; from the repository, `pip install -e .`). It has no runtime dependencies: nothing
from numpy, torch or the core packages reaches your process. The provider SDKs are yours:
whichever of `openai`, `anthropic`, `google-genai` is importable gets patched, the others are
skipped.

## One-line setup

```python
import runlatent
runlatent.auto_instrument()        # reads LATENT_API_URL and LATENT_API_KEY
```

Call it once, before the clients are built or after, it does not matter: the patch is on the
SDK classes (`Completions.create`, `AsyncCompletions.create`, `Responses.create`,
`AsyncResponses.create`; `Messages.create` / `.stream` and the async pair; `Models.generate_content`
/ `.generate_content_stream` and `AsyncModels`), so every client in the process is covered,
including ones a framework builds for you. Keyword arguments:

| argument | default | meaning |
|---|---|---|
| `openai`, `anthropic`, `google_genai` | `True` | which SDKs to patch |
| `api_url` | `$LATENT_API_URL` | the API base; the SDK posts to `<api_url>/score` |
| `api_key` | `$LATENT_API_KEY` | sent as `Authorization: Bearer <key>` |
| `redact` | `None` | `callable(body) -> body or None`, run on every pair before it is queued (below) |
| `timeout` | `5.0` | seconds per POST |
| `max_queue` | `10000` | pairs waiting for the poster thread; a full queue drops, counted |

Idempotent: a second call patches nothing twice, and with the same destination it keeps the
same poster and its queue; a call with a different `api_url` / `api_key` / `timeout` /
`max_queue` swaps the poster at once and drains the old one on its own thread (its final
counters stay in `stats()`), never on the caller and never while the app's calls wait.
`runlatent.uninstrument()` restores every original (classes and wrapped instances). With no URL anywhere, `auto_instrument()` logs one
warning, patches nothing and returns `stats()` with `enabled: False`.

## Per-client wrappers

The same capture on one client object, nothing else in the process touched:

```python
client = runlatent.wrap_openai(OpenAI())          # or AsyncOpenAI()
client = runlatent.wrap_anthropic(Anthropic())    # or AsyncAnthropic()
client = runlatent.wrap_genai(genai.Client(...))  # sync and .aio
```

Each returns the client. Wrapping twice is a no-op; a wrapped instance under a patched class
captures once, not twice.

## Environment

| variable | meaning |
|---|---|
| `LATENT_API_URL` | the API base (the sidecar's address when you run it yourself: `http://sidecar:8090`) |
| `LATENT_API_KEY` | the bearer (the sidecar's `LATENT_SIDECAR_TOKEN`) |
| `LATENT_SIDECAR_URL`, `LATENT_SIDECAR_TOKEN` | accepted as fallbacks when the two above are unset |
| `LATENT_SDK_EXIT_FLUSH_S` | seconds the atexit handler waits for queued pairs at a normal interpreter exit (default `2.0`; `0` disables, see Limits) |

## What is posted

One JSON body per finished call, the sidecar's `/score` contract plus one extra block:

```json
{"request_id": "chatcmpl-abc123",
 "messages": [{"role": "system", "content": "..."}, {"role": "user", "content": "..."}],
 "output": "the full answer text",
 "model": "gpt-4.1", "finish_reason": "stop",
 "client": {"provider": "openai", "api": "chat.completions.create", "sdk": "openai", "sdk_version": "1.x",
            "runlatent": "0.1.0", "captured_at": "2026-10-02T12:00:00.000+00:00", "stream": false,
            "response_id": "chatcmpl-abc123", "finish_reason_raw": "stop",
            "usage": {"input_tokens": 120, "output_tokens": 48, "total_tokens": 168},
            "parts_dropped": 0, "items_dropped": 0}}
```

- `request_id` is the provider's response id (`chatcmpl-…`, `msg_…`, `resp_…`, Gemini's
  `response_id`) when it has one, else a uuid; the raw id is repeated in `client.response_id`.
- `messages` is the request as the app sent it, system prompt included, reduced to
  `{role, content}`: OpenAI `messages` (a `tool` turn is a text message; an assistant turn with
  `tool_calls` carries `[tool_calls omitted]`); the Responses API's `instructions` as `system`
  and its `input` string or message items (non-message items, function call outputs and the like,
  are skipped and counted in `items_dropped`; `previous_response_id` is recorded, the earlier
  turns are not fetched); Anthropic's `system` as `system` and `messages` (`tool_result` text
  kept, `thinking` blocks skipped); GenAI's `config.system_instruction` as `system` and `contents`
  (`model` role becomes `assistant`, loose strings and parts become one user turn).
- `output` is the whole answer text: the first choice's `message.content`; the Responses
  output's `message` items; Anthropic's `text` blocks; the first Gemini candidate's text parts
  (thoughts skipped). Non-text answer parts (a `tool_use` block, a Gemini `function_call`, an
  image) contribute nothing to it, so an answer with no text is not posted: counted
  `skipped_tool_only` when such parts were seen, `skipped_empty_output` otherwise (a blocked
  prompt). A Responses API response with `status: failed` or an `error` is not an answer and is
  never posted (`skipped_failed_response`), in the stream (`response.failed`) or not.
- `messages` is read BEFORE the SDK method runs, into fresh dicts and strings, sync and async
  alike. The pair is posted later (after the response, after a stream, after an await) but from
  that snapshot: an app that appends the assistant turn to its own list while the call is in
  flight or before the stream is exhausted does not leak the answer into the posted prompt. A
  one-shot iterable (a generator for `messages`, Responses `input`, Anthropic `messages` /
  `system`, GenAI `contents`) is not read ahead, which would empty it for the SDK: it reaches the
  SDK through a recording pass-through, so the provider receives exactly the elements the app
  produces (and, if the generator raises, the same exception at the same point; that request is
  counted `skipped_unreadable_request`), and the snapshot is what the SDK consumed.
- `finish_reason` is the plugin's vocabulary, so the service's `finish_to_review` rule reads
  it; the provider's word is kept in `client.finish_reason_raw`; a missing one stays unset so
  the sidecar marks `finish_unknown` rather than guess. The map:

  | posted | OpenAI chat | OpenAI Responses (`status[:reason]`) | Anthropic | Gemini |
  |---|---|---|---|---|
  | `stop` | `stop` | `completed` | `end_turn`, `stop_sequence` | `STOP` |
  | `length` | `length` | `incomplete:max_output_tokens` | `max_tokens`, `model_context_window_exceeded` | `MAX_TOKENS` |
  | `tool_calls` | `tool_calls`, `function_call` | | `tool_use` | `MALFORMED_FUNCTION_CALL` |
  | `content_filter` | `content_filter` | `incomplete:content_filter` | `refusal` | `SAFETY`, `RECITATION`, `BLOCKLIST`, `PROHIBITED_CONTENT`, `SPII`, `IMAGE_SAFETY`, `LANGUAGE` |
  | `abort` | (stream closed early) | `cancelled`, `in_progress`, `queued` | | |
  | passed through | | `failed` | `pause_turn` | `other`, `unknown` (`OTHER`, `FINISH_REASON_UNSPECIFIED`), `blocked:<reason>` for a blocked prompt |

  A word not in the table passes through lower-cased, cut to the sidecar's 32 characters.
- A conversation longer than the sidecar's 256-message limit keeps its system turn and the
  newest turns; the number dropped is in `client.messages_truncated` (without the cap the
  sidecar answers 422 and the pair is lost).
- Multimodal and tool parts in the PROMPT (images, audio, files, `inline_data`, function calls)
  are reduced to a marker in place, `[image_url omitted]`, so the prompt keeps its shape, and
  counted in `client.parts_dropped` and in `stats()["multimodal_parts_dropped"]`; in the answer
  they are dropped without a marker (above). The bytes never leave the process.
- `usage` is normalised to input / output / total tokens.

The sidecar validates the body with its own `ScoreRequest`; the `client` block is outside that
contract today and is ignored by pydantic, not refused (`tests/test_sdk.py` posts the body to
the real app). Token ids are not posted (the SDKs do not expose them).

## What is never posted

- API keys, headers, the request's other parameters (temperature, tools, response formats,
  metadata): only `model` and the message text are read.
- Image, audio and file bytes or URLs (reduced to a marker, see above).
- Thinking / reasoning blocks (Anthropic `thinking`, Gemini `thought` parts, Responses
  `reasoning` items).
- Anything when the call failed: the SDK's exception propagates unchanged and no pair exists.
- Anything the `redact` hook returned `None` for, or raised on (a failed masker never lets the
  unmasked pair through; the drop is counted).

## Streaming

A stream is returned to the app through a thin proxy that yields the same items and delegates
every other attribute (`.response`, `.close()`, `with …`), and the pair is posted **once, at
the end**: when the iterator is exhausted (OpenAI chat chunks and Responses events, Anthropic
raw events, Gemini chunks are assembled into the full text; the last chunk's finish and usage
are taken), or when **the proxy's own** `close()` or context exit runs before that, in which
case the partial text is posted with `finish_reason: "abort"` and `client.stream_complete:
false`. A stream the app simply stops reading (`for chunk in s: break` with no `close()` and no
`with` on `s`) is never posted: it is counted `stream_unfinished` when the proxy is
garbage-collected (a finalizer counts, it does not post). OpenAI's own helper managers,
`client.chat.completions.stream()` and `client.responses.stream()`, go through the patched
`create` and capture once when run to exhaustion; their `__exit__` closes the HTTP response,
not the stream, so an early exit from one of them is the unfinished case above, not `abort`.
One consequence: `chat.completions.stream()` can raise `LengthFinishReasonError` from its own
parsing after the final chunk arrived; the raw stream under it was exhausted, so the pair IS
posted, as `length`, although the app got an exception rather than a completion (E-R149). `abort` is only ever
the app closing a healthy stream early: a stream that raised mid-way, or a `with` block an
exception left (the SDK raising on an `error` event, Anthropic's `overloaded_error` included,
or the app raising inside the block), posts nothing and is counted `stream_errored`. Anthropic's
`messages.stream()` manager hands the app the SDK's `MessageStream` through a thin proxy (same
events, `text_stream`, `get_final_message()` and every attribute delegated) that remembers a
failure raised while reading, so an `overloaded_error` the app catches inside the block still
posts nothing; a clean exit reads `current_message_snapshot` (no extra network read), and a
healthy stream the app stops reading early posts the partial text as `abort`. Only the first
choice / candidate is assembled when `n > 1`.

## Redaction (default: none)

By default the SDK posts the text as the app sent and received it, because masking changes what
Latent reads: E-R143 (2026-10-01) ran a default de-identification tool over logged prompts and
answers masked separately and a third of the good answers turned into judged failures (the
answer named "Felica Cannon", the masked source only "Felica"), the probe lost up to 0.07 AUROC
and the flagged list overlapped the unmasked one about 50%. Identifiers-only masking with ONE
dictionary derived from the source and applied to prompt and answer alike kept the judge at
0.96 and the read within 0.02. If you must mask before posting, do that, and do it in one place:

```python
def redact(body):
    table = my_entities(body["messages"])             # names from the source, one dictionary
    body["messages"] = [dict(m, content=mask(m["content"], table)) for m in body["messages"]]
    body["output"] = mask(body["output"], table)      # the same table on the answer
    return body                                        # or None to skip this pair

runlatent.auto_instrument(redact=redact)
```

The hook sees the full body (the `client` block included) and may rewrite or drop it. Prefer
raw text under the processing terms (retention is the service's `policy.retention`) and mask
after the read where a report needs it.

## Counters and shutdown

```python
runlatent.stats()
# {"enabled": true, "api_url": "...", "patched": ["openai:Completions.create", ...],
#  "captured": 120, "posted": 118, "failed": 0, "dropped": 2, "queued": 0,
#  "skipped_empty_output": 3, "multimodal_parts_dropped": 7, "redact_skipped": 0,
#  "errors": {"post:ConnectionRefusedError": 0}, "poster": {... the poster's own counters ...}}
runlatent.flush(timeout=5.0)   # at shutdown: wait for the queue to be attempted
```

`captured` counts pairs assembled; `posted` / `failed` / `dropped` / `queued` are the poster's
(`latent_vllm.events.SinkPoster` through `latent_sidecar.tap.FireAndForget`: daemon thread,
bounded queue, three attempts per pair, no spill file: a pair that cannot be delivered is
lost and counted). `dropped` includes pairs captured with no destination configured.
`errors` is keyed by stage and exception class (`observe:…` a shape the adapter could not read,
`redact:…`, `post:…`, `stream_feed:…`); nothing in it ever reached the app. Only the class is
logged (at DEBUG), never the exception's message: a masker that raises with its input in the
message must not put answer text in your log.

## Limits

- Only the methods above are patched. `chat.completions.parse()` and `responses.parse()` post
  through their own path and are not captured; `beta.*`, the Assistants and Realtime APIs and
  other providers are not in this version; wrap a gateway or use `latent_sidecar.tap` for those.
  Captured because they route through the patched methods: `chat.completions.stream()` and
  `responses.stream()` (exhaust-to-capture, see Streaming) and Gemini `chats.send_message`.
- Arguments must be passed by keyword (the SDKs require it; a positional `messages` is not read).
- `with_raw_response.create(...)` and `with_streaming_response.create(...)` return the HTTP
  response object, not the parsed completion; reading it would consume a streaming body, so
  those calls are not captured and are counted under `stats()["skipped_raw_response"]`.
- At a normal interpreter exit an atexit handler waits up to `LATENT_SDK_EXIT_FLUSH_S`
  seconds (default 2.0) for the queued pairs to be posted, so a short script that never calls
  `flush()` still delivers; with it set to `0` nothing waits and pairs still queued at exit are
  lost without a counter (the process is gone). A `kill -9` or an `os._exit()` loses them either
  way: call `runlatent.flush()` in your shutdown path when it matters.
- The proxy around a stream is not an instance of the SDK's `Stream` class; `isinstance`
  checks on it fail (attribute access does not).
- There is no on-disk spill: a Latent outage longer than the retries loses those pairs, counted.
- `max_queue` counts pairs, not bytes, behind one serial poster thread (ingest tops out near
  1 / POST latency pairs per second). A full queue of 10,000 pairs with 100 KB RAG prompts is
  about 1 GB held in the process before anything is dropped: size `max_queue` to the burst you
  are willing to hold in memory, and lower it on memory-tight hosts.
- Multi-choice (`n > 1`) and multi-candidate responses post the first only.
- Gemini thinking models count their thought tokens against `max_output_tokens`: a short budget
  yields many empty answers (not posted, `skipped_empty_output`) and truncated ones, posted as
  `length` and routed to review by the service's default `finish_to_review` rule (E-R149's
  first Gemini run at 160 tokens produced 44 of them). Give thinking models a budget that covers
  the thinking, or disable it, before reading the review queue as a quality signal.
