feat(traces): span export queue with retries - #954
Conversation
posthog-python Compliance ReportDate: 2026-09-17 03:42:47 UTC ✅ All Tests Passed!111/111 tests passed Capture_V1 Tests✅ 94/94 tests passed View Details
Feature_Flags Tests✅ 17/17 tests passed View Details
|
Important Files Changed
Prompt To Fix All With AI### Issue 1
posthog/tracing/_export.py:449-452
`timer.start()` can run the callback before `_flush_timer` points to the new timer. A zero-delay flush can then look stale and return. The old timer is canceled, and the dead timer is stored. Publish the replacement before starting it, and restore the old timer if `start()` fails.
---
For each issue above, determine whether it is valid and should be fixed. If so, fix it directly.Reviews (2): Last reviewed commit: "fix(traces): keep the scheduled flush wh..." | Re-trigger Greptile |
d51f793 to
d817b98
Compare
d817b98 to
57fabd6
Compare
57fabd6 to
2eb751c
Compare
2eb751c to
0a4baaf
Compare
dustinbyrne
left a comment
There was a problem hiding this comment.
Looks good for the export-queue slice, with one non-blocking note: _export.py:182 calls warn_if_due(force=True) after every flush, including automatic flushes. Repeated rejected batches can therefore emit a warning per flush rather than at most once per configured interval. Consider reserving forced reporting for shutdown/exit and updating the fixed-clock test that currently expects two warnings in one interval.
The transport and pipeline feedback belongs to their respective stack PRs; this approval does not cover the complete public tracing release.
AI-assisted review using source, tests and existing CI; no new tests were run.
0a4baaf to
d7aa5d3
Compare
|
Thanks. We warn after every flush on purpose, to match Node. It calls |
dustinbyrne
left a comment
There was a problem hiding this comment.
Thanks, confirmed against Node: the flush path calls _warnAboutDrops(), with an explicit once-per-flush test. I withdraw my warning-cadence suggestion. The spec wording should be reconciled separately, not by changing Python alone. Approval stands.
AI-assisted follow-up source review, including Node parity and the rebased transport change; no new tests run locally.
jzhu13
left a comment
There was a problem hiding this comment.
Reviewed against traces/05-pipeline. Tests pass at the head. Lock ordering is consistent, _send runs with no lock held, the 413 halving cannot reach zero or loop, and the Retry-After window handles units, clamp, and negatives correctly. Two items I would fix before merge; both surface as user-visible data loss once #957 wires this in.
Blocking
posthog/tracing/_export.py:348aflush(timeout)with budget remaining never retries aretry-laterbatch:_apply_outcome_lockedreturnsstop=True,_drainpropagates it,flushskips the follow-up pass, and nothing sleeps. Reproduced: one 503 thenok;flush(30.0)made one attempt in 0.000 s and left the queue intact. #957's shutdown then callsclose()and discards the whole backlog with a "Discarding N span(s)" warning, while the events lane retries up to ten times with sleeps. Suggest: whentimeout is not Noneand the pass ended in retry-later,time.sleep(min(self._next_flush_delay_locked(), remaining))and loop_drainwhile budget remains, charging normally. Timer-driven flushes (_timer is not None) keep the current no-sleep behavior.posthog/tracing/_export.py:155the deadline is computed before_flush_lock.acquire(timeout=...), and acquire failure returns silently with zero requests and no log. That contradicts the docstring's "no request starts once it is spent, except the first", which #957 repeats to users inClient.flush. Reproduced: a timer flush parked in_send,flush(0.2)returned at 0.20 s with 0 send attempts and 2 spans still queued. Either start the budget after the lock is acquired, or narrow the contract and debug-log the give-up path.
Non-blocking
posthog/tracing/_export.py:317a server 413 caused by one oversized span collapses_max_export_batch_sizeto 1 and the +1 ramp keeps it small for a long time. Reproduced: 1 giant + 511 small spans at batch 512 took 42 requests (sizes 512, 256, ..., 1, 1, 2, 3, ...) where 2 would do, and left the size at 33 so the depth trigger fired every ~35 spans afterwards. Restore the pre-halving size once the size-1 413 isolates the culprit, and ramp multiplicatively.posthog/tracing/_export.py:350dropping a head batch after 8 windows resets_consecutive_failures, so the exporter immediately sends the next batch to an endpoint it just saw fail 8 times and re-enables the depth trigger. Keep the count on a retry-later drop.posthog/tracing/_export.py:29duplicates the 30 s ceiling thatcapture_v1._MAX_BACKOFF_SECONDSdefines and AGENTS.md calls "the single ceiling for both the exponential backoff and the Retry-After clamp". Import it.- Nits: ten fields track failure and timer state where one
_next_attempt_atdeadline would replace_jitterand_head_batch_chargeable_at;FakeTimer.fire()ignorescancelled, so a test can pass by firing a timer the code already cancelled; shared fixtures are imported fromhelpers.pyand re-exported via__all__rather than aconftest.py.
Reviewed with Claude Code (Claude Fable 5.1). Behaviors above were reproduced against this branch head.
d7aa5d3 to
a451080
Compare
PR overviewAll previously flagged issues have been addressed. No open security concerns remain on this pull request. Security reviewNo open security issues remain on this pull request. Fixed/addressed: 1 · PR risk: 0/10 |
|
Thanks. Done in a451080: item 4 (a dropped batch keeps the backoff), item 5 (ceiling imported from capture_v1) and the |
a451080 to
79869b9
Compare
|
Follow-up, 79869b9:
|
Adds SpanExporter, the in-memory queue ended spans wait in until they are batched and sent, separate from the events queue. A timer flushes every flush_interval and a full batch flushes at once; only one flush runs at a time. A full queue drops the incoming span, never a queued parent. Failures back off exponentially with jitter, floored by Retry-After (clamped to 30 s, extended by a later deadline but never shortened), and automatic sends pause while backing off. A batch refused across 8 backoff windows is dropped, a 413 halves the batch and ramps back, and other 4xx drop it. flush(timeout) always sends the first batch, so a serverless handler with no budget left still ships spans. Not reachable from the client. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TkZAsCciW4PV8ZdcCHmAbA
…to start, and count an unencodable span once The replacement timer now starts before the old one is cancelled, so a thread that cannot be created leaves the earlier flush in place. Records that fail to encode leave the queue before the send is settled, so a failed send no longer counts them a second time.
Dropping the head batch after eight windows reset the failure count, so the next batch went straight to an endpoint that had just failed eight times and the depth trigger came back on. The backoff ceiling is now the events lane's single constant, and FakeTimer refuses to fire a timer the code cancelled.
…e budget once the lock is held A flush with budget left stopped at the first retriable failure, so a 30 s shutdown flush made one attempt and then discarded the backlog. A caller-driven flush now waits out the backoff and retries while budget remains, with one last attempt at the deadline; timer flushes and untimed flushes are unchanged. The budget starts once no other flush is in flight, and a flush that never gets the lock says so at debug. A size-1 413 restores the batch size it halved from, and the ramp doubles instead of adding one. The resource is encoded once per exporter.
79869b9 to
71d931e
Compare
💡 Motivation and Context
Adds
SpanExporter, the in-memory queue ended spans wait in until they are batched and sent, separate from the events queue.flush_interval; a full batch flushes at once; only one flush runs at a time.Retry-After(clamped to 30 s, extended by a later deadline but never shortened). Automatic sends pause while backing off.flush(timeout)always sends the first batch, so a serverless handler with no budget left still ships spans.Not reachable from the client yet.
Stack (PR 6 of 9, based on
traces/05-pipeline):traces/01-ids-traceparenttraces/02-otlp-encodingtraces/03-span-handlestraces/04-transporttraces/05-pipelinetraces/06-export← this PRtraces/07-span-limitstraces/08-before-span-sendtraces/09-client-wiring💚 How did you test it?
Unit tests in
posthog/test/tracing/test_export.pycover timer and size flushes, backoff andRetry-After, 413 halving, drop rules and the boundedflush(timeout).📝 Checklist
If releasing new changes
sampo addto generate a changeset file🤖 Agent context
Autonomy: Human-driven (agent-assisted)
Implemented with Claude Code (Claude Opus 5) against the traces spec, one commit per slice so each PR reviews on its own. Rebased onto main and opened as a stacked draft in a later Claude Code session (Claude Fable 5.1).
🤖 Generated with Claude Code
https://claude.ai/code/session_012o7CtHLfcypjmXL7g9ZGRC