Skip to content
Merged
Show file tree
Hide file tree
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
8 changes: 8 additions & 0 deletions src/config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,14 @@ export const config = {
logFetchLimit: parsePositiveInt(process.env.MOTEL_OTEL_LOG_LIMIT, 80),
retentionHours: parsePositiveInt(process.env.MOTEL_OTEL_RETENTION_HOURS, 168),
maxDbSizeMb: parsePositiveInt(process.env.MOTEL_OTEL_MAX_DB_SIZE_MB, 1024),
// Stored spans are also bounded by count, independent of their size on disk.
maxSpans: parsePositiveInt(process.env.MOTEL_OTEL_MAX_SPANS, 1_000_000),
// One retention pass evicts until every bound holds, within this much writer time.
retentionPassBudgetMs: parsePositiveInt(process.env.MOTEL_OTEL_RETENTION_PASS_BUDGET_MS, 500),
// Exports the workerd collector reads at once, and the largest it accepts. Beyond
// either it refuses the export and counts it rather than buffering it.
maxPendingIngest: parsePositiveInt(process.env.MOTEL_OTEL_MAX_PENDING_INGEST, 16),
maxIngestBytes: parsePositiveInt(process.env.MOTEL_OTEL_MAX_INGEST_BYTES, 16 * 1024 * 1024),
retentionTraceBatch: parsePositiveInt(process.env.MOTEL_OTEL_RETENTION_TRACE_BATCH, 100),
retentionLogBatch: parsePositiveInt(process.env.MOTEL_OTEL_RETENTION_LOG_BATCH, 5_000),
retentionIntervalSeconds: parsePositiveInt(process.env.MOTEL_OTEL_RETENTION_INTERVAL_SECONDS, 10),
Expand Down
147 changes: 94 additions & 53 deletions src/services/TelemetryStore.ts
Original file line number Diff line number Diff line change
Expand Up @@ -478,6 +478,11 @@ export class TelemetryStore extends Context.Service<
readonly ingestTraces: (payload: OtlpTraceExportRequest) => Effect.Effect<{ readonly insertedSpans: number }, Error>
readonly ingestLogs: (payload: OtlpLogExportRequest) => Effect.Effect<{ readonly insertedLogs: number }, Error>
readonly runRetentionNow: Effect.Effect<void, Error>
/**
* Evict now when stored telemetry exceeds `MOTEL_OTEL_MAX_DB_SIZE_MB`, so the store's size
* holds between retention passes rather than only after each. A no-op within the bound.
*/
readonly holdSizeBound: Effect.Effect<void, Error>
}
>()("motel/TelemetryStore") {}

