Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
156 changes: 53 additions & 103 deletions src/services/TelemetryStore.ts
Original file line number Diff line number Diff line change
Expand Up @@ -395,6 +395,13 @@ const percentile = (values: readonly number[], ratio: number) => {
return sorted[index] ?? 0
}

/**
* List-valued SQL inputs bind as one JSON array: workerd allows at most 100 bound parameters per
* statement, and a list must not depend on how many values it holds.
*/
const jsonValues = "(SELECT value FROM json_each(?))"
const jsonPairs = "(SELECT json_extract(value, '$[0]'), json_extract(value, '$[1]') FROM json_each(?))"

const tokenizeFts = (value: string) => value.match(/[A-Za-z0-9_]+/g)?.filter((token) => token.length > 1) ?? []

const toFtsMatchQuery = (value: string) => {
Expand Down Expand Up @@ -777,82 +784,38 @@ export const makeTelemetryStoreEffect = (db: TelemetryDatabase, opts: TelemetryS
})

const deleteSpanAttributes = db.query(`DELETE FROM span_attributes WHERE trace_id = ? AND span_id = ?`)
const insertSpanAttribute = db.query(`INSERT INTO span_attributes (trace_id, span_id, key, value) VALUES (?, ?, ?, ?)`)
const spanAttributeInsertManyByCount = new Map<number, ReturnType<TelemetryDatabase["query"]>>()
// Lists travel as one JSON parameter: workerd binds at most 100 parameters per statement.
const insertSpanAttributes = db.query(`INSERT INTO span_attributes (trace_id, span_id, key, value) SELECT ?, ?, key, value FROM json_each(?)`)
const insertSpanAttributesMany = (traceId: string, spanId: string, attributes: Readonly<Record<string, string>>) => {
const entries = Object.entries(attributes)
if (entries.length === 0) return
if (entries.length === 1) {
const [key, value] = entries[0]!
insertSpanAttribute.run(traceId, spanId, key, value)
return
}
let query = spanAttributeInsertManyByCount.get(entries.length)
if (!query) {
query = db.query(`INSERT INTO span_attributes (trace_id, span_id, key, value) VALUES ${entries.map(() => "(?, ?, ?, ?)").join(", ")}`)
spanAttributeInsertManyByCount.set(entries.length, query)
}
query.run(...entries.flatMap(([key, value]) => [traceId, spanId, key, value]))
if (Object.keys(attributes).length === 0) return
insertSpanAttributes.run(traceId, spanId, JSON.stringify(attributes))
}
const deleteSpanOperationSearch = db.query(`DELETE FROM span_operation_fts WHERE trace_id = ? AND span_id = ?`)
const insertSpanOperationSearch = db.query(`INSERT INTO span_operation_fts (trace_id, span_id, operation_name) VALUES (?, ?, ?)`)
const deleteSpanOperationSearchManyByCount = new Map<number, ReturnType<TelemetryDatabase["query"]>>()
const insertSpanOperationSearchManyByCount = new Map<number, ReturnType<TelemetryDatabase["query"]>>()
const deleteSpanOperationSearch = db.query(`
DELETE FROM span_operation_fts
WHERE (trace_id, span_id) IN (SELECT json_extract(value, '$[0]'), json_extract(value, '$[1]') FROM json_each(?))
`)
const insertSpanOperationSearch = db.query(`
INSERT INTO span_operation_fts (trace_id, span_id, operation_name)
SELECT json_extract(value, '$[0]'), json_extract(value, '$[1]'), json_extract(value, '$[2]') FROM json_each(?)
`)
const updateSpanOperationSearchMany = (operations: ReadonlyArray<readonly [string, string, string]>) => {
if (operations.length === 0) return
if (operations.length === 1) {
const [traceId, spanId, operationName] = operations[0]!
deleteSpanOperationSearch.run(traceId, spanId)
insertSpanOperationSearch.run(traceId, spanId, operationName)
return
}

let deleteQuery = deleteSpanOperationSearchManyByCount.get(operations.length)
if (!deleteQuery) {
deleteQuery = db.query(`DELETE FROM span_operation_fts WHERE ${operations.map(() => "(trace_id = ? AND span_id = ?)").join(" OR ")}`)
deleteSpanOperationSearchManyByCount.set(operations.length, deleteQuery)
}
deleteQuery.run(...operations.flatMap(([traceId, spanId]) => [traceId, spanId]))

let insertQuery = insertSpanOperationSearchManyByCount.get(operations.length)
if (!insertQuery) {
insertQuery = db.query(`INSERT INTO span_operation_fts (trace_id, span_id, operation_name) VALUES ${operations.map(() => "(?, ?, ?)").join(", ")}`)
insertSpanOperationSearchManyByCount.set(operations.length, insertQuery)
}
insertQuery.run(...operations.flatMap(([traceId, spanId, operationName]) => [traceId, spanId, operationName]))
const list = JSON.stringify(operations)
deleteSpanOperationSearch.run(list)
insertSpanOperationSearch.run(list)
}
const insertLogAttribute = db.query(`INSERT INTO log_attributes (log_id, key, value) VALUES (?, ?, ?)`)
const logAttributeInsertManyByCount = new Map<number, ReturnType<TelemetryDatabase["query"]>>()
const insertLogAttributes = db.query(`INSERT INTO log_attributes (log_id, key, value) SELECT ?, key, value FROM json_each(?)`)
const insertLogAttributesMany = (logId: number, attributes: Readonly<Record<string, string>>) => {
const entries = Object.entries(attributes)
if (entries.length === 0) return
if (entries.length === 1) {
const [key, value] = entries[0]!
insertLogAttribute.run(logId, key, value)
return
}
let query = logAttributeInsertManyByCount.get(entries.length)
if (!query) {
query = db.query(`INSERT INTO log_attributes (log_id, key, value) VALUES ${entries.map(() => "(?, ?, ?)").join(", ")}`)
logAttributeInsertManyByCount.set(entries.length, query)
}
query.run(...entries.flatMap(([key, value]) => [logId, key, value]))
if (Object.keys(attributes).length === 0) return
insertLogAttributes.run(logId, JSON.stringify(attributes))
}
const insertLogBodySearch = db.query(`INSERT INTO log_body_fts (log_id, body) VALUES (?, ?)`)
const insertLogBodySearchManyByCount = new Map<number, ReturnType<TelemetryDatabase["query"]>>()
const insertLogBodySearch = db.query(`
INSERT INTO log_body_fts (log_id, body)
SELECT json_extract(value, '$[0]'), json_extract(value, '$[1]') FROM json_each(?)
`)
const insertLogBodySearchMany = (entries: ReadonlyArray<readonly [string, string]>) => {
if (entries.length === 0) return
if (entries.length === 1) {
const [logId, body] = entries[0]!
insertLogBodySearch.run(logId, body)
return
}
let query = insertLogBodySearchManyByCount.get(entries.length)
if (!query) {
query = db.query(`INSERT INTO log_body_fts (log_id, body) VALUES ${entries.map(() => "(?, ?)").join(", ")}`)
insertLogBodySearchManyByCount.set(entries.length, query)
}
query.run(...entries.flatMap(([logId, body]) => [logId, body]))
insertLogBodySearch.run(JSON.stringify(entries))
}

