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
81 changes: 78 additions & 3 deletions src/workerd.integration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16,11 +16,10 @@ const freePort = () =>
})
})

/** Exercise the built collector over HTTP, using a real workerd process and persistent disk. */
test("workerd ingests, searches, retains and recovers telemetry after process loss", async () => {
/** Write a collector config over a fresh data directory, and start workerd from it. */
const prepareCollector = async () => {
const directory = await mkdtemp(path.join(tmpdir(), "motel-workerd-"))
const port = await freePort()
const origin = `http://127.0.0.1:${port}`
await mkdir(path.join(directory, "data"))
await copyFile(path.join(root, "dist/workerd/motel.mjs"), path.join(directory, "motel.mjs"))
const base = await readFile(path.join(root, "workerd/motel.capnp"), "utf8")
Expand All @@ -40,6 +39,12 @@ test("workerd ingests, searches, retains and recovers telemetry after process lo
stdout: Bun.file(path.join(directory, "server.log")),
stderr: Bun.file(path.join(directory, "server.log")),
})
return { directory, origin: `http://127.0.0.1:${port}`, start }
}

/** Exercise the built collector over HTTP, using a real workerd process and persistent disk. */
test("workerd ingests, searches, retains and recovers telemetry after process loss", async () => {
const { directory, origin, start } = await prepareCollector()
let server = start()
const waitFor = async (predicate: () => Promise<boolean>) => {
for (let i = 0; i < 150; i++) {
Expand Down Expand Up @@ -145,3 +150,73 @@ test("workerd ingests, searches, retains and recovers telemetry after process lo
await rm(directory, { recursive: true, force: true })
}
}, 30000)

/**
* workerd binds at most 100 parameters per SQL statement. Exports, spans and query results larger
* than one statement's worth must still be stored and found in full.
*/
test("workerd stores and searches exports larger than one statement's parameters", async () => {
const { directory, origin, start } = await prepareCollector()
const server = start()
const nano = (ms: number) => String(BigInt(ms) * 1000000n)
const hex = (value: number, length: number) => value.toString(16).padStart(length, "0")
const attribute = (key: string, value: string) => ({ key, value: { stringValue: value } })
const post = async (spans: unknown[]) => {
const response = await fetch(origin + "/v1/traces", {
method: "POST",
headers: { "content-type": "application/json" },
body: JSON.stringify({
resourceSpans: [{ resource: { attributes: [attribute("service.name", "workerd-bounds")] }, scopeSpans: [{ spans }] }],
}),
})
expect(response.status).toBe(200)
return response.json()
}
const search = async (query: string) => (await (await fetch(`${origin}/api/spans/search?${query}`)).json()) as { data: unknown[] }
try {
for (let i = 0; i < 150; i++) {
if (await fetch(origin + "/api/health").then((response) => response.ok, () => false)) break
await Bun.sleep(100)
}
const now = Date.now()
// One export with many child spans, one of which carries more attributes than one statement binds.
const trace = hex(1, 32)
const root = hex(1, 16)
const children = Array.from({ length: 59 }, (_, index) => ({
traceId: trace,
spanId: hex(index + 2, 16),
parentSpanId: root,
name: index === 40 ? "build load" : "query step",
startTimeUnixNano: nano(now - 5),
endTimeUnixNano: nano(now - 4),
attributes:
index === 40
? [attribute("app.id", "app-a"), ...Array.from({ length: 40 }, (_, key) => attribute(`extra.${key}`, "value"))]
: [attribute("db.statement", "select 1")],
}))
expect(
await post([
{ traceId: trace, spanId: root, name: "request", startTimeUnixNano: nano(now - 10), endTimeUnixNano: nano(now), attributes: [] },
...children,
]),
).toEqual({ insertedSpans: 60 })
expect(await search("operation=build%20load&attr.app.id=app-a")).toMatchObject({ data: [{ span: { spanId: hex(42, 16) } }] })
expect(await search("operation=build%20load&attr.extra.39=value")).toMatchObject({ data: [{ span: { spanId: hex(42, 16) } }] })
expect((await search("operation=query%20step&limit=500")).data).toHaveLength(58)
// Results that span more traces than one statement binds.
const fanOut = Array.from({ length: 120 }, (_, index) => ({
traceId: hex(index + 100, 32),
spanId: hex(index + 100, 16),
name: "fan out",
startTimeUnixNano: nano(now - 3),
endTimeUnixNano: nano(now - 2),
attributes: [attribute("batch", "fan")],
}))
expect(await post(fanOut)).toEqual({ insertedSpans: 120 })
expect((await search("operation=fan%20out&attr.batch=fan&limit=500")).data).toHaveLength(120)
} finally {
server.kill()
await server.exited
await rm(directory, { recursive: true, force: true })
}
}, 30000)
Loading