diff --git a/packages/backend/src/services/errors/ErrorsService.test.ts b/packages/backend/src/services/errors/ErrorsService.test.ts index 691a4a817..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,6 +164,7 @@ const makeWarehouseStub = ( scanRows: () => ReadonlyArray> = () => [], onScan?: () => void, fingerprintRows?: () => ReadonlyArray>, + nextActivityRows?: () => ReadonlyArray> | WarehouseQueryError, ): WarehouseQueryServiceApi => ({ query: () => Effect.die(new Error("unexpected warehouse query")), rawSqlQuery: () => Effect.succeed([]), @@ -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") { @@ -209,6 +218,7 @@ const makeErrorsLayer = ( edgeBackend?: ReturnType, fingerprintRows?: () => ReadonlyArray>, dispatcher?: (typeof NotificationDispatcher)["Service"], + nextActivityRows?: () => ReadonlyArray> | WarehouseQueryError, ) => { const testDb = createTestDb(createdDbs) const envLive = Env.layer.pipe(Layer.provide(testConfig())) @@ -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, @@ -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> = [] + 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* () { diff --git a/packages/backend/src/services/errors/ErrorsService.ts b/packages/backend/src/services/errors/ErrorsService.ts index ae934537d..d58470363 100644 --- a/packages/backend/src/services/errors/ErrorsService.ts +++ b/packages/backend/src/services/errors/ErrorsService.ts @@ -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, }) diff --git a/packages/query-engine/src/__sql_baseline__/catalog.sql b/packages/query-engine/src/__sql_baseline__/catalog.sql index 2e32fb7d3..3fcd18995 100644 --- a/packages/query-engine/src/__sql_baseline__/catalog.sql +++ b/packages/query-engine/src/__sql_baseline__/catalog.sql @@ -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, 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..a2ce8e880 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,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) diff --git a/packages/query-engine/src/ch/queries/errors.ts b/packages/query-engine/src/ch/queries/errors.ts index 47cd9ccf8..04999c136 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; 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