Skip to content
Closed
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
17 changes: 17 additions & 0 deletions examples/log_ingestion/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
# Policy-aware log ingestion demo

From `py/`, with `BRAINTRUST_API_KEY` configured through mise:

```bash
mise exec -- python ../examples/log_ingestion/demo.py --replay
```

The demo sends a root/child trace, flushes a score update, relogs in, and sends a final score and comment. `--replay` sends the same memoized feedback records again, preserving their row identities. The JSON output includes the trace link and remaining record count. Inspect the trace with:

```bash
bt view trace --object-ref project_logs:<project_id> --trace-id <trace_id> --json
```

Expect two spans, final root score `quality=1`, and one comment even after replay.

`--overflow` lowers the writer's payload limit to exercise signed uploads on servers that advertise `logs3_payload_max_bytes`. The output reports `overflow_supported`; older servers continue using ordinary ingestion. No provider SDK or model calls are needed.
72 changes: 72 additions & 0 deletions examples/log_ingestion/demo.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,72 @@
"""Send nested traces and ordered score updates through the ingestion transport.

From py/: mise exec -- python ../examples/log_ingestion/demo.py
Add --overflow to exercise a signed overflow upload on a supporting server.
"""

import argparse
import json
import uuid

import braintrust
from braintrust.logger import _internal_get_global_state


def main():
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("--project", default="sdk-839-ingestion-demo")
parser.add_argument("--overflow", action="store_true")
parser.add_argument(
"--replay", action="store_true", help="Replay the same final feedback records to check server deduplication"
)
args = parser.parse_args()
logger = braintrust.init_logger(project=args.project)
run_id = str(uuid.uuid4())
state = _internal_get_global_state()
writer = state.global_bg_logger()
writer.sync_flush = True
if args.overflow:
writer._max_request_size_override = 1024
writer._max_request_size_result = None
with logger.start_span(
name="policy-aware-ingestion-demo", input={"run_id": run_id}, metadata={"issue": 839}
) as root:
with root.start_span(name="child", input="hello") as child:
child.log(output="world", scores={"quality": 1})
root.log(
output={"status": "delivered", "payload": "ü" * (2000 if args.overflow else 1)}, scores={"quality": 0}
)
braintrust.flush()
# Two later updates must reach the same row in order, including after relogin.
logger.log_feedback(id=root.id, scores={"quality": 0.5})
braintrust.flush()
braintrust.login(force_login=True)
logger.log_feedback(id=root.id, scores={"quality": 1}, comment="Verified ordered score update after relogin")
replay_records = writer.queue.drain_all(reserve=True) if args.replay else []
if replay_records:
writer.queue.restore(replay_records)
braintrust.flush()
if replay_records:
for record in replay_records:
writer.queue.put(record)
braintrust.flush()
print(
json.dumps(
{
"run_id": run_id,
"project_id": logger.id,
"row_id": root.id,
"trace_id": root.root_span_id,
"link": root.link(),
"pending_rows": writer.pending_count,
"overflow_uploads": writer._overflow_upload_count,
"replayed_feedback": bool(replay_records),
"overflow_supported": (writer._max_request_size_result or {}).get("can_use_overflow"),
},
indent=2,
)
)


if __name__ == "__main__":
main()
51 changes: 51 additions & 0 deletions py/benchmarks/benches/bench_log_ingestion.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,51 @@
"""Preparation and producer cost, independent of HTTP server latency.

Use benchmarks.log_ingestion for end-to-end healthy/throttled measurements.
"""

import contextlib
import os
import pathlib
import sys

import pyperf


if __package__ in (None, ""):
sys.path.insert(0, str(pathlib.Path(__file__).resolve().parents[2]))

from braintrust.logger import _HTTPBackgroundLogger, construct_logs3_data, stringify_with_overflow_meta
from braintrust.util import LazyValue

from benchmarks._utils import disable_pyperf_psutil
from benchmarks.fixtures import ingestion_rows


_ROWS = ingestion_rows()
_RECORDS = [LazyValue(lambda row=row: row, use_mutex=False) for row in _ROWS]


def _prepare() -> None:
construct_logs3_data([stringify_with_overflow_meta(row) for row in _ROWS]).encode("utf-8")


def main(runner: pyperf.Runner | None = None) -> None:
if runner is None:
disable_pyperf_psutil()
runner = pyperf.Runner()
# Pause the real writer's publisher to isolate producer work; no network calls occur.
os.environ["BRAINTRUST_DISABLE_ATEXIT_FLUSH"] = "1"
writer = _HTTPBackgroundLogger(lambda: contextlib.nullcontext(None))
writer.sync_flush = True
writer._start()

def enqueue() -> None:
writer.log(*_RECORDS)
writer.queue.drain_all()

runner.bench_func("logs3[prepare-100-rows]", _prepare)
runner.bench_func("logs3[enqueue-100-rows]", enqueue)


