Skip to content

Commit 71d931e

Browse files
committed
fix(traces): retry a refused batch within a timed flush, and start the 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.
1 parent 9c4e893 commit 71d931e

2 files changed

Lines changed: 196 additions & 20 deletions

File tree

‎posthog/test/tracing/test_export.py‎

Lines changed: 124 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
import threading
2+
import time
23
from types import SimpleNamespace
34
from unittest import mock
45

@@ -152,8 +153,9 @@ def test_a_spent_budget_still_ships_one_batch(self):
152153
assert queued(pipeline) == []
153154

154155
def test_returns_without_draining_when_another_flush_holds_the_lock_past_the_deadline(
155-
self,
156+
self, caplog
156157
):
158+
caplog.set_level("DEBUG", logger="posthog")
157159
pipeline, sender, _ = make_traces()
158160
pipeline.start_span("a").end()
159161
timer = FakeTimer.instances[-1]
@@ -165,6 +167,31 @@ def test_returns_without_draining_when_another_flush_holds_the_lock_past_the_dea
165167
assert sender.payloads == []
166168
assert pipeline._exporter._flush_timer is timer
167169
assert not timer.cancelled
170+
assert "another flush was still in flight" in caplog.text
171+
172+
def test_the_budget_starts_once_the_lock_is_held(self, clock):
173+
pipeline, sender, _ = make_traces(max_export_batch_size=1)
174+
exporter = pipeline._exporter
175+
176+
def send_slowly(client, payload):
177+
sender.payloads.append(payload)
178+
clock["now"] += 0.1
179+
return SendOutcome("ok")
180+
181+
exporter._send = send_slowly
182+
for _ in range(2):
183+
pipeline.start_span("a").end()
184+
exporter._flush_lock.acquire()
185+
186+
def release_after_a_while():
187+
time.sleep(0.05)
188+
clock["now"] += 1.0
189+
exporter._flush_lock.release()
190+
191+
threading.Thread(target=release_after_a_while).start()
192+
pipeline.flush(timeout=0.5)
193+
assert len(sender.payloads) == 2
194+
assert queued(pipeline) == []
168195

169196
def test_re_arms_after_a_timer_that_failed_to_start(self):
170197
class FlakyTimer(FakeTimer):
@@ -263,9 +290,19 @@ def test_ramps_the_batch_size_back_up_after_a_413_shrink(self):
263290
for i in range(8):
264291
pipeline.start_span(str(i)).end()
265292
pipeline.flush()
266-
assert pipeline._exporter._max_export_batch_size == 6
293+
assert pipeline._exporter._max_export_batch_size == 8
267294
assert [len(b) for b in sender.batches()] == [8, 4, 4]
268295

