Skip to content

feat: harden log ingestion delivery - #846

Open
Abhijeet Prasad (AbhiPrasad) wants to merge 4 commits into
mainfrom
abhi/harden-log-ingestion
Open

Abhijeet Prasad (AbhiPrasad) wants to merge 4 commits into
mainfrom
abhi/harden-log-ingestion

Conversation

@AbhiPrasad

@AbhiPrasad Abhijeet Prasad (AbhiPrasad) commented Oct 2, 2026 •

Copy link
Copy Markdown
Member

Summary

Hardens Python SDK log delivery by giving /logs3 a dedicated policy-aware writer. The writer prepares batches once, sends single-attempt requests through bounded ingestion pools, and owns retry scheduling and throttling.

Transient failures are retained and retried in the background through repeated bounded retry cycles. Permanent log endpoint failures (400, 401, 403, 413, and other non-retryable responses) are reported and released so they cannot block newer rows. An explicit flush() waits for an active delivery up to its own timeout, then sends queued rows; a failed flush raises BraintrustLogFlushError with the rows still retained. Application operations that flush as part of their work (including Eval finalization, summarize, and Logger(async_flush=False)) wait for delivery but report telemetry failures without failing the application call.

The queue remains bounded and may drop records when full, following the existing queue policy; this change does not prevent overflow loss or claim an improvement over the merge-base. Logging calls do not raise BufferError behind a retained batch. Prepared batches keep their original auth, routing, and connection-pool lease through retries. Before forced relogin, queued rows are sealed into ordered batches under their current identity, even when older batches are still pending. This preserves the credential boundary without requiring the outage to end before login can complete. Expired signed overflow URLs are refreshed after a 403. Child processes reset inherited writer locks and queues, dropping the child's copy of parent pending rows so it cannot duplicate them.

The destination scheduler shares concurrency and cooldown state by ingestion URL and credential scope. It honors Retry-After seconds or dates, staggers recovery starts, and waits outside worker threads. Ingestion has isolated HTTP pools and single-attempt adapters; non-blocking connection pools avoid waiting forever when urllib3 clears pools during interpreter shutdown.

Configuration

  • BRAINTRUST_LOG_MAX_CONCURRENCY (default 4): concurrent ingestion requests per destination.
  • BRAINTRUST_LOG_FLUSH_TIMEOUT (default 60 seconds): explicit flush scheduling and cooldown budget.
  • BRAINTRUST_HTTP_TIMEOUT (default 60 seconds): ingestion request timeout, reduced to the remaining flush budget.
  • BRAINTRUST_NUM_RETRIES (default 2): writer retries beyond the first attempt.
  • BRAINTRUST_QUEUE_SIZE (default 25000): queue capacity.

Validation in this branch

  • Core nox session: 973 passed, 69 skipped, 12 xfailed.
  • Type nox session: pyright, mypy, and 36 runtime tests passed.
  • nox -s pylint, Ruff format/check, pre-commit, and git diff --check passed.
  • Added local HTTP regressions for permanent 4xx/413 responses, transient outage recovery after the initial retry budget, explicit-flush contention and queue draining, expired signed URL refresh, queue producer behavior, fork state isolation (including per-LazyValue metadata locks), and login identity boundaries. The background permanent-failure test waits for queue reservation release before asserting the writer is drained, and a later explicit flush now reports background-dropped batches.
  • The local HTTP benchmark harness now completes its explicit final flush and delivered all 2,000 rows in each healthy, throttled, and continuous-throttled scenario. This is not a replacement for the real-stack throughput run.

The local-stack proxy bundle referenced in review was not available in this workspace, so these fixes have not yet been replayed against that stack. A live post-fork HTTP test also could not be completed here: Python 3.14 on macOS segfaulted in proxy_bypass_macosx_sysconf when making a network call from a multithreaded forked child. The fork test here verifies child state reset, not delivery of 200 rows. Please rerun both real-stack checks before merging.

Real-stack findings supplied in review

The first proxy run confirmed the prior revision could hang at exit; stop after permanent errors; stop background delivery after an outage; raise BufferError or flush errors into app code; compute Eval/summarize results before rows arrived; and send queued rows under a replacement user's key. This revision addresses those paths. Later real-stack runs verified the explicit-flush wait, Eval/summarize ordering, login behavior during outages, and forked-child delivery. The most recent 26-scenario rerun found no new regressions against the merge-base. Its remaining fork finding was that per-logger metadata LazyValue locks also need resetting; this revision resets those locks in the child.

The earlier real-server benchmark (3 runs × 20,000 spans, large queue) measured 2,317 vs. 2,385 rows/second for the prior revision vs. merge-base, about 3% slower; flush was about 10% slower and peak memory about 8% higher. The follow-up run did not measure throughput. Earlier default-queue burst runs show that records can still be dropped when the bounded queue fills (9,000–24,000 dropped in a 20,000-span burst, comparable to the merge-base); this change makes no queue-loss or throughput improvement claim.

The review also found improvements this revision preserves: Retry-After seconds and date values are honored; a burst of 12 consecutive 503s across 8 threads delivered 1,600/1,600 rows where the merge-base silently lost 200; and the real signed-upload path delivered large rows over 5 MB.

@chatgpt-codex-connector

chatgpt-codex-connector Bot commented Oct 2, 2026 •

Copy link
Copy Markdown

Codex Review Summary

This comment shows the latest Codex review activity on this pull request.

Review Status Commit Review trigger
📝 Code Review ✅ Completed 2026-10-02T17:46:04.084003Z 0ea712f New commits
ℹ️ About Codex in GitHub

Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you

