Skip to content

Commit 2c6f3bb

Browse files
authored
Merge branch 'main' into feat/go-ci-visibility-failnow-teardown
2 parents 12fa701 + e2961fd commit 2c6f3bb

6 files changed

Lines changed: 278 additions & 35 deletions

File tree

contrib/mark3labs/mcp-go/go.mod

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -31,7 +31,7 @@ require (
3131
github.com/DataDog/sketches-go v1.4.8 // indirect
3232
github.com/Microsoft/go-winio v0.6.2 // indirect
3333
github.com/bahlo/generic-list-go v0.2.0 // indirect
34-
github.com/buger/jsonparser v1.1.1 // indirect
34+
github.com/buger/jsonparser v1.1.2 // indirect
3535
github.com/cenkalti/backoff/v5 v5.0.3 // indirect
3636
github.com/cespare/xxhash/v2 v2.3.0 // indirect
3737
github.com/cihub/seelog v0.0.0-20170130134532-f561c5e57575 // indirect

contrib/mark3labs/mcp-go/go.sum

Lines changed: 2 additions & 2 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

ddtrace/tracer/textmap.go

Lines changed: 21 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -1441,27 +1441,34 @@ func (*propagatorBaggage) extractTextMap(reader TextMapReader) (*SpanContext, er
14411441
return &ctx, nil
14421442
}
14431443

1444-
parts := strings.Split(baggageHeader, ",")
1445-
1446-
// 1) validation & single-trim pass
1447-
for i, kv := range parts {
1444+
// Single pass: enforce baggageMaxItems and baggageMaxBytes, validate, and apply.
1445+
ctr := 0
1446+
byteCount := 0
1447+
for kv := range strings.SplitSeq(baggageHeader, ",") {
1448+
itemBytes := len(kv)
1449+
if ctr > 0 {
1450+
itemBytes++ // comma separator
1451+
}
1452+
if ctr >= baggageMaxItems {
1453+
log.Warn("baggage item count exceeded limit (%d), dropping remaining items", baggageMaxItems)
1454+
break
1455+
}
1456+
if byteCount+itemBytes > baggageMaxBytes {
1457+
log.Warn("baggage byte limit exceeded (%d), dropping remaining items", baggageMaxBytes)
1458+
break
1459+
}
14481460
k, v, ok := strings.Cut(kv, "=")
14491461
trimmedK := strings.TrimSpace(k)
14501462
trimmedV := strings.TrimSpace(v)
14511463
if !ok || trimmedK == "" || trimmedV == "" {
14521464
log.Warn("invalid baggage item: %q, dropping entire header", kv)
1453-
return &ctx, nil
1465+
return &SpanContext{}, nil
14541466
}
1455-
// store back the trimmed pair so we don't re-trim below
1456-
parts[i] = trimmedK + "=" + trimmedV
1457-
}
1458-
1459-
// 2) safe to URL-decode & apply
1460-
for _, kv := range parts {
1461-
rawK, rawV, _ := strings.Cut(kv, "=")
1462-
key, _ := url.QueryUnescape(rawK)
1463-
val, _ := url.QueryUnescape(rawV)
1467+
key, _ := url.QueryUnescape(trimmedK)
1468+
val, _ := url.QueryUnescape(trimmedV)
14641469
ctx.setBaggageItem(key, val)
1470+
byteCount += itemBytes
1471+
ctr++
14651472
}
14661473

14671474
return &ctx, nil

ddtrace/tracer/textmap_test.go

Lines changed: 121 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2892,6 +2892,127 @@ func TestExtractBaggagePropagatorMalformedHeader(t *testing.T) {
28922892
})
28932893
}
28942894