const maxDbSizeBytes = config.otel.maxDbSizeMb * 1024 * 1024
Expand Down Expand Up @@ -942,23 +905,21 @@ export const makeTelemetryStoreEffect = (db: TelemetryDatabase, opts: TelemetryS
}
deletedLogs = deletedLogs || deletedLogsNow

// Batch the trace-id list so the IN placeholders stay under
// SQLite's default limit (~999). Each batch wipes every row
// reachable from those trace_ids across the cascade tables.
// Each batch wipes every row reachable from those trace_ids
// across the cascade tables.
const traceIds = Array.from(toEvict.keys())
const BATCH_SIZE = 500
for (let offset = 0; offset < traceIds.length; offset += BATCH_SIZE) {
const batch = traceIds.slice(offset, offset + BATCH_SIZE)
const placeholders = batch.map(() => "?").join(",")
db.query(`DELETE FROM span_attributes WHERE trace_id IN (${placeholders})`).run(...batch)
const batch = JSON.stringify(traceIds.slice(offset, offset + BATCH_SIZE))
db.query(`DELETE FROM span_attributes WHERE trace_id IN ${jsonValues}`).run(batch)
try {
db.query(`DELETE FROM span_operation_fts WHERE trace_id IN (${placeholders})`).run(...batch)
db.query(`DELETE FROM span_operation_fts WHERE trace_id IN ${jsonValues}`).run(batch)
} catch {
// FTS table may not exist on old DBs.
}
db.query(`DELETE FROM spans WHERE trace_id IN (${placeholders})`).run(...batch)
db.query(`DELETE FROM logs WHERE trace_id IN (${placeholders})`).run(...batch)
db.query(`DELETE FROM trace_summaries WHERE trace_id IN (${placeholders})`).run(...batch)
db.query(`DELETE FROM spans WHERE trace_id IN ${jsonValues}`).run(batch)
db.query(`DELETE FROM logs WHERE trace_id IN ${jsonValues}`).run(batch)
db.query(`DELETE FROM trace_summaries WHERE trace_id IN ${jsonValues}`).run(batch)
}