Codex reacts with 👀 while any review is running, comments if it has suggestions, and reacts with 👍 once all reviews finish with no findings.

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

💡 Codex Review

Here are some automated review suggestions for this pull request.

Reviewed commit: 9b3978269d

ℹ️ About Codex in GitHub

Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you

  • Open a pull request for review
  • Mark a draft as ready
  • Comment "Codex (@codex) review".

If Codex has suggestions, it will comment; otherwise it will react with 👍.

Codex can also answer questions or update the PR. Try commenting "Codex (@codex) address that feedback".

Comment on lines +1329 to +1331
if self._pending:
if _background:
return # Keep retained work parked; only explicit flush promises delivery.

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Upload attachments for batches that already succeeded

When a flush is split into multiple batches and one batch remains pending while another succeeds, this early exit skips every attachment because _pending_attachments is global to the whole wave. The successful batch has already been accepted with attachment references, but if the other batch is permanent (for example, HTTP 413), subsequent flushes always take this branch and those attachments are never uploaded, leaving permanently broken references for successfully ingested rows. Track attachments per batch and upload those belonging to completed batches before reporting the retained failures.

Useful? React with 👍 / 👎.

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

💡 Codex Review

Here are some automated review suggestions for this pull request.

Reviewed commit: bd8639376a

ℹ️ About Codex in GitHub

Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you

  • Open a pull request for review
  • Mark a draft as ready
  • Comment "Codex (@codex) review".

If Codex has suggestions, it will comment; otherwise it will react with 👍.

Codex can also answer questions or update the PR. Try commenting "Codex (@codex) address that feedback".

Comment on lines +1712 to +1714
await asyncio.get_event_loop().run_in_executor(
None, lambda: experiment.state.flush_best_effort()
)

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Propagate flush failures before remote trace reads

When a scorer calls trace.get_thread() during an ingestion outage, flush_best_effort() swallows BraintrustLogFlushError, so LocalTrace._ensure_spans_ready() treats the callback as successful and permanently sets _spans_flushed=True. The remote preprocessor then reads the backend before the task spans arrive, potentially producing an empty/incomplete thread and an incorrect score; this synchronization callback should preserve flush failures so a later caller can retry.

Useful? React with 👍 / 👎.

Comment on lines +108 to +112
def restore(self, items: list[T]) -> None:
"""Return unprepared records ahead of newer records without losing their capacity."""
with self._mutex:
self._reserved -= len(items)
self._queue.extendleft(reversed(items))

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Preserve concurrent records when restoring a failed wave

If record preparation fails after drain_all(reserve=True) while producers add records to the bounded queue, deque.extendleft() can exceed maxlen and silently discard those newly queued records from the right. A full drained wave makes even one concurrent record vulnerable, and the loss is neither returned nor counted in _total_dropped; the reservation must prevent this overflow or restoration must explicitly retain/account for displaced records.

Useful? React with 👍 / 👎.

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

💡 Codex Review

Here are some automated review suggestions for this pull request.

Reviewed commit: bdbca87607

ℹ️ About Codex in GitHub

Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you

  • Open a pull request for review
  • Mark a draft as ready
  • Comment "Codex (@codex) review".

If Codex has suggestions, it will comment; otherwise it will react with 👍.

Codex can also answer questions or update the PR. Try commenting "Codex (@codex) address that feedback".

attachment_errors: list[Exception] = []
for attachment in list(self._pending_attachments):
try:
result = attachment.upload()

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Do not retain a cached attachment failure

When an object-store upload fails but /attachment/status succeeds, Attachment.upload() returns an error status through self._uploader.get() (logger.py:3896-3902), and that LazyValue permanently caches the error result. Leaving the attachment in _pending_attachments therefore makes every subsequent flush receive the same cached error without retrying the network operation; _retained_count never clears, and later queued logs cannot be prepared. Either recreate the uploader before retaining the attachment or report and release this non-retryable cached result.

Useful? React with 👍 / 👎.

self._pending_attachments.extend(attachments)
else:
self._pending = prepared
self._pending_attachments = attachments

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Keep attachment uploads on their batch credentials

When queued rows containing attachments are sealed before force_login=True, their prepared batches retain the original ingestion service, but this unscoped attachment list does not retain that credential lease. After the state swap, Attachment.upload() authenticates through the module-global _state (logger.py:3960-3965), so the old tenant's accepted row can reference an attachment registered and uploaded under the replacement tenant. Associate attachments with their prepared batch/service and upload them before releasing that original credential context.

Useful? React with 👍 / 👎.

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

💡 Codex Review

Here are some automated review suggestions for this pull request.

Reviewed commit: 0ea712f243

ℹ️ About Codex in GitHub

Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you

  • Open a pull request for review
  • Mark a draft as ready
  • Comment "Codex (@codex) review".

If Codex has suggestions, it will comment; otherwise it will react with 👍.

Codex can also answer questions or update the PR. Try commenting "Codex (@codex) address that feedback".

Comment on lines +1791 to +1793
_state._fork_login = (_state.app_url, _state.login_token, _state.org_name)
if _state._global_bg_logger.has_succeeded:
_state._global_bg_logger.value.reset_after_fork()

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Reset explicitly created states after fork

The new at-fork callback resets only the module-global _state. When an application creates a separate BraintrustState and passes it through the supported state= argument, then forks after that state's writer has started, the child retains writer.started=True even though its publisher thread no longer exists, along with the parent's pending queue and potentially locked synchronization objects. Subsequent asynchronous logs in the child therefore remain unpublished, while an explicit flush can duplicate inherited parent rows or deadlock; track and reset every live state (or install an instance-aware fork hook), not only the global singleton.

Useful? React with 👍 / 👎.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant