Metadata-Version: 2.4
Name: batchwatch
Version: 0.2.1
Summary: Client for batchwatch.dev - measure queue time on LLM batch APIs without ever blocking your job
Author: Andreas Graae
License: MIT License
        
        Copyright (c) 2026 Andreas Graae
        
        Permission is hereby granted, free of charge, to any person obtaining a copy
        of this software and associated documentation files (the "Software"), to deal
        in the Software without restriction, including without limitation the rights
        to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
        copies of the Software, and to permit persons to whom the Software is
        furnished to do so, subject to the following conditions:
        
        The above copyright notice and this permission notice shall be included in all
        copies or substantial portions of the Software.
        
        THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
        IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
        FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
        AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
        LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
        OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
        SOFTWARE.
        
Project-URL: Homepage, https://batchwatch.dev
Project-URL: Source, https://github.com/batchwatch/client
Keywords: llm,batch,openai,anthropic,queue,latency,observability
Classifier: Development Status :: 3 - Alpha
Classifier: Intended Audience :: Developers
Classifier: License :: OSI Approved :: MIT License
Classifier: Programming Language :: Python :: 3
Classifier: Topic :: System :: Monitoring
Requires-Python: >=3.8
Description-Content-Type: text/markdown
License-File: LICENSE
Provides-Extra: test
Requires-Dist: pytest>=7; extra == "test"
Dynamic: license-file

# batchwatch — Python client