if __name__ == "__main__":
main()
8 changes: 8 additions & 0 deletions py/benchmarks/fixtures.py
Original file line number Diff line number Diff line change
Expand Up @@ -201,3 +201,11 @@ def make_bt_safe_deep_copy_cases() -> list[tuple[str, Any]]:
("circular", make_circular_payload()),
("non-string-keys", make_non_string_key_payload()),
]


def ingestion_rows(count: int = 100) -> list[dict[str, Any]]:
"""Deterministic trace-shaped rows for preparation and HTTP delivery measurements."""
return [
{"id": str(index), "project_id": "benchmark", "input": "hello" * 50, "scores": {"quality": 1}}
for index in range(count)
]
119 changes: 119 additions & 0 deletions py/benchmarks/log_ingestion.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,119 @@
"""End-to-end HTTP ingestion measurements (run from py/).

python -m benchmarks.log_ingestion --output /tmp/ingestion.json
Works on the legacy writer as well, for a same-harness baseline comparison.
"""

import argparse
import contextlib
import json
import os
import statistics
import threading
import time
import tracemalloc

from braintrust.api._test_server import scripted_server
from braintrust.api._transport import HTTPConnection
from braintrust.logger import _HTTPBackgroundLogger
from braintrust.util import LazyValue

from benchmarks.fixtures import ingestion_rows


def measure(throttled, rows=2000, *, continuous=False):
rejected = 0
accepted = 0
arrivals = []
connections = set()
lock = threading.Lock()
started = time.monotonic()

def respond(method, path, body, headers):
nonlocal rejected, accepted
if path == "/version":
return 200, {}, b"{}"
now = time.monotonic()
with lock:
arrivals.append(now - started)
if throttled and now - started < 1:
rejected += 1
return 429, {"Retry-After": "1"}, b"limited"
accepted += len(json.loads(body)["rows"])
return "sleep", 0.005, 200, {}, b"ok"

os.environ["BRAINTRUST_DISABLE_ATEXIT_FLUSH"] = "1"
os.environ["BRAINTRUST_NUM_RETRIES"] = "2"
with scripted_server(respond, persistent=True) as (url, handler):
# Select the branch's ingestion service when available.
try:
from braintrust.api._ingestion import LogIngestionAPI
from braintrust.api._routing import EndpointRouter

connection = LogIngestionAPI(EndpointRouter(app_url=url, api_url=url), "benchmark", concurrency=4)
source = lambda: contextlib.nullcontext(connection)
except ImportError:
connection = HTTPConnection(url)
source = LazyValue(lambda: connection, use_mutex=False)
writer = _HTTPBackgroundLogger(source)
writer.sync_flush = not continuous
writer._max_request_size_result = {"max_request_size": 6_000_000, "can_use_overflow": False}
latencies = []
tracemalloc.start()
started = time.monotonic()
peak_outstanding = 0
for index, row in enumerate(ingestion_rows(rows)):
tick = time.perf_counter()
writer.log(LazyValue(lambda row=row: row, use_mutex=False))
latencies.append(time.perf_counter() - tick)
with lock:
peak_outstanding = max(peak_outstanding, index + 1 - accepted)
if continuous:
time.sleep(0.0005)
queued = writer.queue.size()
tick = time.monotonic()
writer.flush()
flush_seconds = time.monotonic() - tick
elapsed = time.monotonic() - started
_, peak = tracemalloc.get_traced_memory()
tracemalloc.stop()
if accepted != rows:
raise RuntimeError(f"Benchmark delivered {accepted} of {rows} rows")
writer.sync_flush = False
connection.close()
connections.update(getattr(handler, "connections", []))
return {
"rows": rows,
"delivered_rows": accepted,
"rejected_requests": rejected,
"requests": len(arrivals),
"queued_rows_at_flush": queued,
"outstanding_rows_peak": peak_outstanding,
"pending_rows_after_flush": getattr(writer, "pending_count", 0),
"rows_per_second": rows / elapsed,
"flush_seconds": flush_seconds,
"producer_p50_us": statistics.median(latencies) * 1e6,
"producer_p99_us": sorted(latencies)[int(len(latencies) * 0.99)] * 1e6,
"peak_allocated_bytes": peak,
"recovery_seconds": max(arrivals) if arrivals else 0,
"connections": len(connections),
}


def main():
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("--output", required=True)
parser.add_argument("--runs", type=int, default=3)
args = parser.parse_args()
results = {
name: [measure(throttled) for _ in range(args.runs)]
for name, throttled in [("healthy", False), ("throttled", True)]
}
results["continuous_throttled"] = [measure(True, continuous=True) for _ in range(args.runs)]
with open(args.output, "w") as output:
json.dump(results, output, indent=2)
print(json.dumps(results, indent=2))


if __name__ == "__main__":
main()
Loading
Loading