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
2 changes: 1 addition & 1 deletion docs/warehouse-rollups.md
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@ DDL snapshot, `local-schema.sql` for the embedded chDB engine, and the Rust inse

**A query must read the coarsest tier that can answer it.** The routing guards
(`canUseAnnualServiceOverview`, `canUseTracesAggregatesMv`, `canUseServiceOverviewMv`,
`canUseLogsAggregatesHourly`) exist to enforce that, and each one names the tier it unlocks.
`canUseLogsAggregatesHourly`, `canUseTraceFacetsRollup`) exist to enforce that, and each one names the tier it unlocks.

Rollup routes union a **raw edge** with a **rollup interior**: the rollup answers whole
buckets, and the raw table covers the partial buckets at each end of the window. Getting the
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,12 @@ import { TestClock } from "effect/testing"
import { Deferred, Effect, Exit, Fiber, Option, Schema } from "effect"
import { strict as nodeAssert } from "node:assert"
import { MetricName, OrgId, ServiceName, UserId } from "@maple/domain"
import { RawSqlValidationError, WarehouseUpstreamError } from "@maple/domain/http"
import {
RawSqlValidationError,
WarehouseConfigError,
WarehouseQuotaExceededError,
WarehouseUpstreamError,
} from "@maple/domain/http"
import {
baselineWarehouseCapabilities,
type QueryEngineEvaluateRequest,
Expand Down Expand Up @@ -293,6 +298,80 @@ describe("makeQueryEngineExecute", () => {
}),
)

// The sidebar's two reads: a failure on the first one is what either returns.
const traceSidebarStub = (failRollup: () => unknown) => {
const queries: Array<string> = []
const execute = makeQueryEngineExecute(
makeTinybirdStub({
sqlQuery: (_tenant, sql) => {
queries.push(sql)
return sql.includes("trace_facets_hourly")
? Effect.fail(failRollup() as WarehouseConfigError)
: Effect.succeed(
sql.includes("facetType")
? [{ name: "api", count: 3, facetType: "service" }]
: [
{
minDurationMs: 1,
maxDurationMs: 9,
p50DurationMs: 4,
p95DurationMs: 8,
},
],
)
},
}),
)
return { queries, execute }
}
const sidebarRequest = (kind: "facets" | "stats") => ({
startTime: "2026-01-01 00:00:00",
endTime: "2026-01-08 00:00:00",
query: { kind, source: "traces" as const },
})

it.effect("reads trace_list_mv when the cluster lacks trace_facets_hourly", () =>
Effect.gen(function* () {
for (const kind of ["facets", "stats"] as const) {
const { queries, execute } = traceSidebarStub(
() =>
new WarehouseConfigError({
message: "Unknown table expression identifier 'trace_facets_hourly'",
pipeName: "tracesFacets",
clickhouseType: "UNKNOWN_TABLE",
}),
)
const response = yield* execute(tenant, sidebarRequest(kind))

assert.strictEqual(queries.length, 2, kind)
assert.ok(queries[0]!.includes("trace_facets_hourly"), kind)
assert.ok(!queries[1]!.includes("trace_facets_hourly"), kind)
assert.strictEqual(response.result.kind, kind)
}
}),
)

it.effect("surfaces a rollup read that failed for another reason instead of rereading raw", () =>
Effect.gen(function* () {
const { queries, execute } = traceSidebarStub(
() =>
new WarehouseQuotaExceededError({
message: "Timeout exceeded while reading from table default.trace_facets_hourly",
pipeName: "tracesFacets",
setting: "max_execution_time",
}),
)
const exit = yield* Effect.exit(execute(tenant, sidebarRequest("facets")))

assert.isTrue(Exit.isFailure(exit))
assert.strictEqual(queries.length, 1)
assert.strictEqual(
(getFailure(exit) as { _tag?: string } | undefined)?._tag,
"@maple/http/errors/WarehouseQuotaExceededError",
)
}),
)