Expand Down Expand Up @@ -884,73 +889,108 @@ export const makeTelemetryStoreEffect = (db: TelemetryDatabase, opts: TelemetryS
// across traces and corrupted the summary rebuild). Running traces
// are protected — only `active_span_count = 0` summaries are in
// scope for eviction.
const toEvict = new Set<string>()

// Time-based: completed traces whose last span ended before cutoff.
const timeExpired = db.query(
`SELECT trace_id FROM trace_summaries WHERE active_span_count = 0 AND ended_at_ms > 0 AND ended_at_ms < ? ORDER BY ended_at_ms ASC LIMIT ?`,
).all(cutoff, config.otel.retentionTraceBatch) as readonly { trace_id: string }[]
for (const row of timeExpired) toEvict.add(row.trace_id)

// Size-based: if actual data exceeds the target, drop one bounded
// batch of the oldest completed traces. `(page_count - freelist_count)`
// ignores freed-but-not-vacuumed pages so a large freelist doesn't
// trigger a deletion death spiral.
const dbSize = db.usedBytes()
if (dbSize > maxDbSizeBytes) {
const oldest = db.query(
`SELECT trace_id FROM trace_summaries WHERE active_span_count = 0 ORDER BY started_at_ms ASC LIMIT ?`,
).all(config.otel.retentionTraceBatch) as readonly { trace_id: string }[]
// Set.add dedupes overlap with the time-expired batch above.
for (const row of oldest) toEvict.add(row.trace_id)
}
//
// Every bound holds after each pass, not eventually: one pass evicts
// batches until the store is back within its age, span-count and size
// bounds, so a sustained ingest rate above one batch per interval
// cannot outgrow them. A time budget keeps the pass from holding the
// single writer for long; the next interval continues from there.
const passStarted = Date.now()
let spans = Number((db.query(`SELECT COALESCE(SUM(span_count), 0) AS value FROM trace_summaries`).get() as { value: number }).value)
let dbSize = db.usedBytes()
let evictedTraces = 0
let deletedLogs = false
let deletedOrphans = false
while (Date.now() - passStarted < config.otel.retentionPassBudgetMs) {
const toEvict = new Map<string, number>()

// Time-based: completed traces whose last span ended before cutoff.
const timeExpired = db.query(
`SELECT trace_id, span_count FROM trace_summaries WHERE active_span_count = 0 AND ended_at_ms > 0 AND ended_at_ms < ? ORDER BY ended_at_ms ASC LIMIT ?`,
).all(cutoff, config.otel.retentionTraceBatch) as readonly { trace_id: string; span_count: number }[]
for (const row of timeExpired) toEvict.set(row.trace_id, row.span_count)

// Count and size: drop the oldest completed traces until the excess is
// gone. Size is converted to spans at the store's current bytes per
// span; `usedBytes` ignores freed-but-not-vacuumed pages so a large
// freelist doesn't trigger a deletion death spiral.
const excess = Math.max(
spans - config.otel.maxSpans,
dbSize > maxDbSizeBytes && spans > 0 ? Math.ceil((dbSize - maxDbSizeBytes) / (dbSize / spans)) : 0,
dbSize > maxDbSizeBytes ? 1 : 0,
)
if (excess > 0) {
const oldest = db.query(
`SELECT trace_id, span_count FROM trace_summaries WHERE active_span_count = 0 ORDER BY started_at_ms ASC LIMIT ?`,
).all(config.otel.retentionTraceBatch) as readonly { trace_id: string; span_count: number }[]
let selected = 0
for (const row of oldest) {
if (selected >= excess) break
if (!toEvict.has(row.trace_id)) selected += row.span_count
toEvict.set(row.trace_id, row.span_count)
}
}

// Logs have their own retention boundary. A correlated log may refer
// to a trace that was sampled elsewhere or never reached Motel, so
// tying log eviction to trace_summaries lets those rows grow forever.
const expiredLogs = db.query(`DELETE FROM logs WHERE id IN (SELECT id FROM logs WHERE timestamp_ms < ? ORDER BY timestamp_ms ASC LIMIT ?)`).run(cutoff, config.otel.retentionLogBatch)
let deletedLogs = Number(expiredLogs.changes) > 0
if (dbSize > maxDbSizeBytes) {
const oversizedLogs = db.query(`DELETE FROM logs WHERE id IN (SELECT id FROM logs ORDER BY timestamp_ms ASC LIMIT ?)`).run(config.otel.retentionLogBatch)
deletedLogs = deletedLogs || Number(oversizedLogs.changes) > 0
}
// Logs have their own retention boundary. A correlated log may refer
// to a trace that was sampled elsewhere or never reached Motel, so
// tying log eviction to trace_summaries lets those rows grow forever.
const expiredLogs = db.query(`DELETE FROM logs WHERE id IN (SELECT id FROM logs WHERE timestamp_ms < ? ORDER BY timestamp_ms ASC LIMIT ?)`).run(cutoff, config.otel.retentionLogBatch)
let deletedLogsNow = Number(expiredLogs.changes) > 0
if (dbSize > maxDbSizeBytes) {
const oversizedLogs = db.query(`DELETE FROM logs WHERE id IN (SELECT id FROM logs ORDER BY timestamp_ms ASC LIMIT ?)`).run(config.otel.retentionLogBatch)
deletedLogsNow = deletedLogsNow || Number(oversizedLogs.changes) > 0
}
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.
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)
try {
db.query(`DELETE FROM span_operation_fts WHERE trace_id IN (${placeholders})`).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)
}

