Skip to content

Commit 24391d1

Browse files
authored
JITSU-138: fix Datadog reserved-status collision; synthesize otlp message lines (#1475)
Follow-up to `JITSU-138`, from end-to-end testing against Datadog's agentless OTLP intake (which merges the structured body's top-level keys into root log attributes). **Architecture (reworked per review): the otlp destination stays fully payload-agnostic** — `otlp.go` is untouched relative to `newjitsu`. Body adaptation happens where envelopes are produced to the Kafka topic, and **only for bodies the pipeline constructs itself**: 1. **Go producer** (`eventslog/kafka_events_log.go`), for `bulker_batch` / `bulker_stream` records (bulker-built `bulker.State` / stream-status shapes): - top-level `status` → `record_status`: `COMPLETED`/`FAILED` values collide with Datadog's reserved status attribute and made **every** record render as `critical` regardless of severity (reproduced in isolation with controlled probe payloads) - synthesized `message` when absent/empty/null: `bulker_batch COMPLETED: 2 rows → events39`, with a rune-safe-truncated error preview on failures - the eventId hash is computed from the unadapted body — ids and billing dedup unchanged 2. **Rotor producer** (`kafka-events-store.ts`): dead-letter bodies (rotor-built `{payload, error}`) get their `message` line at construction. **Function-log bodies are user data and are never touched** — no key renames, no synthesized fields. Body stays structured (kvlist) — no stringification. Covered by `TestKafkaEventsLogAdaptsOwnedBodies` (rename, synthesis, truncation, function-body-untouched) and rotor `deadLetterMessage` tests (incl. surrogate-pair safety). Docs note: jitsucom/websites#60. 🤖 Generated with [Claude Code](https://claude.com/claude-code)
2 parents 793d02a + 0892e5b commit 24391d1

4 files changed

Lines changed: 161 additions & 2 deletions

File tree

‎bulker/eventslog/kafka_events_log.go‎

Lines changed: 77 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,11 +3,14 @@ package eventslog
33
import (
44
"crypto/sha256"
55
"encoding/hex"
6+
"encoding/json"
67
"fmt"
8+
"strings"
79
"time"
810

911
"github.com/jitsucom/bulker/jitsubase/jsonorder"
1012
"github.com/jitsucom/bulker/jitsubase/logging"
13+
"github.com/jitsucom/bulker/jitsubase/types"
1114
)
1215

