diff --git a/src/config.ts b/src/config.ts index 5c5010f..e95ce3c 100644 --- a/src/config.ts +++ b/src/config.ts @@ -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), diff --git a/src/services/TelemetryStore.ts b/src/services/TelemetryStore.ts index 6879735..b3a1ea4 100644 --- a/src/services/TelemetryStore.ts +++ b/src/services/TelemetryStore.ts @@ -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 + /** + * 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 } >()("motel/TelemetryStore") {} @@ -884,65 +889,100 @@ 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() - - // 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() + + // 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 @@ -950,7 +990,7 @@ export const makeTelemetryStoreEffect = (db: TelemetryDatabase, opts: TelemetryS // 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 @@ -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)), }) }) diff --git a/src/workerd.ts b/src/workerd.ts index a3fa047..5fe50b1 100644 --- a/src/workerd.ts +++ b/src/workerd.ts @@ -1,5 +1,5 @@ import type { DurableObjectState } from "@cloudflare/workers-types" -import { Effect, Layer, Schema } from "effect" +import { Cause, Effect, Layer, Schema } from "effect" import { HttpRouter, HttpServer } from "effect/unstable/http" import { config } from "./config.js" import { motelApi } from "./httpServer.js" @@ -8,6 +8,7 @@ import { IngestError } from "./services/ingestRpc.js" import { makeTelemetryStoreEffect, TelemetryStoreReadonly, type TelemetryStore } from "./services/TelemetryStore.js" import { workerdDatabase } from "./services/TelemetryStoreWorkerd.js" import { TraceExport, LogExport } from "./otlpSchema.js" +import { decodeProtobufLogs, decodeProtobufTraces } from "./otlpProtobuf.js" import documents from "motel:documents" import pkg from "../package.json" with { type: "json" } @@ -27,54 +28,234 @@ const health = () => ({ version: pkg.version, }) +const ingestPaths = new Set(["/v1/traces", "/v1/logs"]) +const refusalLogInterval = 60_000 + +/** Read a body without a declared length, giving up once it exceeds `limit` bytes. */ +const boundedText = async (request: Request, limit: number): Promise => { + if (request.body === null) return "" + const reader = request.body.getReader() + const decoder = new TextDecoder() + let text = "" + let bytes = 0 + for (;;) { + const { done, value } = await reader.read() + if (done) return text + decoder.decode() + bytes += value.byteLength + if (bytes > limit) { + await reader.cancel() + return undefined + } + text += decoder.decode(value, { stream: true }) + } +} + +/** Why the collector did not store an export. Each is counted; see `MotelCollector.ingestStatus`. */ +type Loss = "queueFull" | "tooLarge" | "invalid" | "storeFailed" +const losses: ReadonlyArray = ["queueFull", "tooLarge", "invalid", "storeFailed"] +interface LossCounts { + /** Exports not stored, by reason. */ + readonly refused: Record + /** Declared bytes of those exports, where the request declared a length. */ + readonly refusedBytes: number + readonly lastRefusedAt: string | null +} +const lossKey = "motel:ingest-losses" + /** One durable SQLite owner for OTLP ingestion, queries and bounded alarm maintenance. */ export class MotelCollector { private readonly state: DurableObjectState private readonly store: Promise + private readonly handler: Promise<(request: Request) => Promise> + /** Exports being read or stored now. Bounded by `maxPendingIngest`. */ + private pending = 0 + /** + * Exports the collector did not store. Kept in the actor's storage, so the counts survive + * eviction and restarts; `GET /api/ingest` reports them and a warning names them in the log. + */ + private losses: LossCounts + private lastLossLog = 0 constructor(state: DurableObjectState) { this.state = state + const saved = state.storage.kv.get(lossKey) + this.losses = { + refused: Object.fromEntries(losses.map((loss) => [loss, saved?.refused[loss] ?? 0])) as Record, + refusedBytes: saved?.refusedBytes ?? 0, + lastRefusedAt: saved?.lastRefusedAt ?? null, + } // The actor owns one connection and bootstrap. No native handle or detached fiber // escapes this scoped construction; durable alarms own subsequent maintenance. this.store = state.blockConcurrencyWhile(() => Effect.runPromise(Effect.scoped(makeTelemetryStoreEffect(workerdDatabase(state.storage), { readonly: false, runRetention: false }))), ) + // A failed bootstrap is reported by each request that needs the store. + this.store.catch(() => {}) + // One router for the actor's lifetime. Building the typed API per request allocates + // far more than the request itself and leaves it for a later collection. + this.handler = this.store.then((store) => { + const ingest = Layer.succeed(AsyncIngest, { + ingestTraces: ({ payload }) => + Schema.decodeUnknownEffect(TraceExport)(payload).pipe( + Effect.mapError(() => new IngestError({ message: "Invalid trace payload" })), + Effect.flatMap((parsed) => + store.ingestTraces(parsed).pipe(Effect.mapError(() => new IngestError({ message: "Trace storage failed" }))), + ), + ), + ingestLogs: ({ payload }) => + Schema.decodeUnknownEffect(LogExport)(payload).pipe( + Effect.mapError(() => new IngestError({ message: "Invalid log payload" })), + Effect.flatMap((parsed) => + store.ingestLogs(parsed).pipe(Effect.mapError(() => new IngestError({ message: "Log storage failed" }))), + ), + ), + }) + const { handler } = HttpRouter.toWebHandler( + motelApi({ health, docs: documents }).pipe( + HttpRouter.provideRequest(ingest), + HttpRouter.provideRequest(Layer.succeed(TelemetryStoreReadonly, store)), + Layer.provide(HttpServer.layerServices), + ), + { disableLogger: true }, + ) + return (request: Request) => handler(request) + }) + this.handler.catch(() => {}) } /** Serve through the same typed HTTP routes as the native collector. */ async fetch(request: Request): Promise { - const store = await this.store + const url = new URL(request.url) + if (request.method === "GET" && url.pathname === "/api/ingest") return Response.json(this.ingestStatus()) + if (request.method === "POST" && ingestPaths.has(url.pathname)) return this.ingest(request) + return this.serve(request) + } + + private async serve(request: Request): Promise { + const handler = await this.handler if ((await this.state.storage.getAlarm()) === null) { await this.state.storage.setAlarm(Date.now() + config.otel.retentionIntervalSeconds * 1000) } - const ingest = Layer.succeed(AsyncIngest, { - ingestTraces: ({ payload }) => - Schema.decodeUnknownEffect(TraceExport)(payload).pipe( - Effect.mapError(() => new IngestError({ message: "Invalid trace payload" })), - Effect.flatMap((parsed) => - store.ingestTraces(parsed).pipe(Effect.mapError(() => new IngestError({ message: "Trace storage failed" }))), - ), - ), - ingestLogs: ({ payload }) => - Schema.decodeUnknownEffect(LogExport)(payload).pipe( - Effect.mapError(() => new IngestError({ message: "Invalid log payload" })), - Effect.flatMap((parsed) => - store.ingestLogs(parsed).pipe(Effect.mapError(() => new IngestError({ message: "Log storage failed" }))), - ), - ), - }) - const handler = HttpRouter.toWebHandler( - motelApi({ health, docs: documents }).pipe( - HttpRouter.provideRequest(ingest), - HttpRouter.provideRequest(Layer.succeed(TelemetryStoreReadonly, store)), - Layer.provide(HttpServer.layerServices), - ), - { disableLogger: true }, - ) + return handler(request) + } + + /** + * Accept an export only while the queue has room and only up to the size bound, so the + * collector's memory for exports never exceeds `maxPendingIngest * maxIngestBytes`. Every + * export it does not store is counted. OTLP exporters retry 429, 500 and 503 and drop 400 + * and 413, so a malformed or oversized export is answered as final and a storage failure as + * transient. + */ + private async ingest(request: Request): Promise { + const declared = request.headers.get("content-length") + const bytes = declared === null ? 0 : Number(declared) + if (bytes > config.otel.maxIngestBytes) { + this.discard(request) + return this.refuse("tooLarge", bytes, "Telemetry export exceeds the collector's size limit", 413) + } + if (this.pending >= config.otel.maxPendingIngest) { + this.discard(request) + return this.refuse("queueFull", bytes, "Collector ingest queue is full", 429, { "retry-after": "1" }) + } + this.pending += 1 try { - return await handler.handler(request) + const signal = new URL(request.url).pathname === "/v1/traces" ? "traces" : "logs" + const protobuf = /application\/(x-)?protobuf/i.test(request.headers.get("content-type") ?? "") + if (declared === null && protobuf) return Response.json({ error: "Protobuf exports need a Content-Length" }, { status: 411 }) + let payload: unknown + try { + if (declared === null) { + const body = await boundedText(request, config.otel.maxIngestBytes) + if (body === undefined) return this.refuse("tooLarge", 0, "Telemetry export exceeds the collector's size limit", 413) + payload = JSON.parse(body) + } else if (protobuf) { + const bytes = new Uint8Array(await request.arrayBuffer()) + payload = signal === "traces" ? decodeProtobufTraces(bytes) : decodeProtobufLogs(bytes) + } else { + // Parsed from the body without keeping its text: a whole export is a large string, + // and one held while it is stored outlives the young generation. + payload = await request.json() + } + } catch (cause) { + return this.refuse("invalid", bytes, `Invalid ${signal} export`, 400, {}, cause) + } + let store: TelemetryStore["Service"] + try { + store = await this.store + } catch (cause) { + return this.refuse("storeFailed", bytes, "Telemetry store is unavailable", 503, { "retry-after": "10" }, cause) + } + return await this.write(store, signal, payload, bytes) } finally { - await handler.dispose() + this.pending -= 1 + } + } + + /** Store one decoded export. Spans and logs are written before the response is sent. */ + private async write(store: TelemetryStore["Service"], signal: "traces" | "logs", payload: unknown, bytes: number): Promise { + let stored: Effect.Effect + if (signal === "traces") { + const decoded = Schema.decodeUnknownExit(TraceExport)(payload) + if (decoded._tag === "Failure") return this.refuse("invalid", bytes, "Invalid traces export", 400, {}, Cause.squash(decoded.cause)) + stored = store.ingestTraces(decoded.value) + } else { + const decoded = Schema.decodeUnknownExit(LogExport)(payload) + if (decoded._tag === "Failure") return this.refuse("invalid", bytes, "Invalid logs export", 400, {}, Cause.squash(decoded.cause)) + stored = store.ingestLogs(decoded.value) + } + try { + if ((await this.state.storage.getAlarm()) === null) { + await this.state.storage.setAlarm(Date.now() + config.otel.retentionIntervalSeconds * 1000) + } + // Each export is stored before the response, then the size bound is held again. + const result = await Effect.runPromise(Effect.tap(stored, () => store.holdSizeBound)) + return Response.json(result) + } catch (cause) { + return this.refuse("storeFailed", bytes, `The ${signal} export could not be stored`, 500, {}, cause) + } + } + + /** Discard an unread body; reading it is the work being refused. */ + private discard(request: Request) { + request.body?.cancel().catch(() => {}) + } + + private refuse( + loss: Loss, + bytes: number, + error: string, + status: number, + headers: Record = {}, + cause?: unknown, + ): Response { + this.losses = { + refused: { ...this.losses.refused, [loss]: this.losses.refused[loss] + 1 }, + refusedBytes: this.losses.refusedBytes + bytes, + lastRefusedAt: new Date().toISOString(), + } + this.state.storage.kv.put(lossKey, this.losses) + const now = Date.now() + // A storage failure is logged every time with its cause; refusals under load once a minute. + if (loss === "storeFailed" || now - this.lastLossLog >= refusalLogInterval) { + this.lastLossLog = now + console.warn( + JSON.stringify({ + message: "motel: telemetry exports not stored", + reason: loss, + ...(cause === undefined ? {} : { cause: String(cause instanceof Error ? cause.message : cause) }), + ...this.ingestStatus(), + }), + ) + } + return Response.json({ error }, { status, headers }) + } + + private ingestStatus() { + return { + pending: this.pending, + maxPending: config.otel.maxPendingIngest, + maxBytes: config.otel.maxIngestBytes, + ...this.losses, } } diff --git a/workerd/README.md b/workerd/README.md index 3b913f8..f793bef 100644 --- a/workerd/README.md +++ b/workerd/README.md @@ -26,6 +26,29 @@ read-only disk service named `ASSETS`, and persistent local disk storage. Settings use the same `MOTEL_OTEL_*` names, supplied as text bindings. Health uses PID 0 because a worker does not own an operating-system process. +The collector's memory and store are bounded by configuration. It reads at most +`MOTEL_OTEL_MAX_PENDING_INGEST` exports at once (default 16), each at most +`MOTEL_OTEL_MAX_INGEST_BYTES` (default 16 MiB). It answers a further export with +429 and `Retry-After: 1`, and an oversized one with 413. A malformed export gets +400. When the store cannot be opened or written it answers 503 or 500, which OTLP +exporters retry. Every export it does not store is counted by reason, with its +declared bytes, in the collector's own storage, so the counts survive restarts. +`GET /api/ingest` reports the queue and the counts, and the log names each +storage failure and, at most once a minute, other refusals. An exporter that gives +up after its retries may drop further telemetry it never sends; the collector +cannot count that. + +Spans are written to SQLite as each export arrives. The stored telemetry is kept +within `MOTEL_OTEL_RETENTION_HOURS` (default 168), `MOTEL_OTEL_MAX_SPANS` (default +1,000,000) and `MOTEL_OTEL_MAX_DB_SIZE_MB` (default 1024), whichever is reached +first, by evicting the oldest completed traces. The size bound is enforced after +every export, so the database file exceeds it by at most one export. Age and count +are enforced by retention passes every `MOTEL_OTEL_RETENTION_INTERVAL_SECONDS`, +each spending at most `MOTEL_OTEL_RETENTION_PASS_BUDGET_MS` (default 500) of +writer time. workerd owns the SQLite write-ahead log beside the file; its size is +not configurable from the worker. A busy server reaches the size bound long +before the others: at about 3 KiB per span, 1 GiB holds about 350,000 spans. + `bun run workerd:test` runs the built worker as a real process. It verifies HTTP trace/log ingestion, searches, seven-day retrieval, alarm retention, malformed payload refusal, and retained data after SIGKILL. Run `bun run workerd:build`