// 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.
const traceIds = Array.from(toEvict)
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)
// Log-side orphans (log_attributes + FTS) are keyed by log.id,
// so prune what no longer has a parent log row.
const orphanAttributes = db.query(`DELETE FROM log_attributes WHERE rowid IN (SELECT log_attributes.rowid FROM log_attributes WHERE NOT EXISTS (SELECT 1 FROM logs WHERE logs.id = log_attributes.log_id) LIMIT ?)`).run(config.otel.retentionLogBatch)
deletedOrphans = deletedOrphans || Number(orphanAttributes.changes) > 0
try {
db.query(`DELETE FROM span_operation_fts WHERE trace_id IN (${placeholders})`).run(...batch)
const orphanFts = db.query(`DELETE FROM log_body_fts WHERE rowid IN (SELECT rowid FROM log_body_fts WHERE NOT EXISTS (SELECT 1 FROM logs WHERE logs.id = CAST(log_body_fts.log_id AS INTEGER)) LIMIT ?)`).run(config.otel.retentionLogBatch)
deletedOrphans = deletedOrphans || Number(orphanFts.changes) > 0
} 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)
}

// Log-side orphans (log_attributes + FTS) are keyed by log.id,
// so prune what no longer has a parent log row.
const orphanAttributes = db.query(`DELETE FROM log_attributes WHERE rowid IN (SELECT log_attributes.rowid FROM log_attributes WHERE NOT EXISTS (SELECT 1 FROM logs WHERE logs.id = log_attributes.log_id) LIMIT ?)`).run(config.otel.retentionLogBatch)
let deletedOrphans = Number(orphanAttributes.changes) > 0
try {
const orphanFts = db.query(`DELETE FROM log_body_fts WHERE rowid IN (SELECT rowid FROM log_body_fts WHERE NOT EXISTS (SELECT 1 FROM logs WHERE logs.id = CAST(log_body_fts.log_id AS INTEGER)) LIMIT ?)`).run(config.otel.retentionLogBatch)
deletedOrphans = deletedOrphans || Number(orphanFts.changes) > 0
} catch {
// FTS table may not exist on old DBs.
for (const count of toEvict.values()) spans -= count
evictedTraces += toEvict.size
const sizeBefore = dbSize
dbSize = db.usedBytes()
// Nothing left to evict, or every bound holds again.
if (toEvict.size === 0 && !deletedLogsNow) break
const backlog = timeExpired.length === config.otel.retentionTraceBatch || Number(expiredLogs.changes) === config.otel.retentionLogBatch
if (spans <= config.otel.maxSpans && dbSize <= maxDbSizeBytes && !backlog) break
// Some deletions free pages only after a later merge or checkpoint. When a
// pass freed nothing measurable, wait for that rather than evict more.
if (spans <= config.otel.maxSpans && !backlog && dbSize >= sizeBefore) break
}

// Checkpoint after a big delete pass so the freed pages land
// in the main DB file and become eligible for incremental
// vacuum. Use RESTART (not PASSIVE): PASSIVE silently no-ops
// when readers are active, which is the documented mechanism
// behind WAL/freelist starvation when ingest is busy.
if (toEvict.size === 0 && !deletedLogs && !deletedOrphans) return
if (evictedTraces === 0 && !deletedLogs && !deletedOrphans) return
db.checkpoint("RESTART")

// Incremental FTS5 merge — DELETE on an FTS5-indexed row
Expand Down Expand Up @@ -2411,5 +2451,6 @@ export const makeTelemetryStoreEffect = (db: TelemetryDatabase, opts: TelemetryS
getAiCall,
aiCallStats,
runRetentionNow: Effect.andThen(backfillBatch, Effect.andThen(reconcileTraceSummaries, cleanupExpired())),
holdSizeBound: Effect.suspend(() => (db.usedBytes() > maxDbSizeBytes ? cleanupExpired() : Effect.void)),
})
})
Loading
Loading