296+
def test_restores_the_batch_size_once_a_413_isolates_the_oversized_span(self):
297+
sender = FakeSender(*([SendOutcome("too-large")] * 4), SendOutcome("ok"))
298+
pipeline, _, _ = make_traces(sender=sender, max_export_batch_size=8)
299+
for i in range(8):
300+
pipeline.start_span(str(i)).end()
301+
pipeline.flush()
302+
assert [len(b) for b in sender.batches()] == [8, 4, 2, 1, 7]
303+
assert pipeline._exporter._max_export_batch_size == 8
304+
assert queued(pipeline) == []
305+
269306
def test_a_batch_measured_too_large_locally_splits_only_that_drain(self):
270307
sender = FakeSender(TOO_LARGE_LOCALLY, SendOutcome("ok"))
271308
pipeline, _, _ = make_traces(sender=sender, max_export_batch_size=8)
@@ -911,3 +948,88 @@ def test_a_forked_child_drops_the_inherited_queue_and_timer(self):
911948
)
912949
pipeline.start_span("child-span").end()
913950
assert [r.name for r in queued(pipeline)] == ["child-span"]
951+
952+
953+
def waits_advance(clock, pipeline):
954+
"""Make the exporter's backoff wait move the fake clock instead of sleeping."""
955+
waited = []
956+
957+
def wait(seconds):
958+
waited.append(seconds)
959+
clock["now"] += seconds
960+
return False
961+
962+
pipeline._exporter._wait_for_retry = wait
963+
return waited
964+
965+
966+
class TestRetryWithinBudget:
967+
def test_retries_a_retriable_failure_after_its_backoff(self, clock):
968+
sender = FakeSender(SendOutcome("retry-later"), SendOutcome("ok"))
969+
pipeline, _, _ = make_traces(sender=sender)
970+
waited = waits_advance(clock, pipeline)
971+
pipeline.start_span("a").end()
972+
pipeline.flush(timeout=30)
973+
assert len(sender.payloads) == 2
974+
assert waited == [5]
975+
assert queued(pipeline) == []
976+
977+
def test_makes_a_last_attempt_at_the_deadline(self, clock):
978+
sender = FakeSender(SendOutcome("retry-later"))
979+
pipeline, _, _ = make_traces(sender=sender)
980+
waited = waits_advance(clock, pipeline)
981+
pipeline.start_span("a").end()
982+
pipeline.flush(timeout=12)
983+
# Attempts at 0, 5 and 12: the second backoff of 10 is cut to the budget.
984+
assert len(sender.payloads) == 3
985+
assert waited == [5, 7]
986+
assert len(queued(pipeline)) == 1
987+
988+
def test_honours_a_retry_after_within_the_budget(self, clock):
989+
sender = FakeSender(SendOutcome("retry-later", 60), SendOutcome("ok"))
990+
pipeline, _, _ = make_traces(sender=sender)
991+
waited = waits_advance(clock, pipeline)
992+
pipeline.start_span("a").end()
993+
pipeline.flush(timeout=40)
994+
assert waited == [MAX_RETRY_AFTER_SECONDS]
995+
assert queued(pipeline) == []
996+
997+
def test_a_timer_flush_does_not_retry(self, clock):
998+
sender = FakeSender(SendOutcome("retry-later"))
999+
pipeline, _, _ = make_traces(sender=sender)
1000+
waited = waits_advance(clock, pipeline)
1001+
pipeline.start_span("a").end()
1002+
FakeTimer.instances[-1].fire()
1003+
assert len(sender.payloads) == 1
1004+
assert waited == []
1005+
1006+
def test_a_flush_without_a_timeout_does_not_retry(self, clock):
1007+
sender = FakeSender(SendOutcome("retry-later"))
1008+
pipeline, _, _ = make_traces(sender=sender)
1009+
waited = waits_advance(clock, pipeline)
1010+
pipeline.start_span("a").end()
1011+
pipeline.flush()
1012+
assert len(sender.payloads) == 1
1013+
assert waited == []
1014+
1015+
def test_close_cuts_the_wait_short(self, clock):
1016+
sender = FakeSender(SendOutcome("retry-later"))
1017+
pipeline, _, _ = make_traces(sender=sender)
1018+
exporter = pipeline._exporter
1019+
pipeline.start_span("a").end()
1020+
1021+
def close_during_wait(seconds):
1022+
exporter.close()
1023+
return True
1024+
1025+
exporter._wait_for_retry = close_during_wait
1026+
pipeline.flush(timeout=30)
1027+
assert len(sender.payloads) == 1
1028+
1029+
def test_the_real_wait_returns_when_close_is_called(self):
1030+
pipeline, _, _ = make_traces()
1031+
exporter = pipeline._exporter
1032+
threading.Thread(target=lambda: (time.sleep(0.05), exporter.close())).start()
1033+
started = time.monotonic()
1034+
assert exporter._wait_for_retry(5) is True
1035+
assert time.monotonic() - started < 2

‎posthog/tracing/_export.py‎

Lines changed: 72 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@
2020
build_otlp_span,
2121
build_resource_attributes,
2222
build_traces_payload,
23+
to_resource_key_value_list,
2324
)
2425
from ._transport import SendOutcome, send_traces_batch
2526

