Skip to content
Open
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
123 changes: 122 additions & 1 deletion packages/backend/src/services/errors/ErrorsService.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import {
IssueSeverityListCursor,
OrgId,
UserId,
WarehouseQueryError,
} from "@maple/domain/http"
import {
ActorId,
Expand Down Expand Up @@ -163,6 +164,7 @@ const makeWarehouseStub = (
scanRows: () => ReadonlyArray<Record<string, unknown>> = () => [],
onScan?: () => void,
fingerprintRows?: () => ReadonlyArray<Record<string, unknown>>,
nextActivityRows?: () => ReadonlyArray<Record<string, unknown>> | WarehouseQueryError,
): WarehouseQueryServiceApi => ({
query: () => Effect.die(new Error("unexpected warehouse query")),
rawSqlQuery: () => Effect.succeed([]),
Expand All @@ -186,6 +188,13 @@ const makeWarehouseStub = (
if (options?.context === "errorIssueEnvFingerprints") {
return Effect.orDie(compiledQueryOf(compiled).decodeRows(fingerprintRows?.() ?? []))
}
// What lies behind a lagging cursor's empty window. Empty by default: no
// later errors, so the cursor may go straight to the cutoff.
if (options?.context === "errorTickNextActivity") {
const rows = nextActivityRows?.() ?? []
if (rows instanceof WarehouseQueryError) return Effect.fail(rows)
return Effect.orDie(compiledQueryOf(compiled).decodeRows(rows))
}
// Active-org discovery reads the same data the scan does, so model that
// consistency: surface the org iff it currently has error rows.
if (options?.context === "errorActiveOrgsDiscovery") {
Expand All @@ -209,6 +218,7 @@ const makeErrorsLayer = (
edgeBackend?: ReturnType<typeof makeMemoryBackend>,
fingerprintRows?: () => ReadonlyArray<Record<string, unknown>>,
dispatcher?: (typeof NotificationDispatcher)["Service"],
nextActivityRows?: () => ReadonlyArray<Record<string, unknown>> | WarehouseQueryError,
) => {
const testDb = createTestDb(createdDbs)
const envLive = Env.layer.pipe(Layer.provide(testConfig()))
Expand All @@ -225,7 +235,7 @@ const makeErrorsLayer = (
const errorPolicyLive = ErrorPolicyService.layer.pipe(Layer.provide(databaseLive))
const warehouseLive = Layer.succeed(
WarehouseQueryService,
makeWarehouseStub(scanRows, onScan, fingerprintRows),
makeWarehouseStub(scanRows, onScan, fingerprintRows, nextActivityRows),
)
const errorIssueReadModelsLive = Layer.effect(
ErrorIssueReadModelsService,
Expand Down Expand Up @@ -1664,6 +1674,117 @@ describe("ErrorsService.runTick", () => {
}).pipe(Effect.provide(makeErrorsLayer(() => burst)))
})

it.effect("a lagging cursor crosses a quiet stretch in one tick", () => {
let issueScans = 0
return Effect.gen(function* () {
const errors = yield* ErrorsService
const database = yield* Database
yield* TestClock.setTime(TICK_MS)
yield* seedIssue(asIssueId(randomUUID()))
yield* errors.runTick()

// Three days with no ticks and no errors: 864 five-minute windows.
const laterMs = TICK_MS + 3 * 24 * 60 * 60_000
yield* TestClock.setTime(laterMs)
yield* errors.runTick()

assert.strictEqual(issueScans, 2)
const cursor = yield* database.execute((db) =>
db.select().from(errorTickStates).where(eq(errorTickStates.orgId, ORG)),
)
assert.strictEqual(cursor[0]?.processedThrough.getTime(), laterMs - 60_000)
}).pipe(
Effect.provide(
makeErrorsLayer(
() => [],
() => {
issueScans += 1
},
),
),
)
})

it.effect("a lagging cursor still advances one window when it cannot look ahead", () => {
return Effect.gen(function* () {
const errors = yield* ErrorsService
const database = yield* Database
yield* TestClock.setTime(TICK_MS)
yield* seedIssue(asIssueId(randomUUID()))
yield* errors.runTick()

yield* TestClock.setTime(TICK_MS + 3 * 24 * 60 * 60_000)
yield* errors.runTick()

const cursor = yield* database.execute((db) =>
db.select().from(errorTickStates).where(eq(errorTickStates.orgId, ORG)),
)
// The claimed five-minute window commits; nothing is skipped.
assert.strictEqual(cursor[0]?.processedThrough.getTime(), TICK_MS - 60_000 + 5 * 60_000)
}).pipe(
Effect.provide(
makeErrorsLayer(
() => [],
undefined,
undefined,
undefined,
undefined,
() =>
new WarehouseQueryError({
message: "Query execution timed out after 30 seconds",
pipeName: "errorTickNextActivity",
}),
),
),
)
})

it.effect("a lagging cursor stops at the next minute that has errors", () => {
const laterMs = TICK_MS + 3 * 24 * 60 * 60_000
// Errors resume twenty minutes before the tick that finds the backlog.
const resumeMs = laterMs - 20 * 60_000
let rows: ReadonlyArray<Record<string, unknown>> = []
return Effect.gen(function* () {
const errors = yield* ErrorsService
const database = yield* Database
const cursorMs = () =>
database
.execute((db) => db.select().from(errorTickStates).where(eq(errorTickStates.orgId, ORG)))
.pipe(Effect.map((states) => states[0]?.processedThrough.getTime()))
yield* TestClock.setTime(TICK_MS)
yield* seedIssue(asIssueId(randomUUID()))
yield* errors.runTick()

yield* TestClock.setTime(laterMs)
yield* errors.runTick()
// Not past the errors: the window that holds them is still to be applied.
assert.strictEqual(yield* cursorMs(), resumeMs)
assert.lengthOf(yield* loadIssuesByFingerprint(SCAN_FINGERPRINT), 0)

rows = [
scanRow({
firstSeen: formatWarehouseDateTime(resumeMs + 1_000),
lastSeen: formatWarehouseDateTime(resumeMs + 2_000),
}),
]
yield* TestClock.setTime(laterMs + 60_000)
yield* errors.runTick()
assert.strictEqual(yield* cursorMs(), resumeMs + 5 * 60_000)
assert.lengthOf(yield* loadIssuesByFingerprint(SCAN_FINGERPRINT), 1)
}).pipe(
Effect.provide(
makeErrorsLayer(
() => rows,
undefined,
undefined,
undefined,
undefined,
() => [{ nextMinute: formatWarehouseDateTime(resumeMs) }],
),
),
)
})

it.effect("a window applied without the cursor claim commits nothing", () => {
const rows = [scanRow()]
return Effect.gen(function* () {
Expand Down
42 changes: 42 additions & 0 deletions packages/backend/src/services/errors/ErrorsService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1205,9 +1205,51 @@ const make: Effect.Effect<
}),
)
}
// A lagging cursor crossing a quiet stretch one window per cron gains four
// minutes a tick. An org with no issue state is not scanned while it has no
// errors, so its cursor stays where it went quiet, and without this the
// errors that bring it back a month later wait weeks behind empty windows.
// When the claimed window is empty and backlog remains, apply everything up
// to the next minute that has errors (or the cutoff) as one empty window.
// Not after a split: that window was narrowed because dense data follows it.
let fastForwarded = false
if (issuesRaw.length === 0 && splits === 0 && !tickWindow.isBootstrap && windowEndMs < cutoffMs) {
const next = yield* warehouse
.compiledQuery(
tenant,
CH.compile(CH.errorTickNextActivityQuery(), {
orgId,
startTime: formatWarehouseDateTime(windowEndMs),
endTime: formatWarehouseDateTime(cutoffMs),
}),
{ profile: "aggregation", context: "errorTickNextActivity" },
)
.pipe(
// Skipping ahead only saves time. If the read fails, apply the window
// as claimed: failing the tick here would hold the cursor on a window
// that could have committed.
Effect.catch((error) =>
Effect.logWarning("Error tick could not look past an empty window").pipe(
Effect.annotateLogs({ orgId, windowStartMs, windowEndMs, error: error.message }),
Effect.as(undefined),
),
),
)
const nextActivityMs =
next === undefined
? windowEndMs
: next[0] === undefined
? cutoffMs
: parseWarehouseDateTime(String(next[0].nextMinute))
if (nextActivityMs > windowEndMs) {
windowEndMs = Math.min(nextActivityMs, cutoffMs)
fastForwarded = true
}
}
yield* Effect.annotateCurrentSpan({
windowEndMs,
windowSplits: splits,
windowFastForwarded: fastForwarded,
scanFingerprints: issuesRaw.length,
})

Expand Down
11 changes: 11 additions & 0 deletions packages/query-engine/src/__sql_baseline__/catalog.sql
Original file line number Diff line number Diff line change
Expand Up @@ -751,6 +751,17 @@ SELECT
GROUP BY fingerprintHash
FORMAT JSON

-- builder:errors:errorTickNextActivityQuery:lagging-cursor [d115967d]
SELECT
error_fingerprints_minutely.Minute AS nextMinute
FROM error_fingerprints_minutely
WHERE error_fingerprints_minutely.OrgId = 'org_sql_catalog'
AND error_fingerprints_minutely.Minute >= '2026-01-01 10:30:00'
AND error_fingerprints_minutely.Minute < '2026-01-03 14:15:00'
ORDER BY nextMinute ASC
LIMIT 1
FORMAT JSON

-- builder:errors:spanDetailQuery:default [02aceae6]
SELECT
trace_detail_spans.TraceId AS traceId,
Expand Down
7 changes: 7 additions & 0 deletions packages/query-engine/src/benchmark/builders.ts
Original file line number Diff line number Diff line change
Expand Up @@ -802,6 +802,13 @@ export const builderFixtures: ReadonlyArray<BuilderFixture> = [
label: "cursor-window",
compile: () => CH.compileUnsafe(CH.errorTickIssuesQuery(), window),
},
{
// ErrorsService fast-forward: the next minute with errors behind a lagging cursor.
module: "errors",
name: "errorTickNextActivityQuery",
label: "lagging-cursor",
compile: () => CH.compileUnsafe(CH.errorTickNextActivityQuery(), window),
},
{
// ErrorsService one-time cursor bootstrap from retained canonical events.
module: "errors",
Expand Down
1 change: 1 addition & 0 deletions packages/query-engine/src/ch/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -355,6 +355,7 @@ export {
errorIssuesQuery,
errorTickBootstrapIssuesQuery,
errorTickIssuesQuery,
errorTickNextActivityQuery,
errorFingerprintsQuery,
errorIssueTimeseriesQuery,
errorIssueSampleTracesQuery,
Expand Down
15 changes: 15 additions & 0 deletions packages/query-engine/src/ch/queries/errors.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import {
errorIssuesQuery,
errorTickBootstrapIssuesQuery,
errorTickIssuesQuery,
errorTickNextActivityQuery,
errorFingerprintsQuery,
tracesDurationStatsQuery,
tracesFacetsQuery,
Expand Down Expand Up @@ -451,6 +452,20 @@ describe("errorTickIssuesQuery", () => {
})
})

describe("errorTickNextActivityQuery", () => {
it("reads the earliest minute with errors from the minute rollup", () => {
const { sql } = compileUnsafe(errorTickNextActivityQuery(), baseParams)

expect(sql).toContain("FROM error_fingerprints_minutely")
expect(sql).toContain("OrgId = 'org_1'")
expect(sql).toContain("Minute >= '2024-01-01 00:00:00'")
expect(sql).toContain("Minute < '2024-01-02 00:00:00'")
expect(sql).toContain("error_fingerprints_minutely.Minute AS nextMinute")
expect(sql).toContain("ORDER BY nextMinute ASC")
expect(sql).toContain("LIMIT 1")
})
})

describe("errorTickBootstrapIssuesQuery", () => {
it("bootstraps once from raw events without a truncating limit", () => {
const { sql } = compileUnsafe(errorTickBootstrapIssuesQuery(), baseParams)
Expand Down
20 changes: 20 additions & 0 deletions packages/query-engine/src/ch/queries/errors.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1163,6 +1163,26 @@ export function errorTickIssuesQuery() {
.format("JSON")
}

/**
* The earliest minute bucket holding errors in a window; no row when the window
* holds none. A lagging cursor whose claimed window is empty reads it to cross
* the rest of a quiet stretch in one step instead of one window per cron.
* Ordered by the rollup's sorting key and limited to one row, so it stops at
* the first match rather than reading the whole window.
*/
export function errorTickNextActivityQuery() {
return from(ErrorFingerprintsMinutely)
.select(($) => ({ nextMinute: $.Minute }))
.where(($) => [
$.OrgId.eq(orgIdParam),
$.Minute.gte(param.dateTimeSeconds("startTime")),
$.Minute.lt(param.dateTimeSeconds("endTime")),
])
.orderBy(["nextMinute", "asc"])
.limit(1)
.format("JSON")
}

/**
* One-time cursor bootstrap against the existing per-occurrence projection.
* Incremental materialized views do not backfill historical rows, so a newly
Expand Down
Loading