it.effect("fills missing buckets while preserving existing traces values", () =>
Effect.gen(function* () {
const execute = makeQueryEngineExecute(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,8 @@

import { afterAll, assert, beforeAll, describe, it } from "@effect/vitest"
import { migrations, renderStatementFull } from "@maple/domain/clickhouse"
import * as CH from "@maple/query-engine/ch"
import { normalizeSqlForClickHouseClient } from "@maple/query-engine/execution"
import {
applyRealMigrations,
clickhouseE2eEnabled,
Expand Down Expand Up @@ -53,7 +55,15 @@ const BATCH_1: ReadonlyArray<SeedSpan> = [
// Not a root span: never reaches trace_list_mv, so never the rollup.
{ traceId: "t1", ms: BASE_MS + 61_000, service: "api", durationNs: 1_000, parentSpanId: "root-t1" },
]
const BATCH_2: ReadonlyArray<SeedSpan> = [root("t5", BASE_MS + 240_000, "api", 3_000_000)]
const BATCH_2: ReadonlyArray<SeedSpan> = [
root("t5", BASE_MS + 240_000, "api", 3_000_000),
// The longest span sits in a whole hour, so only the hourly tier can report it.
root("t6", BASE_MS + HOUR_MS + 120_000, "api", 80_000_000),
root("t7", BASE_MS + 2 * HOUR_MS + 300_000, "api", 4_000_000),
root("t8", BASE_MS + 2 * HOUR_MS + 900_000, "api", 6_000_000),
// The only failing staging root inside a whole hour: proves filters reach the hourly tier.
root("t9", BASE_MS + HOUR_MS + 180_000, "worker", 30_000_000, { status: "Error", env: "staging" }),
]

const insert = async (spans: ReadonlyArray<SeedSpan>): Promise<void> => {
const rows = spans
Expand Down Expand Up @@ -114,7 +124,7 @@ describe.skipIf(!clickhouseE2eEnabled)("trace_facets_hourly materialization", ()
const expected = await fromTraceList()
assert.strictEqual(
expected.reduce((total, row) => total + Number(row.traces), 0),
5,
9,
)
assert.deepStrictEqual(await fromRollup(), expected)

Expand All @@ -125,4 +135,42 @@ describe.skipIf(!clickhouseE2eEnabled)("trace_facets_hourly materialization", ()
await clickhouseExec(renderStatementFull(backfill, database), database)
assert.deepStrictEqual(await fromRollup(), expected)
})

// Starts mid-hour after t1 and ends mid-hour between t7 and t8, so both raw
// edges hold rows in and out of the window; the aligned window has an empty
// raw tier, whose zero extremes must not reach the result.
it("answers the sidebar identically from the splice and from trace_list_mv alone", async () => {
const run = (sql: string) => runJson(normalizeSqlForClickHouseClient(sql))
const byFacet = (rows: ReadonlyArray<Record<string, unknown>>) =>
[...rows].sort((a, b) => `${a.facetType}:${a.name}`.localeCompare(`${b.facetType}:${b.name}`))

for (const [filters, longestMs] of [
[{}, 80],
[{ hasError: true, deploymentEnvs: ["staging"] }, 30],
] as const) {
for (const [startMs, endMs] of [
[BASE_MS + 90_000, BASE_MS + 2 * HOUR_MS + 600_000],
[BASE_MS + HOUR_MS, BASE_MS + 2 * HOUR_MS],
] as const) {
const window = { orgId: ORG_ID, startTime: chDateTime(startMs), endTime: chDateTime(endMs) }
const facets = (rawOnly: boolean) =>
CH.compileUnionUnsafe(CH.tracesFacetsQuery({ ...filters, rawOnly }), window).sql
const stats = (rawOnly: boolean) =>
CH.compileUnsafe(CH.tracesDurationStatsQuery({ ...filters, rawOnly }), window).sql
assert.include(facets(false), "trace_facets_hourly")
assert.notInclude(facets(true), "trace_facets_hourly")
assert.notInclude(stats(true), "trace_facets_hourly")

assert.deepStrictEqual(byFacet(await run(facets(false))), byFacet(await run(facets(true))))
const [spliced] = await run(stats(false))
const [raw] = await run(stats(true))
assert.strictEqual(spliced!.minDurationMs, raw!.minDurationMs)
assert.strictEqual(spliced!.maxDurationMs, longestMs)
assert.strictEqual(raw!.maxDurationMs, longestMs)
for (const key of ["p50DurationMs", "p95DurationMs"]) {
assert.closeTo(Number(spliced![key]), Number(raw![key]), Number(raw![key]) * 0.01)
}
}
}
})
})
Loading
Loading