Client for [batchwatch.dev](https://batchwatch.dev): crowdsourced measurement
of queue time on LLM batch APIs.

Batch endpoints cost 50% of the synchronous ones, but "completes within 24
hours" is impossible to plan around. batchwatch measures what the queue
actually does and answers one question: *should I use batch for this job?*

Standard library only. No dependencies, and none planned.

## Install

Install straight from the repo:

    pip install git+https://github.com/batchwatch/client#subdirectory=python

or copy `src/batchwatch/` into your project — it is two files. A PyPI release is
on the way.

## Two lines

```python
from batchwatch import Batchwatch

bw = Batchwatch(token="bw_...")        # token optional; falls back to $BATCHWATCH_TOKEN

# 1. before you submit — does this belong in the queue?
if bw.should_batch("gpt-5.6-sol", max_wait="15m"):
    job = client.batches.create(...)
else:
    answer = client.chat.completions.create(...)

# 2. measure it, so the next person gets a better answer
with bw.track("gpt-5.6-sol", input_tokens=9720) as t:
    result = wait_for(job)
    t.done(output_tokens=result.usage.completion_tokens)
```

Get a key with no email and no card:

    curl -X POST https://batchwatch.dev/v1/keys -d '{"label":"my pipeline"}'

## The deadline guard — batch when you can, sync when you must

Batch is half the price, but a queue that misses your deadline can take down a
product. `bw.batch(...)` gets you both: it runs your batch, watches the clock,
and if the batch has not finished by your deadline it cancels it and runs your
synchronous fallback instead — so your job always gets an answer, on time.

```python
with bw.batch("gpt-5.6-sol", deadline="15m",
              on_deadline=lambda: client.chat.completions.create(...)) as job:
    job.submit(lambda: client.batches.create(...))
result = job.result()
```

You hand over two callables — the batch-create and the sync fallback — and the
client runs them. It never sees or builds your provider payload; there is no
field for it, exactly as with `track()`. If the batch finishes in time you get
its result; if the deadline fires you get the fallback's result.

Every fallback is a **measured prediction outcome**: "batching would have missed
— running sync was right." It goes down the same accuracy path `should_batch()`
already feeds, so the server can score how often the guard was needed. Nothing
new leaves the machine.

`deadline` speaks the same duration language as `should_batch`'s `max_wait` —
`"15m"`, `"6h"`, `"30s"`, a bare number of seconds, or `None` for no guard. The
defaults duck-type the OpenAI / Anthropic batch shape; a different provider
passes `poll=` / `cancel=` / `result_of=` callables.

### The poll loop is ours, not yours

`result()` owns the wait so you do not write the same sleep/backoff loop
everyone else does. It polls with **exponential backoff and jitter**, never
faster than a **rate-limit floor** (a naive one-second loop against a 24-hour
job is 86,400 requests and an angry provider) and never slower than a ceiling.
The **first** interval is informed by the model's own **measured p50** — no
reason to poll every five seconds against a model whose median queue time is
forty minutes — and reading that p50 **fails open**: if batchwatch is
unreachable, polling simply continues on the fixed fallback schedule.

The 24-hour **expiry is a distinct terminal state**, never silently a timeout:
`job.expired` is `True` when the batch hit the cutoff, so you can tell "the
queue was slow" from "the batch failed".

Tune the cadence if you need to; the defaults are sane:

```python
with bw.batch("gpt-5.6-sol", deadline="6h",
              on_deadline=lambda: client.chat.completions.create(...),
              poll_base=30, poll_floor=5, poll_ceiling=900) as job:
    job.submit(lambda: client.batches.create(...))
result = job.result()
```

Driving an async or worker loop yourself instead of blocking a thread? Step the
same state machine one poll at a time:

```python
state = job.poll_once()          # PollState.RUNNING / DONE / EXPIRED / FAILED
if state is PollState.RUNNING:
    await asyncio.sleep(job.next_interval)   # backoff + jitter, already applied
```

### Partial completion — landed, failed, expired

A batch of 20,000 requests is not binary: some land, some fail per-request, and
some never return before the 24-hour expiry. `job.split()` separates the three,
mapped back to your own objects **by `custom_id`** (never by index — provider
ordering is not guaranteed):

```python
result = job.split(results=downloaded_lines, custom_ids=every_submitted_id)
result.landed        # count that came back clean
result.failed        # count that failed per-request
result.expired       # count still outstanding at the 24h cutoff
result.complete      # True only if EVERYTHING landed — never a silent success
result.result_for("my-request-42")     # your object for one id, mapped correctly

# retry only what failed — idempotent, so a second call submits nothing new
child = job.retry_failed(lambda failed_ids: client.batches.create(...))
```

An expired job is reported to batchwatch as `status="expired"`, which is
recorded but kept out of the percentiles (only `completed` rows count) — so a
slow queue neither pollutes p90 nor loses the "the queue was slow" signal. That
measurement rule lives with the ingest contract on the server, under
`/v1/calls/complete` → "Partial completion", not only here in the client.

## It fails open, always

If batchwatch is down, slow, or broken, your job must not notice. That is the
first requirement, ahead of collecting any data at all.

- Every submission runs on a daemon thread. `track()` does no network I/O on
  your thread.
- Two-second timeout by default (`BATCHWATCH_TIMEOUT`).
- Every batchwatch error is swallowed and logged at `DEBUG` on the
  `batchwatch` logger. Nothing is printed unless you ask for it.
- `should_batch()` is the one synchronous call, because you are waiting for
  the answer. If it cannot answer, you get your own `default` back — never a
  guess. The default is `False`, "run it synchronously": being wrong that way
  costs money, being wrong the other way blows a deadline.
- An exception raised inside your own `with` block is recorded as `failed`
  and re-raised untouched. We swallow our errors, never yours.

`tests/test_fail_open.py` proves it against a port nothing listens on and
against a socket that accepts but never answers.

## It never sends your content

No prompts, no completions, no system prompts, no tool calls, no file names.
The request body is built from a fixed allowlist — provider, model, mode,
endpoint, request count, token counts, timestamps, status — and everything
else is dropped in `_scrub()` on the way out. There is no field to put text
in.

`tests/test_no_content.py` asserts it on the bytes a real HTTP server
received, and includes a positive control so the test cannot pass by the
client simply sending nothing.

## `output_tokens` defaults to `None`, never `0`

You know your input tokens. You cannot know your output tokens before the
model has answered. So the default is absence, not zero.

Zero is not a harmless placeholder here: output costs five to six times as
much as input, so a saving computed on zero output is systematically too
low — measured at 3.4x too low on a real model — and nothing in the response
would tell you. If you know a ceiling, pass `max_tokens` instead and the
answer comes back labelled as a ceiling.

## Spooling

When a measurement cannot be delivered, the completed record is appended to a
JSONL file and replayed later through `POST /v1/calls/complete`. Losing
measurements exactly when the network is bad means losing them exactly when
they are most interesting.

- Default path: `$BATCHWATCH_SPOOL`, or `batchwatch-spool.jsonl` in the
  system temp directory. Set `BATCHWATCH_SPOOL=""` or pass `spool=None` to
  turn it off.
- The spool is replayed automatically, at most once a minute, right after a
  successful call — that is the moment we know the network is up. Call
  `bw.flush_spool()` yourself from a shutdown hook if you want it drained on
  exit.
- **Spooling requires a token.** `/v1/calls/complete` takes your own
  timestamps, so it is closed to anonymous callers; without a key a spool
  file could never be sent, and writing one would just leak disk. Without a
  token, undeliverable measurements are dropped and logged at `DEBUG`.
- The file is capped at 5 MB. Beyond that, measurements are dropped rather
  than filling your disk.
- A replayed measurement can arrive twice if the original `PATCH` reached the
  server but the response did not. That is deliberate: a duplicate is visible
  in the dataset, a lost measurement is not.
- Threads are handled. Two *processes* sharing one spool file may send a
  record twice — give each process its own `BATCHWATCH_SPOOL` if that matters.

## Configuration

| Argument | Environment | Default |
|---|---|---|
| `token` | `BATCHWATCH_TOKEN` | none (anonymous) |
| `base_url` | `BATCHWATCH_URL` | `https://batchwatch.dev` |
| `timeout` | `BATCHWATCH_TIMEOUT` | `2.0` seconds |
| `spool` | `BATCHWATCH_SPOOL` | `<tempdir>/batchwatch-spool.jsonl` |
| `enabled` | — | `True` |

`enabled=False` turns every network call into a no-op, which is what you want
in CI.

## API

- `should_batch(model, max_wait=None, default=False, **kw) -> bool`
- `advice(model, max_wait=None, provider="openai", input_tokens=None, output_tokens=None, max_tokens=None, risk="p90") -> dict | None`
- `wait_now(model, provider="openai", mode="batch") -> dict | None`
- `track(model, provider="openai", mode="batch", requests=1, input_tokens=None, endpoint=None)` — context manager
  - `t.done(output_tokens=None, status="completed", ttfb_ms=None)`
  - `t.failed()`
  - `t.started(input_tokens=...)` when the count is only known after submission
- `batch(model, deadline=None, on_deadline=None, provider="openai", **kw) -> BatchJob` — the deadline-guarded, self-polling job (context manager)
  - `job.submit(create)` — run the caller's batch-create callable, remember the handle
  - `job.result()` — block, polling with backoff + jitter, or fall back to `on_deadline()` at the deadline
  - `job.poll_once() -> PollState` — one non-blocking step for an async/worker loop; `job.next_interval` is the seconds to sleep before the next
  - `job.fell_back` — `True` once the guard has fired; `job.expired` — `True` if the batch hit the 24h expiry; `job.poll_count` — polls made
  - `job.split(results=None, custom_ids=None) -> BatchResult` — split landed / failed / expired, mapped by `custom_id`
  - `job.retry_failed(resubmit, result=None) -> BatchJob | None` — resubmit only the failed subset, idempotently
  - poll cadence: `poll_base` / `poll_floor` / `poll_ceiling` / `poll_backoff` / `poll_jitter` / `use_p50_cadence`
  - override the provider shape with `poll=` / `cancel=` / `result_of=` / `succeeded=` / `expired=`
- `flush(timeout=5.0) -> bool` — wait for outstanding submissions before exit
- `flush_spool(timeout=None) -> int` — send what is on disk, returns accepted

## Tests

    python -m pytest -q

91 tests, no network beyond loopback. They start real HTTP servers on
ephemeral ports rather than monkeypatching `urllib`: the thing under test is
network behaviour, so the network should be in the test. The poll-loop tests
use an injected clock, sleep and rng, so the backoff schedule is pinned exactly
with no real sleeps.

## Licence

MIT
