Skip to content

Commit e82eca1

Browse files
carderneTrigger.dev RepoOps
authored andcommitted
fix(webapp): preserve log search batches with invalid JSON escapes
Prevent log search batches from being lost when a message contains an invalid JSON escape. Search inserts now use the same lenient escape setting as source inserts and share the existing bounded sanitize/row-skip recovery, so an unrepairable row does not discard its neighbors. Drop counters start at zero for each reason, and recovered batches count only skipped rows as dropped. Mono-RevId: 353071db24c77a392baaa7cb96340f6ab64d8e3e
1 parent 84dd4a4 commit e82eca1

6 files changed

Lines changed: 204 additions & 26 deletions

File tree

‎apps/webapp/app/v3/eventRepository/clickhouseEventRepository.server.ts‎

Lines changed: 50 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -71,6 +71,7 @@ import type {
7171
TraceEventOptions,
7272
TraceSummary,
7373
} from "./eventRepository.types";
74+
import { insertLogsSearchRows } from "./insertLogsSearchRows.server";
7475
import {
7576
insertWithBadRowSkip,
7677
type JsonParseRecoveryOutcome,
@@ -233,6 +234,19 @@ export class ClickhouseEventRepository implements IEventRepository {
233234
"logs_search.dual_write.rows_dropped",
234235
{ unit: "rows" }
235236
);
237+
for (const reason of [
238+
"limiter_full",
239+
"shutdown",
240+
"source_recovered",
241+
"mapping_failed",
242+
"insert_failed",
243+
"parse_failed",
244+
]) {
245+
this._logsSearchRowsDroppedCounter.add(0, {
246+
...this._logsSearchMetricAttributes,
247+
reason,
248+
});
249+
}
236250
this._logsSearchBatchesCounter = meter.createCounter("logs_search.dual_write.batches", {
237251
unit: "batches",
238252
});
@@ -575,34 +589,51 @@ export class ClickhouseEventRepository implements IEventRepository {
575589

576590
async #insertLogsSearchRows(flushId: string, rows: TaskEventSearchV2Input[]): Promise<void> {
577591
const startedAt = Date.now();
578-
let lastError: { clickhouseErrorType?: string } | undefined;
592+
let lastError: unknown;
579593

580594
for (let attempt = 1; attempt <= 2; attempt++) {
581-
const [error] = await this._logsSearchClickhouse.taskEventsSearch.insert(rows, {
582-
params: {
583-
clickhouse_settings: {
584-
async_insert: 0,
585-
insert_deduplication_token: flushId,
586-
},
587-
},
588-
});
589-
590-
if (!error) {
591-
this._logsSearchRowsLandedCounter.add(rows.length, this._logsSearchMetricAttributes);
595+
try {
596+
const outcome = await insertLogsSearchRows(
597+
this._logsSearchClickhouse.taskEventsSearch.insert,
598+
flushId,
599+
rows,
600+
logger
601+
);
602+
const dropped = outcome.kind === "recovered" ? outcome.rowsDropped : 0;
603+
if (dropped > 0) {
604+
this._logsSearchRowsDroppedCounter.add(dropped, {
605+
...this._logsSearchMetricAttributes,
606+
reason: "parse_failed",
607+
});
608+
}
609+
if (outcome.kind !== "recovered" || outcome.rowsDroppedExact) {
610+
this._logsSearchRowsLandedCounter.add(
611+
rows.length - dropped,
612+
this._logsSearchMetricAttributes
613+
);
614+
}
592615
this._logsSearchBatchesCounter.add(1, {
593616
...this._logsSearchMetricAttributes,
594-
outcome: attempt === 1 ? "ok" : "retried_ok",
617+
outcome: landedNothing(outcome, rows.length)
618+
? "failed"
619+
: outcome.kind === "recovered"
620+
? outcome.rowsDroppedExact
621+
? "recovered"
622+
: "recovered_unknown"
623+
: attempt === 1 && outcome.kind === "inserted"
624+
? "ok"
625+
: "retried_ok",
595626
});
596627
this._logsSearchFlushDurationHistogram.record(
597628
Date.now() - startedAt,
598629
this._logsSearchMetricAttributes
599630
);
600631
return;
601-
}
602-
603-
lastError = error;
604-
if (attempt === 1) {
605-
await new Promise((resolve) => setTimeout(resolve, 500));
632+
} catch (error) {
633+
lastError = error;
634+
if (attempt === 1) {
635+
await new Promise((resolve) => setTimeout(resolve, 500));
636+
}
606637
}
607638
}
608639

@@ -621,7 +652,7 @@ export class ClickhouseEventRepository implements IEventRepository {
621652
logger.error("Logs search dual-write insert failed", {
622653
flushId,
623654
rows: rows.length,
624-
clickhouseErrorType: lastError?.clickhouseErrorType,
655+
error: lastError,
625656
});
626657
}
627658

Lines changed: 48 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,48 @@
1+
import type { ClickHouse, ClickHouseSettings, TaskEventSearchV2Input } from "@internal/clickhouse";
2+
import {
3+
insertWithBadRowSkip,
4+
isClickHouseJsonParseError,
5+
type JsonParseRecoveryLogger,
6+
} from "./sanitizeRowsOnParseError.server";
7+
8+
export function insertLogsSearchRows(
9+
insertRows: ClickHouse["taskEventsSearch"]["insert"],
10+
flushId: string,
11+
rows: TaskEventSearchV2Input[],
12+
logger: JsonParseRecoveryLogger
13+
) {
14+
const insert = async (batch: TaskEventSearchV2Input[], settings?: ClickHouseSettings) => {
15+
const [error, result] = await insertRows(batch, {
16+
params: {
17+
clickhouse_settings: {
18+
async_insert: 0,
19+
insert_deduplication_token: flushId,
20+
...settings,
21+
},
22+
},
23+
});
24+
if (error) throw error;
25+
return result;
26+
};
27+
28+
return insertWithBadRowSkip({
29+
rows,
30+
contextLabel: "task_events_search_v2",
31+
logger,
32+
logContext: { flushId },
33+
hasMaterializedViews: false,
34+
isParseError: (error) =>
35+
isClickHouseJsonParseError(error) ||
36+
(typeof error === "object" &&
37+
error !== null &&
38+
"clickhouseErrorType" in error &&
39+
error.clickhouseErrorType === "CANNOT_PARSE_ESCAPE_SEQUENCE"),
40+
insert: (batch) => insert(batch),
41+
insertAllowingBadRows: (batch) =>
42+
insert(batch, {
43+
input_format_parallel_parsing: 0,
44+
input_format_allow_errors_num: String(batch.length),
45+
input_format_allow_errors_ratio: 1,
46+
}),
47+
});
48+
}

‎apps/webapp/app/v3/eventRepository/sanitizeRowsOnParseError.server.ts‎

Lines changed: 12 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -435,12 +435,13 @@ export async function insertWithLimitedStrip<T extends object>(params: {
435435
*/
436436
async function tryInsertAllowingBadRows<T extends object>(
437437
insertAllowingBadRows: (rows: T[]) => Promise<unknown>,
438-
rows: T[]
438+
rows: T[],
439+
isParseError: (error: unknown) => boolean = isClickHouseJsonParseError
439440
): Promise<[unknown, undefined] | [undefined, unknown]> {
440441
try {
441442
return [undefined, await insertAllowingBadRows(rows)];
442443
} catch (error) {
443-
if (!isClickHouseJsonParseError(error)) throw error;
444+
if (!isParseError(error)) throw error;
444445
return [error, undefined];
445446
}
446447
}
@@ -565,14 +566,16 @@ export async function insertWithBadRowSkip<T extends object>(params: {
565566
insert: (rows: T[]) => Promise<unknown>;
566567
insertAllowingBadRows: (rows: T[]) => Promise<unknown>;
567568
hasMaterializedViews?: boolean;
569+
isParseError?: (error: unknown) => boolean;
568570
}): Promise<JsonParseRecoveryOutcome> {
569571
const { rows, contextLabel, logger, logContext, insert, insertAllowingBadRows } = params;
570572
const hasMaterializedViews = params.hasMaterializedViews ?? true;
573+
const isParseError = params.isParseError ?? isClickHouseJsonParseError;
571574

572575
try {
573576
return { kind: "inserted", insertResult: await insert(rows) };
574577
} catch (firstError) {
575-
if (!isClickHouseJsonParseError(firstError)) throw firstError;
578+
if (!isParseError(firstError)) throw firstError;
576579

577580
const firstMessage = errorMessage(firstError);
578581
const { rowsTouched, fieldsSanitized } = sanitizeRows(rows);
@@ -590,11 +593,15 @@ export async function insertWithBadRowSkip<T extends object>(params: {
590593
try {
591594
return { kind: "sanitized", insertResult: await insert(rows) };
592595
} catch (retryError) {
593-
if (!isClickHouseJsonParseError(retryError)) throw retryError;
596+
if (!isParseError(retryError)) throw retryError;
594597
}
595598
}
596599

597-
const [skipError, insertResult] = await tryInsertAllowingBadRows(insertAllowingBadRows, rows);
600+
const [skipError, insertResult] = await tryInsertAllowingBadRows(
601+
insertAllowingBadRows,
602+
rows,
603+
isParseError
604+
);
598605

599606
if (skipError) {
600607
return wholeBatchDropped({

‎apps/webapp/test/clickhouseEventRepositoryDualWrite.test.ts‎

Lines changed: 79 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,11 +1,22 @@
1-
import { ClickHouse, type TaskEventV2Input } from "@internal/clickhouse";
1+
import {
2+
ClickHouse,
3+
TASK_EVENT_SEARCH_V2_INSERT_COLUMNS,
4+
toTaskEventSearchV2Row,
5+
type TaskEventSearchV2Input,
6+
type TaskEventV2Input,
7+
} from "@internal/clickhouse";
28
import { clickhouseTest } from "@internal/testcontainers";
39
import { describe, expect, vi } from "vitest";
410
import { z } from "zod";
511
import {
612
ClickhouseEventRepository,
713
logsSearchRolloutSelectedRowCount,
814
} from "~/v3/eventRepository/clickhouseEventRepository.server";
15+
import { insertLogsSearchRows } from "~/v3/eventRepository/insertLogsSearchRows.server";
16+
import {
17+
INVALID_UTF16_SENTINEL,
18+
insertWithBadRowSkip,
19+
} from "~/v3/eventRepository/sanitizeRowsOnParseError.server";
920
import { latestMetrics, metricSum } from "./otlpMetrics.helpers";
1021
import { createInMemoryMetrics } from "./utils/tracing";
1122

@@ -86,7 +97,10 @@ describe("ClickhouseEventRepository logs search dual writer", () => {
8697

8798
try {
8899
const allowedEvents = Array.from({ length: 1_000 }, (_, index) =>
89-
event({ span_id: `span_allowed_${index}` })
100+
event({
101+
span_id: `span_allowed_${index}`,
102+
message: index === 500 ? "broken \uD800 escape" : "source write survives",
103+
})
90104
);
91105
(repository as any).addToBatch([
92106
...allowedEvents,
@@ -127,6 +141,63 @@ describe("ClickhouseEventRepository logs search dual writer", () => {
127141
60_000
128142
);
129143

144+
clickhouseTest(
145+
"recovers a strict search insert without losing neighboring rows",
146+
async ({ clickhouseContainer }) => {
147+
const clickhouse = new ClickHouse({
148+
url: clickhouseContainer.getConnectionUrl(),
149+
logLevel: "error",
150+
});
151+
const strictInsert = clickhouse.writer.insertUnsafe<TaskEventSearchV2Input>({
152+
name: "strict-search-insert",
153+
table: "trigger_dev.task_events_search_v2",
154+
columns: TASK_EVENT_SEARCH_V2_INSERT_COLUMNS,
155+
settings: { input_format_json_throw_on_bad_escape_sequence: 1 },
156+
});
157+
const rows = ["clean prefix", "broken \uD800 escape", "clean suffix"].map((message, index) =>
158+
toTaskEventSearchV2Row(event({ message, span_id: `span_${index}` }), new Date())
159+
);
160+
const readMessages = clickhouse.reader.query({
161+
name: "read-recovered-search-messages",
162+
query: `SELECT message FROM trigger_dev.task_events_search_v2
163+
WHERE environment_id = {environmentId: String} ORDER BY span_id`,
164+
params: z.object({ environmentId: z.string() }),
165+
schema: z.object({ message: z.string() }),
166+
});
167+
168+
try {
169+
const [error] = await strictInsert(rows);
170+
expect(error?.clickhouseErrorType).toBe("CANNOT_PARSE_ESCAPE_SEQUENCE");
171+
const insert = async (batch: TaskEventSearchV2Input[]) => {
172+
const [insertError, result] = await strictInsert(batch);
173+
if (insertError) throw insertError;
174+
return result;
175+
};
176+
await expect(
177+
insertWithBadRowSkip({
178+
rows,
179+
contextLabel: "default-recovery",
180+
logger: console,
181+
insert,
182+
insertAllowingBadRows: insert,
183+
})
184+
).rejects.toMatchObject({ clickhouseErrorType: "CANNOT_PARSE_ESCAPE_SEQUENCE" });
185+
const outcome = await insertLogsSearchRows(strictInsert, "strict-recovery", rows, console);
186+
expect(outcome.kind).toBe("sanitized");
187+
const [queryError, messages] = await readMessages({ environmentId: "env_dual_write_test" });
188+
expect(queryError).toBeNull();
189+
expect(messages).toEqual([
190+
{ message: "clean prefix" },
191+
{ message: INVALID_UTF16_SENTINEL },
192+
{ message: "clean suffix" },
193+
]);
194+
} finally {
195+
await clickhouse.close();
196+
}
197+
},
198+
60_000
199+
);
200+
130201
clickhouseTest(
131202
"records mapping failures as dropped search rows",
132203
async ({ clickhouseContainer }) => {
@@ -154,6 +225,12 @@ describe("ClickhouseEventRepository logs search dual writer", () => {
154225
});
155226

156227
try {
228+
expect(
229+
metricSum(await latestMetrics(metrics), "logs_search.dual_write.rows_dropped", {
230+
table: "task_events_search_v2",
231+
reason: "mapping_failed",
232+
})
233+
).toBe(0);
157234
(repository as any).addToBatch([
158235
event({ start_time: new Date().toISOString().replace("T", " ").replace("Z", "") }),
159236
]);

‎apps/webapp/test/sanitizeRowsOnParseError.test.ts‎

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -74,6 +74,18 @@ describe("isClickHouseJsonParseError", () => {
7474
expect(isClickHouseJsonParseError(err)).toBe(true);
7575
});
7676

77+
it("does not classify escape errors as native JSON parse errors", () => {
78+
expect(
79+
isClickHouseJsonParseError({ type: "CANNOT_PARSE_ESCAPE_SEQUENCE", message: "bad escape" })
80+
).toBe(false);
81+
expect(
82+
isClickHouseJsonParseError({
83+
clickhouseErrorType: "CANNOT_PARSE_ESCAPE_SEQUENCE",
84+
message: "bad escape",
85+
})
86+
).toBe(false);
87+
});
88+
7789
it("returns false for unrelated errors", () => {
7890
expect(isClickHouseJsonParseError(new Error("Connection refused"))).toBe(false);
7991
expect(

‎internal-packages/clickhouse/src/taskEventsSearch.ts‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -120,6 +120,9 @@ export function insertTaskEventsSearchV2(
120120
name: "insertTaskEventsSearchV2",
121121
table: "trigger_dev.task_events_search_v2",
122122
columns: TASK_EVENT_SEARCH_V2_INSERT_COLUMNS,
123+
settings: {
124+
input_format_json_throw_on_bad_escape_sequence: 0,
125+
},
123126
});
124127
}
125128

0 commit comments

Comments
 (0)