2895+
func TestExtractBaggagePropagatorMaxItems(t *testing.T) {
2896+
tracer, err := newTracer()
2897+
assert.NoError(t, err)
2898+
defer tracer.Stop()
2899+
2900+
var b strings.Builder
2901+
for i := range baggageMaxItems + 5 {
2902+
if i > 0 {
2903+
b.WriteByte(',')
2904+
}
2905+
iStr := strconv.Itoa(i)
2906+
b.WriteString("key" + iStr + "=val" + iStr)
2907+
}
2908+
2909+
headers := TextMapCarrier{
2910+
DefaultTraceIDHeader: "4",
2911+
DefaultParentIDHeader: "1",
2912+
DefaultBaggageHeader: b.String(),
2913+
}
2914+
s, err := tracer.Extract(headers)
2915+
assert.NoError(t, err)
2916+
2917+
got := make(map[string]string)
2918+
s.ForeachBaggageItem(func(k, v string) bool {
2919+
got[k] = v
2920+
return true
2921+
})
2922+
assert.Len(t, got, baggageMaxItems)
2923+
for i := range baggageMaxItems {
2924+
iStr := strconv.Itoa(i)
2925+
assert.Equal(t, "val"+iStr, got["key"+iStr])
2926+
}
2927+
for i := baggageMaxItems; i < baggageMaxItems+5; i++ {
2928+
iStr := strconv.Itoa(i)
2929+
_, present := got["key"+iStr]
2930+
assert.False(t, present, "key%s should not be present", iStr)
2931+
}
2932+
}
2933+
2934+
func TestExtractBaggagePropagatorMaxBytes(t *testing.T) {
2935+
tracer, err := newTracer()
2936+
assert.NoError(t, err)
2937+
defer tracer.Stop()
2938+
2939+
// 12 items, each "keyN=" + 1000 'a's = 1005 wire bytes. Including comma
2940+
// separators, the first 8 fit under baggageMaxBytes (8192); the 9th would
2941+
// push the running total to 9053 > 8192, so items 8..11 are dropped.
2942+
const itemValLen = 1000
2943+
const numItems = 12
2944+
const expectedKept = 8
2945+
val := strings.Repeat("a", itemValLen)
2946+
2947+
var b strings.Builder
2948+
for i := range numItems {
2949+
if i > 0 {
2950+
b.WriteByte(',')
2951+
}
2952+
b.WriteString("key" + strconv.Itoa(i) + "=" + val)
2953+
}
2954+
2955+
headers := TextMapCarrier{
2956+
DefaultTraceIDHeader: "4",
2957+
DefaultParentIDHeader: "1",
2958+
DefaultBaggageHeader: b.String(),
2959+
}
2960+
s, err := tracer.Extract(headers)
2961+
assert.NoError(t, err)
2962+
2963+
got := make(map[string]string)
2964+
s.ForeachBaggageItem(func(k, v string) bool {
2965+
got[k] = v
2966+
return true
2967+
})
2968+
assert.Len(t, got, expectedKept)
2969+
for i := range expectedKept {
2970+
assert.Equal(t, val, got["key"+strconv.Itoa(i)])
2971+
}
2972+
for i := expectedKept; i < numItems; i++ {
2973+
_, present := got["key"+strconv.Itoa(i)]
2974+
assert.False(t, present, "key%d should not be present", i)
2975+
}
2976+
}
2977+
2978+
func TestExtractBaggagePropagatorMalformedPastLimit(t *testing.T) {
2979+
tracer, err := newTracer()
2980+
assert.NoError(t, err)
2981+
defer tracer.Stop()
2982+
2983+
// baggageMaxItems valid entries followed by a malformed entry. Because
2984+
// the malformed entry sits past the items limit it is never inspected,
2985+
// so the valid prefix is kept (regression check on the single-pass design).
2986+
var b strings.Builder
2987+
for i := range baggageMaxItems {
2988+
if i > 0 {
2989+
b.WriteByte(',')
2990+
}
2991+
iStr := strconv.Itoa(i)
2992+
b.WriteString("key" + iStr + "=val" + iStr)
2993+
}
2994+
b.WriteString(",malformed_no_equals_sign")
2995+
2996+
headers := TextMapCarrier{
2997+
DefaultTraceIDHeader: "4",
2998+
DefaultParentIDHeader: "1",
2999+
DefaultBaggageHeader: b.String(),
3000+
}
3001+
s, err := tracer.Extract(headers)
3002+
assert.NoError(t, err)
3003+
3004+
got := make(map[string]string)
3005+
s.ForeachBaggageItem(func(k, v string) bool {
3006+
got[k] = v
3007+
return true
3008+
})
3009+
assert.Len(t, got, baggageMaxItems)
3010+
for i := range baggageMaxItems {
3011+
iStr := strconv.Itoa(i)
3012+
assert.Equal(t, "val"+iStr, got["key"+iStr])
3013+
}
3014+
}
3015+
28953016
func TestExtractOnlyBaggage(t *testing.T) {
28963017
t.Setenv("DD_TRACE_PROPAGATION_STYLE", "baggage")
28973018
headers := TextMapCarrier(map[string]string{

internal/llmobs/llmobs.go

Lines changed: 47 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -146,9 +146,10 @@ type LLMObs struct {
146146
// lifecycle
147147
mu sync.Mutex
148148
running bool
149-
wg sync.WaitGroup
150-
stopCh chan struct{} // signal stop
151-
flushNowCh chan struct{}
149+
sendWg sync.WaitGroup // tracks in-flight batchSend goroutines
150+
workerDone chan struct{} // closed when the worker loop exits
151+
stopCh chan struct{} // signal stop
152+
flushNowCh chan chan struct{}
152153
flushInterval time.Duration
153154
}
154155

@@ -187,8 +188,9 @@ func newLLMObs(cfg *config.Config, tracer Tracer) (*LLMObs, error) {
187188
Tracer: tracer,
188189
spanEventsCh: make(chan *transport.LLMObsSpanEvent),
189190
evalMetricsCh: make(chan *transport.LLMObsMetric),
191+
workerDone: make(chan struct{}),
190192
stopCh: make(chan struct{}),
191-
flushNowCh: make(chan struct{}, 1),
193+
flushNowCh: make(chan chan struct{}, 1),
192194
flushInterval: defaultFlushInterval,
193195
}, nil
194196
}
@@ -245,6 +247,13 @@ func Flush() {
245247
}
246248
}
247249

250+
// FlushSync forces a flush of all buffered LLMObs data and blocks until the flush completes.
251+
func FlushSync() {
252+
if activeLLMObs != nil {
253+
activeLLMObs.FlushSync()
254+
}
255+
}
256+
248257
// Run starts the worker loop that processes span events and metrics.
249258
func (l *LLMObs) Run() {
250259
l.mu.Lock()
@@ -255,7 +264,8 @@ func (l *LLMObs) Run() {
255264
l.running = true
256265
l.mu.Unlock()
257266

258-
l.wg.Go(func() {
267+
go func() {
268+
defer close(l.workerDone)
259269
// this goroutine should be the only one writing to the internal buffers
260270

261271
ticker := time.NewTicker(l.flushInterval)
@@ -268,9 +278,7 @@ func (l *LLMObs) Run() {
268278
if l.bufSpanEventsSize+evSize > sizeLimitEVPEvent {
269279
log.Debug("llmobs: span events buffer size limit reached, flushing before adding new event")
270280
params := l.clearBuffersNonLocked()
271-
l.wg.Go(func() {
272-
l.batchSend(params)
273-
})
281+
l.sendWg.Go(func() { l.batchSend(params) })
274282
}
275283
l.bufSpanEvents = append(l.bufSpanEvents, ev)
276284
l.bufSpanEventsSize += evSize
@@ -280,16 +288,22 @@ func (l *LLMObs) Run() {
280288

281289
case <-ticker.C:
282290
params := l.clearBuffersNonLocked()
283-
l.wg.Go(func() {
284-
l.batchSend(params)
285-
})
291+
l.sendWg.Go(func() { l.batchSend(params) })
286292

287-
case <-l.flushNowCh:
293+
case done := <-l.flushNowCh:
288294
log.Debug("llmobs: on-demand flush signal")
289295
params := l.clearBuffersNonLocked()
290-
l.wg.Go(func() {
296+
l.sendWg.Add(1)
297+
go func() {
298+
defer func() {
299+
l.sendWg.Done()
300+
if done != nil {
301+
l.sendWg.Wait()
302+
close(done)
303+
}
304+
}()
291305
l.batchSend(params)
292-
})
306+
}()
293307

294308
case <-l.stopCh:
295309
log.Debug("llmobs: stop signal")
@@ -299,7 +313,7 @@ func (l *LLMObs) Run() {
299313
return
300314
}
301315
}
302-
})
316+
}()
303317
}
304318

305319
// clearBuffersNonLocked clears the internal buffers and returns the corresponding batchSendParams to send to the backend.
@@ -320,11 +334,24 @@ func (l *LLMObs) clearBuffersNonLocked() batchSendParams {
320334
func (l *LLMObs) Flush() {
321335
// non-blocking edge trigger so multiple calls coalesce
322336
select {
323-
case l.flushNowCh <- struct{}{}:
337+
case l.flushNowCh <- nil:
324338
default:
325339
}
326340
}
327341

342+
// FlushSync forces an immediate flush and blocks until the flush completes.
343+
func (l *LLMObs) FlushSync() {
344+
done := make(chan struct{})
345+
select {
346+
case l.flushNowCh <- done:
347+
select {
348+
case <-done:
349+
case <-l.stopCh:
350+
}
351+
case <-l.stopCh:
352+
}
353+
}
354+
328355
// Stop requests shutdown, drains what’s already in the channels, flushes, and waits.
329356
func (l *LLMObs) Stop() {
330357
l.mu.Lock()
@@ -342,8 +369,10 @@ func (l *LLMObs) Stop() {
342369
close(l.stopCh)
343370
}
344371

345-
// Wait for the main worker to exit (it will do a final flush)
346-
l.wg.Wait()
372+
// Wait for the worker loop to exit (it does a final synchronous flush),
373+
// then wait for any async batchSend goroutines still in flight.
374+
<-l.workerDone
375+
l.sendWg.Wait()
347376
}
348377

349378
// drainChannels pulls everything currently buffered in the channels into our in-memory buffers.

0 commit comments

Comments
 (0)