// Log-side orphans (log_attributes + FTS) are keyed by log.id,
Expand Down Expand Up @@ -1225,12 +1186,11 @@ export const makeTelemetryStoreEffect = (db: TelemetryDatabase, opts: TelemetryS

const loadTracesByIds = (traceIds: readonly string[]) => {
if (traceIds.length === 0) return [] as readonly TraceItem[]
const placeholders = traceIds.map(() => "?").join(", ")
const rows = db.query(`
SELECT * FROM spans
WHERE trace_id IN (${placeholders})
WHERE trace_id IN ${jsonValues}
ORDER BY start_time_ms ASC
`).all(...traceIds) as SpanRow[]
`).all(JSON.stringify(traceIds)) as SpanRow[]

const grouped = new Map<string, SpanRow[]>()
for (const row of rows) {
Expand Down Expand Up @@ -1502,12 +1462,11 @@ export const makeTelemetryStoreEffect = (db: TelemetryDatabase, opts: TelemetryS
// touched by the candidate set. One indexed scan per trace_id
// is much cheaper than a per-span lookup loop while computing
// depth, and we get the trace-root lookup in the same pass.
const placeholders = traceIds.map(() => "?").join(", ")
const allSpanRows = db.query(`
SELECT trace_id, span_id, parent_span_id, operation_name, start_time_ms
FROM spans
WHERE trace_id IN (${placeholders})
`).all(...traceIds) as Array<{ trace_id: string; span_id: string; parent_span_id: string | null; operation_name: string; start_time_ms: number }>
WHERE trace_id IN ${jsonValues}
`).all(JSON.stringify(traceIds)) as Array<{ trace_id: string; span_id: string; parent_span_id: string | null; operation_name: string; start_time_ms: number }>

const rootOperationByTraceId = new Map<string, { operationName: string; startTimeMs: number }>()
for (const row of allSpanRows) {
Expand Down Expand Up @@ -1559,10 +1518,9 @@ export const makeTelemetryStoreEffect = (db: TelemetryDatabase, opts: TelemetryS
// using SQLite's row-value `IN` syntax, then parseSpanRow per
// kept row. Result order follows `filteredLite` so the caller
// sees the same ordering the candidate scan produced.
const keptValues = filteredLite.map(() => "(?, ?)").join(", ")
const fullRows = db.query(`
SELECT * FROM spans WHERE (trace_id, span_id) IN (VALUES ${keptValues})
`).all(...filteredLite.flatMap((row) => [row.trace_id, row.span_id])) as SpanRow[]
SELECT * FROM spans WHERE (trace_id, span_id) IN ${jsonPairs}
`).all(JSON.stringify(filteredLite.map((row) => [row.trace_id, row.span_id]))) as SpanRow[]
const fullRowByKey = new Map<string, SpanRow>()
for (const row of fullRows) {
fullRowByKey.set(keyOf(row.trace_id, row.span_id), row)
Expand Down Expand Up @@ -1687,13 +1645,12 @@ export const makeTelemetryStoreEffect = (db: TelemetryDatabase, opts: TelemetryS
const attrKey = input.groupBy.slice(5)
const traceIds = summaries.map((s) => s.traceId)
if (traceIds.length > 0) {
const placeholders = traceIds.map(() => "?").join(", ")
const rows = db.query(`
SELECT trace_id, value
FROM span_attributes
WHERE key = ? AND trace_id IN (${placeholders})
WHERE key = ? AND trace_id IN ${jsonValues}
GROUP BY trace_id
`).all(attrKey, ...traceIds) as Array<{ trace_id: string; value: string }>
`).all(attrKey, JSON.stringify(traceIds)) as Array<{ trace_id: string; value: string }>

attrLookup = new Map()
for (const row of rows) {
Expand Down Expand Up @@ -2129,16 +2086,12 @@ export const makeTelemetryStoreEffect = (db: TelemetryDatabase, opts: TelemetryS
/** Load attribute values for a set of spans by key */
const loadSpanAttrValues = (spans: ReadonlyArray<{ trace_id: string; span_id: string }>, keys: readonly string[]): Map<string, Map<string, string>> => {
if (spans.length === 0 || keys.length === 0) return new Map()
const spanPlaceholders = spans.map(() => "(?, ?)").join(", ")
const keyPlaceholders = keys.map(() => "?").join(", ")
const spanParams = spans.flatMap((s) => [s.trace_id, s.span_id])

const rows = db.query(`
SELECT trace_id, span_id, key, value
FROM span_attributes
WHERE (trace_id, span_id) IN (VALUES ${spanPlaceholders})
AND key IN (${keyPlaceholders})
`).all(...spanParams, ...keys) as Array<{ trace_id: string; span_id: string; key: string; value: string }>
WHERE (trace_id, span_id) IN ${jsonPairs}
AND key IN ${jsonValues}
`).all(JSON.stringify(spans.map((s) => [s.trace_id, s.span_id])), JSON.stringify(keys)) as Array<{ trace_id: string; span_id: string; key: string; value: string }>

const result = new Map<string, Map<string, string>>()
for (const row of rows) {
Expand Down Expand Up @@ -2184,15 +2137,13 @@ export const makeTelemetryStoreEffect = (db: TelemetryDatabase, opts: TelemetryS
const attrMap = loadSpanAttrValues(rows, summaryAttrKeys)

// Count tool call child spans per AI span
const spanPlaceholders = rows.map(() => "(?, ?)").join(", ")
const spanParams = rows.flatMap((r) => [r.trace_id, r.span_id])
const toolCountRows = db.query(`
SELECT parent_span_id, COUNT(*) AS cnt
FROM spans
WHERE (trace_id, parent_span_id) IN (VALUES ${spanPlaceholders})
WHERE (trace_id, parent_span_id) IN ${jsonPairs}
AND operation_name LIKE 'ai.toolCall%'
GROUP BY trace_id, parent_span_id
`).all(...spanParams) as Array<{ parent_span_id: string; cnt: number }>
`).all(JSON.stringify(rows.map((r) => [r.trace_id, r.span_id]))) as Array<{ parent_span_id: string; cnt: number }>
const toolCounts = new Map(toolCountRows.map((r) => [r.parent_span_id, r.cnt]))

return rows.map((row): AiCallSummary => {
Expand Down Expand Up @@ -2394,12 +2345,11 @@ export const makeTelemetryStoreEffect = (db: TelemetryDatabase, opts: TelemetryS
const tokenKey = input.agg === "total_input_tokens" ? AI_ATTR_MAP.inputTokens : AI_ATTR_MAP.outputTokens
const allSpanIds = [...groups.values()].flatMap((b) => b.spanIds)
if (allSpanIds.length > 0) {
const placeholders = allSpanIds.map(() => "?").join(", ")
const tokenRows = db.query(`
SELECT span_id, CAST(value AS REAL) AS tokens
FROM span_attributes
WHERE key = ? AND span_id IN (${placeholders})
`).all(tokenKey, ...allSpanIds) as Array<{ span_id: string; tokens: number }>
WHERE key = ? AND span_id IN ${jsonValues}
`).all(tokenKey, JSON.stringify(allSpanIds)) as Array<{ span_id: string; tokens: number }>

const tokenBySpan = new Map(tokenRows.map((r) => [r.span_id, r.tokens]))

Expand Down
Loading