From fa688dc2ea9ef2fc2b0dd1647c134f6b7b70c402 Mon Sep 17 00:00:00 2001 From: JeremyFunk Date: Wed, 7 Oct 2026 00:56:48 +0200 Subject: [PATCH 1/2] fix(errors): let a lagging error-tick cursor cross quiet stretches in one tick An org with no issue state is not scanned while it has no errors, so its cursor stays where it went quiet. When errors return, the tick resumed from that cursor at five minutes of event time per scan, and only on minutes the org looked active: an org quiet for a month needed weeks to reach the errors that brought it back, and produced no issues or notifications meanwhile. When a claimed window is empty and backlog remains, the tick now reads the next minute that has errors from the minute rollup and applies everything up to it (or to the cutoff) as one empty window. --- .../src/services/errors/ErrorsService.test.ts | 86 ++++++++++++++++++- .../src/services/errors/ErrorsService.ts | 34 ++++++++ .../src/__sql_baseline__/catalog.sql | 10 +++ .../query-engine/src/benchmark/builders.ts | 7 ++ packages/query-engine/src/ch/index.ts | 1 + .../src/ch/queries/errors.test.ts | 14 +++ .../query-engine/src/ch/queries/errors.ts | 20 +++++ 7 files changed, 171 insertions(+), 1 deletion(-) diff --git a/packages/backend/src/services/errors/ErrorsService.test.ts b/packages/backend/src/services/errors/ErrorsService.test.ts index 691a4a817..59640bd78 100644 --- a/packages/backend/src/services/errors/ErrorsService.test.ts +++ b/packages/backend/src/services/errors/ErrorsService.test.ts @@ -163,6 +163,7 @@ const makeWarehouseStub = ( scanRows: () => ReadonlyArray> = () => [], onScan?: () => void, fingerprintRows?: () => ReadonlyArray>, + nextActivityRows?: () => ReadonlyArray>, ): WarehouseQueryServiceApi => ({ query: () => Effect.die(new Error("unexpected warehouse query")), rawSqlQuery: () => Effect.succeed([]), @@ -186,6 +187,11 @@ 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") { + return Effect.orDie(compiledQueryOf(compiled).decodeRows(nextActivityRows?.() ?? [])) + } // 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") { @@ -209,6 +215,7 @@ const makeErrorsLayer = ( edgeBackend?: ReturnType, fingerprintRows?: () => ReadonlyArray>, dispatcher?: (typeof NotificationDispatcher)["Service"], + nextActivityRows?: () => ReadonlyArray>, ) => { const testDb = createTestDb(createdDbs) const envLive = Env.layer.pipe(Layer.provide(testConfig())) @@ -225,7 +232,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, @@ -1664,6 +1671,83 @@ 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 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> = [] + 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), bucketCount: 4 }], + ), + ), + ) + }) + it.effect("a window applied without the cursor claim commits nothing", () => { const rows = [scanRow()] return Effect.gen(function* () { diff --git a/packages/backend/src/services/errors/ErrorsService.ts b/packages/backend/src/services/errors/ErrorsService.ts index ae934537d..06af635cb 100644 --- a/packages/backend/src/services/errors/ErrorsService.ts +++ b/packages/backend/src/services/errors/ErrorsService.ts @@ -1205,9 +1205,43 @@ 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( + Effect.mapError(makePersistenceError), + Effect.tapError(() => releaseTickClaim(orgId, tickWindow.claimToken, nowMs)), + ) + const first = next[0] + const nextActivityMs = + first !== undefined && Number(first.bucketCount) > 0 + ? parseWarehouseDateTime(String(first.nextMinute)) + : cutoffMs + if (nextActivityMs > windowEndMs) { + windowEndMs = Math.min(nextActivityMs, cutoffMs) + fastForwarded = true + } + } yield* Effect.annotateCurrentSpan({ windowEndMs, windowSplits: splits, + windowFastForwarded: fastForwarded, scanFingerprints: issuesRaw.length, }) diff --git a/packages/query-engine/src/__sql_baseline__/catalog.sql b/packages/query-engine/src/__sql_baseline__/catalog.sql index 2e32fb7d3..fd05e52e7 100644 --- a/packages/query-engine/src/__sql_baseline__/catalog.sql +++ b/packages/query-engine/src/__sql_baseline__/catalog.sql @@ -751,6 +751,16 @@ SELECT GROUP BY fingerprintHash FORMAT JSON +-- builder:errors:errorTickNextActivityQuery:lagging-cursor [2d9a9812] +SELECT + min(error_fingerprints_minutely.Minute) AS nextMinute, + count() AS bucketCount + 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' + FORMAT JSON + -- builder:errors:spanDetailQuery:default [02aceae6] SELECT trace_detail_spans.TraceId AS traceId, diff --git a/packages/query-engine/src/benchmark/builders.ts b/packages/query-engine/src/benchmark/builders.ts index 78d4a8109..2e98f9d83 100644 --- a/packages/query-engine/src/benchmark/builders.ts +++ b/packages/query-engine/src/benchmark/builders.ts @@ -802,6 +802,13 @@ export const builderFixtures: ReadonlyArray = [ 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", diff --git a/packages/query-engine/src/ch/index.ts b/packages/query-engine/src/ch/index.ts index c28c1b193..61669f492 100644 --- a/packages/query-engine/src/ch/index.ts +++ b/packages/query-engine/src/ch/index.ts @@ -355,6 +355,7 @@ export { errorIssuesQuery, errorTickBootstrapIssuesQuery, errorTickIssuesQuery, + errorTickNextActivityQuery, errorFingerprintsQuery, errorIssueTimeseriesQuery, errorIssueSampleTracesQuery, diff --git a/packages/query-engine/src/ch/queries/errors.test.ts b/packages/query-engine/src/ch/queries/errors.test.ts index c4fc62cee..bb50b7354 100644 --- a/packages/query-engine/src/ch/queries/errors.test.ts +++ b/packages/query-engine/src/ch/queries/errors.test.ts @@ -10,6 +10,7 @@ import { errorIssuesQuery, errorTickBootstrapIssuesQuery, errorTickIssuesQuery, + errorTickNextActivityQuery, errorFingerprintsQuery, tracesDurationStatsQuery, tracesFacetsQuery, @@ -451,6 +452,19 @@ describe("errorTickIssuesQuery", () => { }) }) +describe("errorTickNextActivityQuery", () => { + it("reads the earliest minute and bucket count 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("min(error_fingerprints_minutely.Minute) AS nextMinute") + expect(sql).not.toContain("GROUP BY") + }) +}) + describe("errorTickBootstrapIssuesQuery", () => { it("bootstraps once from raw events without a truncating limit", () => { const { sql } = compileUnsafe(errorTickBootstrapIssuesQuery(), baseParams) diff --git a/packages/query-engine/src/ch/queries/errors.ts b/packages/query-engine/src/ch/queries/errors.ts index 47cd9ccf8..532ed547e 100644 --- a/packages/query-engine/src/ch/queries/errors.ts +++ b/packages/query-engine/src/ch/queries/errors.ts @@ -1163,6 +1163,26 @@ export function errorTickIssuesQuery() { .format("JSON") } +/** + * The earliest minute bucket holding errors in a window, and how many buckets + * the window holds (zero means `nextMinute` is the epoch, not a real minute). + * 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. + */ +export function errorTickNextActivityQuery() { + return from(ErrorFingerprintsMinutely) + .select(($) => ({ + nextMinute: CH.min_($.Minute), + bucketCount: CH.count(), + })) + .where(($) => [ + $.OrgId.eq(orgIdParam), + $.Minute.gte(param.dateTimeSeconds("startTime")), + $.Minute.lt(param.dateTimeSeconds("endTime")), + ]) + .format("JSON") +} + /** * One-time cursor bootstrap against the existing per-occurrence projection. * Incremental materialized views do not backfill historical rows, so a newly From 84365d96659876c85a1f69acdd8689f68bf54d02 Mon Sep 17 00:00:00 2001 From: JeremyFunk Date: Wed, 7 Oct 2026 01:01:24 +0200 Subject: [PATCH 2/2] fix(errors): first-match look-ahead query, fall back to the claimed window when it fails --- .../src/services/errors/ErrorsService.test.ts | 45 +++++++++++++++++-- .../src/services/errors/ErrorsService.ts | 20 ++++++--- .../src/__sql_baseline__/catalog.sql | 7 +-- .../src/ch/queries/errors.test.ts | 7 +-- .../query-engine/src/ch/queries/errors.ts | 16 +++---- 5 files changed, 71 insertions(+), 24 deletions(-) diff --git a/packages/backend/src/services/errors/ErrorsService.test.ts b/packages/backend/src/services/errors/ErrorsService.test.ts index 59640bd78..9ca0fc42f 100644 --- a/packages/backend/src/services/errors/ErrorsService.test.ts +++ b/packages/backend/src/services/errors/ErrorsService.test.ts @@ -10,6 +10,7 @@ import { IssueSeverityListCursor, OrgId, UserId, + WarehouseQueryError, } from "@maple/domain/http" import { ActorId, @@ -163,7 +164,7 @@ const makeWarehouseStub = ( scanRows: () => ReadonlyArray> = () => [], onScan?: () => void, fingerprintRows?: () => ReadonlyArray>, - nextActivityRows?: () => ReadonlyArray>, + nextActivityRows?: () => ReadonlyArray> | WarehouseQueryError, ): WarehouseQueryServiceApi => ({ query: () => Effect.die(new Error("unexpected warehouse query")), rawSqlQuery: () => Effect.succeed([]), @@ -190,7 +191,9 @@ const makeWarehouseStub = ( // 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") { - return Effect.orDie(compiledQueryOf(compiled).decodeRows(nextActivityRows?.() ?? [])) + 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. @@ -215,7 +218,7 @@ const makeErrorsLayer = ( edgeBackend?: ReturnType, fingerprintRows?: () => ReadonlyArray>, dispatcher?: (typeof NotificationDispatcher)["Service"], - nextActivityRows?: () => ReadonlyArray>, + nextActivityRows?: () => ReadonlyArray> | WarehouseQueryError, ) => { const testDb = createTestDb(createdDbs) const envLive = Env.layer.pipe(Layer.provide(testConfig())) @@ -1702,6 +1705,40 @@ describe("ErrorsService.runTick", () => { ) }) + 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. @@ -1742,7 +1779,7 @@ describe("ErrorsService.runTick", () => { undefined, undefined, undefined, - () => [{ nextMinute: formatWarehouseDateTime(resumeMs), bucketCount: 4 }], + () => [{ nextMinute: formatWarehouseDateTime(resumeMs) }], ), ), ) diff --git a/packages/backend/src/services/errors/ErrorsService.ts b/packages/backend/src/services/errors/ErrorsService.ts index 06af635cb..d58470363 100644 --- a/packages/backend/src/services/errors/ErrorsService.ts +++ b/packages/backend/src/services/errors/ErrorsService.ts @@ -1225,14 +1225,22 @@ const make: Effect.Effect< { profile: "aggregation", context: "errorTickNextActivity" }, ) .pipe( - Effect.mapError(makePersistenceError), - Effect.tapError(() => releaseTickClaim(orgId, tickWindow.claimToken, nowMs)), + // 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 first = next[0] const nextActivityMs = - first !== undefined && Number(first.bucketCount) > 0 - ? parseWarehouseDateTime(String(first.nextMinute)) - : cutoffMs + next === undefined + ? windowEndMs + : next[0] === undefined + ? cutoffMs + : parseWarehouseDateTime(String(next[0].nextMinute)) if (nextActivityMs > windowEndMs) { windowEndMs = Math.min(nextActivityMs, cutoffMs) fastForwarded = true diff --git a/packages/query-engine/src/__sql_baseline__/catalog.sql b/packages/query-engine/src/__sql_baseline__/catalog.sql index fd05e52e7..3fcd18995 100644 --- a/packages/query-engine/src/__sql_baseline__/catalog.sql +++ b/packages/query-engine/src/__sql_baseline__/catalog.sql @@ -751,14 +751,15 @@ SELECT GROUP BY fingerprintHash FORMAT JSON --- builder:errors:errorTickNextActivityQuery:lagging-cursor [2d9a9812] +-- builder:errors:errorTickNextActivityQuery:lagging-cursor [d115967d] SELECT - min(error_fingerprints_minutely.Minute) AS nextMinute, - count() AS bucketCount + 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] diff --git a/packages/query-engine/src/ch/queries/errors.test.ts b/packages/query-engine/src/ch/queries/errors.test.ts index bb50b7354..a2ce8e880 100644 --- a/packages/query-engine/src/ch/queries/errors.test.ts +++ b/packages/query-engine/src/ch/queries/errors.test.ts @@ -453,15 +453,16 @@ describe("errorTickIssuesQuery", () => { }) describe("errorTickNextActivityQuery", () => { - it("reads the earliest minute and bucket count from the minute rollup", () => { + 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("min(error_fingerprints_minutely.Minute) AS nextMinute") - expect(sql).not.toContain("GROUP BY") + expect(sql).toContain("error_fingerprints_minutely.Minute AS nextMinute") + expect(sql).toContain("ORDER BY nextMinute ASC") + expect(sql).toContain("LIMIT 1") }) }) diff --git a/packages/query-engine/src/ch/queries/errors.ts b/packages/query-engine/src/ch/queries/errors.ts index 532ed547e..04999c136 100644 --- a/packages/query-engine/src/ch/queries/errors.ts +++ b/packages/query-engine/src/ch/queries/errors.ts @@ -1164,22 +1164,22 @@ export function errorTickIssuesQuery() { } /** - * The earliest minute bucket holding errors in a window, and how many buckets - * the window holds (zero means `nextMinute` is the epoch, not a real minute). - * 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. + * 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: CH.min_($.Minute), - bucketCount: CH.count(), - })) + .select(($) => ({ nextMinute: $.Minute })) .where(($) => [ $.OrgId.eq(orgIdParam), $.Minute.gte(param.dateTimeSeconds("startTime")), $.Minute.lt(param.dateTimeSeconds("endTime")), ]) + .orderBy(["nextMinute", "asc"]) + .limit(1) .format("JSON") }