1316
// Live Events observability export fan-out (JITSU-138): a KafkaEventsLogService
@@ -156,6 +159,15 @@ func (k *KafkaEventsLogService) PostAsync(event *ActorEvent) {
156159
envelope.ConnectionId = event.ActorId
157160
case EventTypeBatch, EventTypeProcessed:
158161
envelope.DestinationId = event.ActorId
162+
// bulker_batch / bulker_stream bodies are constructed by bulker itself
163+
// (bulker.State / stream status shapes), so we own their top-level keys
164+
// and adapt them for observability backends here, at the producer —
165+
// the otlp destination stays payload-agnostic. Function-log bodies are
166+
// user data and are never touched. The eventId hash above is computed
167+
// from the unadapted body, so the id is unaffected
168+
if adapted := adaptOwnedBody(body, string(event.EventType)); adapted != nil {
169+
envelope.Body = adapted
170+
}
159171
}
160172
payload, err := jsonorder.Marshal(&envelope)
161173
if err != nil {
@@ -174,6 +186,71 @@ func (k *KafkaEventsLogService) PostAsync(event *ActorEvent) {
174186
k.produceBillingRecord(&envelope, timestamp)
175187
}
176188

189+
// adaptOwnedBody adapts a bulker-owned record body for observability
190+
// backends: `status` is renamed to `record_status` (Datadog's agentless OTLP
191+
// intake merges top-level body keys into root log attributes, where `status`
192+
// is reserved for severity — COMPLETED/FAILED values made every record render
193+
// as `critical`), and a short human-readable `message` is synthesized when
194+
// absent or empty so log list views show a line instead of a blank. Returns
195+
// nil (caller keeps the original body) when the body can't be re-parsed
196+
func adaptOwnedBody(bodyJSON []byte, eventType string) types.Json {
197+
obj := types.NewJson(0)
198+
if err := jsonorder.Unmarshal(bodyJSON, &obj); err != nil || obj == nil {
199+
return nil
200+
}
201+
if v, ok := obj.Get("status"); ok {
202+
obj.Delete("status")
203+
obj.Set("record_status", v)
204+
}
205+
if v, ok := obj.Get("message"); !ok || v == "" || v == nil {
206+
obj.Set("message", stateMessage(eventType, obj))
207+
}
208+
return obj
209+
}
210+
211+
// stateMessage: e.g. "bulker_batch COMPLETED: 2 rows → events39"
212+
func stateMessage(eventType string, body types.Json) string {
213+
var b strings.Builder
214+
b.WriteString(eventType)
215+
if s := body.GetS("record_status"); s != "" {
216+
b.WriteString(" " + s)
217+
}
218+
if v, ok := body.Get("processedRows"); ok {
219+
if n, isNum := asInt64(v); isNum {
220+
fmt.Fprintf(&b, ": %d rows", n)
221+
}
222+
}
223+
if rep, ok := body.Get("representation"); ok {
224+
if repObj, isObj := rep.(types.Json); isObj {
225+
if name := repObj.GetS("name"); name != "" {
226+
b.WriteString(" → " + name)
227+
}
228+
}
229+
}
230+
if e := body.GetS("error"); e != "" {
231+
if r := []rune(e); len(r) > 140 {
232+
e = string(r[:140]) + "…"
233+
}
234+
b.WriteString(" — " + e)
235+
}
236+
return b.String()
237+
}
238+
239+
func asInt64(v any) (int64, bool) {
240+
switch n := v.(type) {
241+
case json.Number:
242+
i, err := n.Int64()
243+
return i, err == nil
244+
case float64:
245+
return int64(n), true
246+
case int64:
247+
return n, true
248+
case int:
249+
return int64(n), true
250+
}
251+
return 0, false
252+
}
253+
177254
// produceBillingRecord emits one active_incoming record per exported envelope.
178255
// The composed key {eventId}_0_{secondsWithinHour} with an hour-truncated
179256
// timestamp follows ingest's buildSyncMetrics; dedup happens downstream via

‎bulker/eventslog/kafka_events_log_test.go‎

Lines changed: 49 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@ package eventslog
33
import (
44
"encoding/json"
55
"fmt"
6+
"strings"
67
"testing"
78
"time"
89

@@ -203,6 +204,54 @@ func TestKafkaEventsLogOrderedMapBodyFidelity(t *testing.T) {
203204
require.Contains(t, payload, `"lastMappedRow":{"message_id":"m1","nested":{"a":1}}`)
204205
}
205206

207+
func TestKafkaEventsLogAdaptsOwnedBodies(t *testing.T) {
208+
service, produced := newTestKafkaService("", map[string]bool{"ws1": true})
209+
210+
// bulker_batch: status renamed (Datadog reserved-attribute collision) and a
211+
// message line synthesized from status/rows/target table
212+
batch := testActorEvent()
213+
batch.EventType = EventTypeBatch
214+
batch.ActorId = "dst1"
215+
rep := types.NewJson(1)
216+
rep.Set("name", "events39")
217+
body := types.NewJson(3)
218+
body.Set("status", "COMPLETED")
219+
body.Set("processedRows", 2)
220+
body.Set("representation", rep)
221+
batch.Event = body
222+
service.PostAsync(batch)
223+
require.Len(t, *produced, 2)
224+
envelope := map[string]any{}
225+
require.NoError(t, json.Unmarshal((*produced)[0].payload, &envelope))
226+
exported := envelope["body"].(map[string]any)
227+
require.NotContains(t, exported, "status")
228+
require.Equal(t, "COMPLETED", exported["record_status"])
229+
require.Equal(t, "bulker_batch COMPLETED: 2 rows → events39", exported["message"])
230+
231+
// failed stream: error preview appended, long errors truncated
232+
stream := testActorEvent()
233+
stream.EventType = EventTypeProcessed
234+
stream.ActorId = "dst1"
235+
stream.Event = map[string]any{"status": "FAILED", "error": strings.Repeat("x", 200)}
236+
service.PostAsync(stream)
237+
require.NoError(t, json.Unmarshal((*produced)[2].payload, &envelope))
238+
exported = envelope["body"].(map[string]any)
239+
msg := exported["message"].(string)
240+
require.True(t, strings.HasPrefix(msg, "bulker_stream FAILED — "))
241+
require.True(t, strings.HasSuffix(msg, "…"))
242+
require.Less(t, len(msg), 200)
243+
244+
// function logs are user data: status untouched, no message synthesized
245+
fn := testActorEvent()
246+
fn.Event = map[string]any{"type": "log-info", "status": "custom-user-value"}
247+
service.PostAsync(fn)
248+
require.NoError(t, json.Unmarshal((*produced)[4].payload, &envelope))
249+
exported = envelope["body"].(map[string]any)
250+
require.Equal(t, "custom-user-value", exported["status"])
251+
require.NotContains(t, exported, "record_status")
252+
require.NotContains(t, exported, "message")
253+
}
254+
206255
func TestKafkaEventsLogNoBillingOnEnqueueFailure(t *testing.T) {
207256
attempted := &[]string{}
208257
service := NewKafkaEventsLogService(KafkaEventsLogConfig{

‎services/rotor/__tests__/kafka-events-store.test.ts‎

Lines changed: 21 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,25 @@
11
import { describe, expect, it } from "vitest";
2-
import { exportEventId } from "../src/lib/kafka-events-store";
2+
import { deadLetterMessage, exportEventId } from "../src/lib/kafka-events-store";
3+
4+
describe("dead-letter message synthesis", () => {
5+
it("uses string errors and Error-like objects, falls back to bare label", () => {
6+
expect(deadLetterMessage("boom")).toBe("dead-letter — boom");
7+
expect(deadLetterMessage(new Error("kaput"))).toBe("dead-letter — kaput");
8+
expect(deadLetterMessage({ message: "obj msg" })).toBe("dead-letter — obj msg");
9+
expect(deadLetterMessage(null)).toBe("dead-letter");
10+
expect(deadLetterMessage({ code: 42 })).toBe("dead-letter");
11+
});
12+
13+
it("truncates long errors rune-safely", () => {
14+
const msg = deadLetterMessage("x".repeat(200));
15+
expect(msg).toMatch(/…$/);
16+
expect(Array.from(msg).length).toBeLessThan(200 - 140 + 160);
17+
// astral characters must not be split into surrogate halves
18+
const emoji = deadLetterMessage("💥".repeat(150));
19+
expect(emoji.includes("�")).toBe(false);
20+
expect(Array.from(emoji).slice(-2).join("")).toBe("💥…");
21+
});
22+
});
323

424
describe("export eventId content hash", () => {
525
it("matches the pinned cross-language vector (same recipe as Go eventslog.ExportEventId)", () => {

‎services/rotor/src/lib/kafka-events-store.ts‎

Lines changed: 14 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -54,6 +54,15 @@ type ExportRecord = {
5454
body: Record<string, any>;
5555
};
5656

57+
// mirrors the Go producer's stateMessage: short human line for backends that
58+
// surface a message field, with a rune-safe 140-char error preview
59+
export function deadLetterMessage(error: any): string {
60+
const errText = typeof error === "string" ? error : error && typeof error.message === "string" ? error.message : "";
61+
const runes = Array.from(errText);
62+
const preview = runes.length > 140 ? runes.slice(0, 140).join("") + "…" : errText;
63+
return preview ? `dead-letter — ${preview}` : "dead-letter";
64+
}
65+
5766
export function createKafkaEventsStore(): EventsStore {
5867
const prefix = serverEnv.KAFKA_TOPIC_PREFIX || "";
5968
// unset = billing/metrics emission disabled for this deployment, matching
@@ -189,7 +198,11 @@ export function createKafkaEventsStore(): EventsStore {
189198
level: "error",
190199
actorId: connectionId,
191200
connectionId,
192-
body: { payload, error },
201+
// dead-letter bodies are rotor-owned, so the human-readable message
202+
// line for observability backends is synthesized here at the producer —
203+
// the otlp destination stays payload-agnostic. Function-log bodies
204+
// already carry their own message and are never touched
205+
body: { payload, error, message: deadLetterMessage(error) },
193206
});
194207
},
195208
close() {

0 commit comments

Comments
 (0)