@@ -101,20 +102,28 @@ def __init__(
101102
self._config = config
102103
self._drops = drops
103104
self._send = send
104-
self._resource_attributes = build_resource_attributes(
105-
config.service_name,
106-
config.service_version,
107-
config.environment,
108-
config.resource_attributes,
105+
# Encoded once: the resource is the same for every batch.
106+
self._resource = to_resource_key_value_list(
107+
build_resource_attributes(
108+
config.service_name,
109+
config.service_version,
110+
config.environment,
111+
config.resource_attributes,
112+
)
109113
)
110114

111115
self._lock = threading.Lock()
112116
self._flush_lock = threading.Lock()
117+
# Set by close(), so a flush waiting out a backoff returns at once.
118+
self._closing = threading.Event()
113119
self._closed = False
114120
self._queue: List[SpanRecord] = []
115121
self._flush_timer: Optional[threading.Timer] = None
116122
self._flush_timer_fires_at = 0.0
117123
self._max_export_batch_size = config.max_export_batch_size
124+
# The size before a server 413 halved it, restored once the oversized
125+
# span is isolated and dropped.
126+
self._batch_size_before_halving: Optional[int] = None
118127
self._consecutive_failures = 0
119128
# Drawn once per failure, so the timer and the retry-budget charge see
120129
# the same delay.
@@ -151,15 +160,24 @@ def flush(
151160
) -> None:
152161
"""Drain the queue: one pass over what was queued, then one follow-up pass.
153162
154-
With a ``timeout``, no request starts once it is spent, except the
155-
first. A request already in flight is bounded by the client's timeout.
163+
With a ``timeout``, the budget starts once no other flush is in flight;
164+
one still in flight after that long a wait is left to finish and
165+
nothing is sent here. No request starts once the budget is spent,
166+
except the first, and a retriable failure is retried after its backoff
167+
while budget remains, once more at the deadline. Without a timeout a
168+
retriable failure is left to the timer. A request already in flight is
169+
bounded by the client's timeout.
156170
"""
157-
deadline = None if timeout is None else time.monotonic() + timeout
158171
# -1 is Lock.acquire's unbounded form.
159172
if not self._flush_lock.acquire(
160173
timeout=-1 if timeout is None else max(0.0, timeout)
161174
):
175+
log.debug(
176+
"Skipping a span flush: another flush was still in flight after %ss",
177+
timeout,
178+
)
162179
return
180+
deadline = None if timeout is None else time.monotonic() + timeout
163181
try:
164182
with self._lock:
165183
# A timer that waited behind another flush may have been
@@ -168,25 +186,50 @@ def flush(
168186
return
169187
self._clear_timer_locked()
170188
try:
171-
removed, stop = self._drain(deadline)
172-
if (
173-
removed
174-
and not stop
175-
and self._queue
176-
and (deadline is None or time.monotonic() < deadline)
177-
):
178-
self._drain(deadline)
189+
self._drain_within_budget(deadline, retry=_timer is None)
179190
finally:
180191
with self._lock:
181192
self._rearm_after_pass_locked()
182193
finally:
183194
self._flush_lock.release()
184195
self._drops.warn_if_due(force=True)
185196

197+
def _drain_within_budget(self, deadline: Optional[float], retry: bool) -> None:
198+
while True:
199+
removed, stop = self._drain(deadline)
200+
if not stop:
201+
if (
202+
removed
203+
and self._queue
204+
and (deadline is None or time.monotonic() < deadline)
205+
):
206+
self._drain(deadline)
207+
return
208+
if deadline is None or not retry:
209+
return
210+
wait = self._retry_wait_locked()
211+
remaining = deadline - time.monotonic()
212+
if wait is None or remaining <= 0:
213+
return
214+
if self._wait_for_retry(min(wait, remaining)):
215+
return
216+
217+
def _retry_wait_locked(self) -> Optional[float]:
218+
"""The backoff to wait out before retrying, or ``None`` when there is nothing to retry."""
219+
with self._lock:
220+
if self._closed or not self._queue or not self._consecutive_failures:
221+
return None
222+
return self._next_flush_delay_locked()
223+
224+
def _wait_for_retry(self, seconds: float) -> bool:
225+
"""Wait out a backoff; ``True`` when close() cut the wait short."""
226+
return self._closing.wait(seconds)
227+
186228
def close(self) -> None:
187229
"""Stop exporting and discard what is still queued. Called at shutdown."""
188230
with self._lock:
189231
self._closed = True
232+
self._closing.set()
190233
self._clear_timer_locked()
191234
discarded = len(self._queue)
192235
self._queue = []
@@ -212,10 +255,12 @@ def reinit_after_fork(self) -> None:
212255
# in the child, so they are replaced rather than acquired.
213256
self._lock = threading.Lock()
214257
self._flush_lock = threading.Lock()
258+
self._closing = threading.Event()
215259
self._flush_timer = None
216260
self._flush_timer_fires_at = 0.0
217261
self._queue = []
218262
self._max_export_batch_size = self._config.max_export_batch_size
263+
self._batch_size_before_halving = None
219264
self._retry_after.reset()
220265
self._end_failure_sequence_locked()
221266

@@ -281,7 +326,7 @@ def _drain(self, deadline: Optional[float]) -> Tuple[int, bool]:
281326
sent_any = True
282327
try:
283328
outcome = self._send(
284-
self._client, build_traces_payload(spans, self._resource_attributes)
329+
self._client, build_traces_payload(spans, self._resource)
285330
)
286331
except Exception:
287332
log.debug("Span batch send failed", exc_info=True)
@@ -312,20 +357,29 @@ def _apply_outcome_locked(
312357
if outcome.kind == "ok":
313358
del self._queue[:size]
314359
self._end_failure_sequence_locked()
360+
self._batch_size_before_halving = None
315361
if self._max_export_batch_size < self._config.max_export_batch_size:
316-
self._max_export_batch_size += 1
362+
self._max_export_batch_size = min(
363+
self._config.max_export_batch_size, self._max_export_batch_size * 2
364+
)
317365
return size, False, None
318366

319367
if outcome.kind == "too-large":
320368
if size == 1:
321369
del self._queue[:1]
322370
self._end_failure_sequence_locked()
323371
self._drops.record(1, "it is too large for the ingestion endpoint")
372+
# The oversized span is gone; the batches after it are not suspect.
373+
if self._batch_size_before_halving is not None:
374+
self._max_export_batch_size = self._batch_size_before_halving
375+
self._batch_size_before_halving = None
324376
return 1, False, None
325377
# Halve the refused batch, not the configured size: a shallow queue
326378
# would otherwise resend the same body.
327379
halved = max(1, size // 2)
328380
if not outcome.measured_locally:
381+
if self._batch_size_before_halving is None:
382+
self._batch_size_before_halving = self._max_export_batch_size
329383
self._max_export_batch_size = halved
330384
self._reset_head_batch_budget_locked()
331385
return (

0 commit comments

Comments
 (0)