Skip to content

Commit a1d8e27

Browse files
Break Lapdog payloads into chunks of 20 spans each (#404)
* Make Lapdog respect DD_TAGS * Replace DD_TAGS with session-scoped Lapdog tags Allow Claude to tag its active Lapdog session through an explicit CLI command, avoiding detached-process environment limitations while updating existing and future spans. Co-authored-by: Cursor <cursoragent@cursor.com> * Add a release note and extend the change to Codex * Correctly handle parallel Codex sessions * Also add support for Pi and address remaining inline comments * Add one more comment * Make Lapdog split large traces into chunks of at most 20 spans before sending * Add a release note --------- Co-authored-by: Cursor <cursoragent@cursor.com>
1 parent 1fbb60e commit a1d8e27

3 files changed

Lines changed: 99 additions & 15 deletions

File tree

ddapm_test_agent/claude_hooks.py

Lines changed: 37 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -71,6 +71,7 @@
7171
"user_name",
7272
}
7373
)
74+
_MAX_SPANS_PER_BACKEND_REQUEST = 20
7475

7576
# Models with 1M token context windows (native, no beta header needed).
7677
# All other models default to 200k.
@@ -1768,6 +1769,36 @@ async def _post_to_backend(self, url: str, headers: Dict[str, str], data: bytes,
17681769
except Exception as e:
17691770
log.warning("Error trying to %s: %s", description, e)
17701771

1772+
async def _post_spans_to_backend(
1773+
self,
1774+
url: str,
1775+
headers: Dict[str, str],
1776+
spans: List[Dict[str, Any]],
1777+
description: str,
1778+
) -> None:
1779+
"""POST spans in ordered chunks that stay within the backend request limit."""
1780+
chunk_count = (len(spans) + _MAX_SPANS_PER_BACKEND_REQUEST - 1) // _MAX_SPANS_PER_BACKEND_REQUEST
1781+
if chunk_count > 1:
1782+
log.info(
1783+
"Splitting %d spans into %d backend requests of at most %d spans",
1784+
len(spans),
1785+
chunk_count,
1786+
_MAX_SPANS_PER_BACKEND_REQUEST,
1787+
)
1788+
1789+
for chunk_index, start in enumerate(range(0, len(spans), _MAX_SPANS_PER_BACKEND_REQUEST), start=1):
1790+
chunk = spans[start : start + _MAX_SPANS_PER_BACKEND_REQUEST]
1791+
payload = {
1792+
"_dd.stage": "raw",
1793+
"event_type": "span",
1794+
"spans": chunk,
1795+
}
1796+
data = gzip.compress(msgpack.packb(payload))
1797+
chunk_description = description
1798+
if chunk_count > 1:
1799+
chunk_description = f"{description} (chunk {chunk_index}/{chunk_count}, {len(chunk)} spans)"
1800+
await self._post_to_backend(url, headers, data, chunk_description)
1801+
17711802
async def _forward_span_update(self, spans: List[Dict[str, Any]]) -> None:
17721803
"""Forward span updates to the DD backend via the update endpoint."""
17731804
if not spans:
@@ -1778,13 +1809,7 @@ async def _forward_span_update(self, spans: List[Dict[str, Any]]) -> None:
17781809
return
17791810
url, headers = target
17801811

1781-
payload = {
1782-
"_dd.stage": "raw",
1783-
"event_type": "span",
1784-
"spans": spans,
1785-
}
1786-
data = gzip.compress(msgpack.packb(payload))
1787-
await self._post_to_backend(url, headers, data, f"forward {len(spans)} span updates")
1812+
await self._post_spans_to_backend(url, headers, spans, f"forward {len(spans)} span updates")
17881813

17891814
async def _forward_trace_to_backend(
17901815
self, session_id: str, trace_id: Optional[str] = None, span_source: str = "Claude hooks"
@@ -1817,14 +1842,11 @@ async def _forward_trace_to_backend(
18171842
span["tags"] = tags + ["lapdog_forwarded:true"]
18181843
forwarded_spans.append(span)
18191844

1820-
payload = {
1821-
"_dd.stage": "raw",
1822-
"event_type": "span",
1823-
"spans": forwarded_spans,
1824-
}
1825-
data = gzip.compress(msgpack.packb(payload))
1826-
await self._post_to_backend(
1827-
url, headers, data, f"forward {len(spans)} {span_source} spans for trace {trace_id}"
1845+
await self._post_spans_to_backend(
1846+
url,
1847+
headers,
1848+
forwarded_spans,
1849+
f"forward {len(spans)} {span_source} spans for trace {trace_id}",
18281850
)
18291851

18301852
async def _forward_eval_metrics_to_backend(self, session_id: str, trace_id: Optional[str] = None) -> None:
Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
1+
---
2+
fixes:
3+
- |
4+
When forwarding payloads to the Datadog backend, send them in chunks
5+
of at most 20 spans each to get around payload size limits on the backend.

tests/test_claude_hooks.py

Lines changed: 57 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,11 @@
1+
import gzip
12
import json
23
import os
34
import subprocess
45
import tempfile
56

7+
import msgpack
8+
69
from ddapm_test_agent.claude_hooks import ClaudeHooksAPI
710
from ddapm_test_agent.claude_link_tracker import ClaudeLinkTracker
811
from ddapm_test_agent.claude_proxy import ClaudeProxyAPI
@@ -37,6 +40,60 @@ async def test_hook_missing_session_id(agent):
3740
assert "session_id" in body["error"]
3841

3942

43+
async def test_trace_forwarding_chunks_payloads_at_twenty_spans(monkeypatch):
44+
cases = [
45+
(20, [20]),
46+
(21, [20, 1]),
47+
(45, [20, 20, 5]),
48+
]
49+
for span_count, expected_chunk_sizes in cases:
50+
hooks = ClaudeHooksAPI()
51+
session = hooks._get_or_create_session(f"chunk-session-{span_count}")
52+
trace_id = f"chunk-trace-{span_count}"
53+
session.trace_id = trace_id
54+
hooks._assembled_spans = [
55+
{
56+
"span_id": str(index),
57+
"trace_id": trace_id,
58+
"tags": [],
59+
}
60+
for index in range(span_count)
61+
]
62+
posted_payloads = []
63+
descriptions = []
64+
65+
monkeypatch.setattr(
66+
hooks,
67+
"_resolve_backend_target",
68+
lambda *args, **kwargs: ("http://backend.example/api/v2/llmobs", {"Content-Type": "application/msgpack"}),
69+
)
70+
71+
async def fake_post_to_backend(url, headers, data, description):
72+
assert url == "http://backend.example/api/v2/llmobs"
73+
assert headers["Content-Type"] == "application/msgpack"
74+
posted_payloads.append(msgpack.unpackb(gzip.decompress(data), raw=False))
75+
descriptions.append(description)
76+
77+
monkeypatch.setattr(hooks, "_post_to_backend", fake_post_to_backend)
78+
79+
await hooks._forward_trace_to_backend(session.session_id)
80+
81+
assert [len(payload["spans"]) for payload in posted_payloads] == expected_chunk_sizes
82+
assert all(payload["_dd.stage"] == "raw" for payload in posted_payloads)
83+
assert all(payload["event_type"] == "span" for payload in posted_payloads)
84+
forwarded_spans = [span for payload in posted_payloads for span in payload["spans"]]
85+
assert [span["span_id"] for span in forwarded_spans] == [str(index) for index in range(span_count)]
86+
assert all("lapdog_forwarded:true" in span["tags"] for span in forwarded_spans)
87+
assert all("lapdog_forwarded:true" not in span["tags"] for span in hooks._assembled_spans)
88+
if span_count <= 20:
89+
assert all("(chunk " not in description for description in descriptions)
90+
else:
91+
assert all(
92+
f"chunk {index + 1}/{len(expected_chunk_sizes)}" in description
93+
for index, description in enumerate(descriptions)
94+
)
95+
96+
4097
async def test_session_tags_apply_to_existing_and_future_spans(agent):
4198
session_id = "sess-custom-tags"
4299
session_token = "launch-token"

0 commit comments

Comments
 (0)