diff --git a/packages/backend/src/services/errors/ErrorsService.test.ts b/packages/backend/src/services/errors/ErrorsService.test.ts index ed0b4aabd..691a4a817 100644 --- a/packages/backend/src/services/errors/ErrorsService.test.ts +++ b/packages/backend/src/services/errors/ErrorsService.test.ts @@ -1622,6 +1622,48 @@ describe("ErrorsService.runTick", () => { ) }) + it.effect("a window with more fingerprints than one statement can bind still commits", () => { + // Postgres caps a statement at 32,767 bind parameters. 3,100 new + // fingerprints is 37,200 on the candidate upsert alone, and a window that + // cannot commit never advances the cursor. + const burst = Array.from({ length: 3_100 }, (_, i) => scanRow({ fingerprintHash: `7${i}` })) + return Effect.gen(function* () { + const errors = yield* ErrorsService + const database = yield* Database + yield* TestClock.setTime(TICK_MS) + yield* seedIssue(asIssueId(randomUUID())) + + const result = yield* errors.runTick() + assert.strictEqual(result.issuesTouched, 3_100) + assert.strictEqual(result.incidentsOpened, 3_100) + + const cursor = yield* database.execute((db) => + db.select().from(errorTickStates).where(eq(errorTickStates.orgId, ORG)), + ) + assert.strictEqual(cursor[0]?.processedThrough.getTime(), TICK_MS - 60_000) + // Every chunk landed, not just the ones a counter reports: one `created` + // event per issue, and no promoted candidate left behind. + const created = yield* database.execute((db) => + db.select().from(errorIssueEvents).where(eq(errorIssueEvents.type, "created")), + ) + assert.lengthOf(created, 3_100) + assert.lengthOf(yield* database.execute((db) => db.select().from(errorFingerprintCandidates)), 0) + + // The same burst a minute later lands on the incidents it just opened. + yield* TestClock.setTime(TICK_MS + 60_000) + const ongoing = yield* errors.runTick() + assert.strictEqual(ongoing.issuesTouched, 3_100) + assert.strictEqual(ongoing.incidentsOpened, 0) + + // Thirty quiet minutes later every one of those incidents goes stale in + // the same window. + burst.length = 0 + yield* TestClock.setTime(TICK_MS + 33 * 60_000) + const quiet = yield* runTicksUntilCaughtUp() + assert.strictEqual(quiet.incidentsResolved, 3_100) + }).pipe(Effect.provide(makeErrorsLayer(() => burst))) + }) + 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/error-tick-persistence.ts b/packages/backend/src/services/errors/error-tick-persistence.ts index 579d17916..3a8d8f7b2 100644 --- a/packages/backend/src/services/errors/error-tick-persistence.ts +++ b/packages/backend/src/services/errors/error-tick-persistence.ts @@ -25,7 +25,7 @@ import type { } from "@maple/domain/http" import { FINGERPRINT_VERSION } from "@maple/domain/tinybird/fingerprint" import { and, eq, inArray, isNull, lt, sql } from "drizzle-orm" -import { Effect, Schema } from "effect" +import { Array as Arr, Effect, Schema } from "effect" export interface ErrorTickScanRow { readonly fingerprintHash: string @@ -120,6 +120,24 @@ const MAX_TRACKED_VERSIONS = 50 /** The transaction handle `persistErrorTickWindow` runs its statements on. */ type ErrorTickTransaction = MapleTx +/** + * Rows per statement for the window's bulk reads and writes. Postgres caps a + * statement at 32,767 bind parameters and the widest write here binds 26 per + * row, so a single statement per window stops committing past ~1,200 + * fingerprints. A steady-state window is far below this and still runs every + * write as one statement. + */ +const BULK_CHUNK_ROWS = 500 + +const chunks = (items: ReadonlyArray) => Arr.chunksOf(items, BULK_CHUNK_ROWS) + +/** Run a row-returning statement once per chunk and concatenate what came back. */ +const inChunks = ( + items: ReadonlyArray, + run: (chunk: Array) => Effect.Effect, E>, +): Effect.Effect, E> => + Effect.forEach(chunks(items), run).pipe(Effect.map((results) => results.flat())) + /** * Accumulate not-yet-promoted fingerprints and return the ones that have now * earned an Issue, carrying their ACCUMULATED totals rather than this window's. @@ -136,37 +154,38 @@ const promoteCandidates = ( Effect.gen(function* () { if (unknownRows.length === 0) return [] - const upserted = yield* tx - .insert(errorFingerprintCandidates) - .values( - unknownRows.map((row) => ({ - orgId: input.orgId, - fingerprintHash: row.fingerprintHash, - serviceName: row.serviceName, - exceptionType: row.exceptionType, - exceptionMessage: row.exceptionMessage, - errorLabel: row.errorLabel, - topFrame: row.topFrame, - serviceVersionsJson: mergeVersions([], row.serviceVersions), - occurrenceCount: row.count, - firstSeenAt: msToDate(row.firstSeenMs), - lastSeenAt: msToDate(row.lastSeenMs), - updatedAt: windowEnd, - })), - ) - .onConflictDoUpdate({ - target: [errorFingerprintCandidates.orgId, errorFingerprintCandidates.fingerprintHash], - set: { - serviceName: sql`excluded.service_name`, - exceptionType: sql`excluded.exception_type`, - exceptionMessage: sql`excluded.exception_message`, - errorLabel: sql`excluded.error_label`, - topFrame: sql`excluded.top_frame`, - // Union rather than overwrite: a candidate accumulates across ticks and - // the builds it was seen from seed the issue it is promoted into. The - // cap mirrors MAX_TRACKED_VERSIONS; which build is dropped past it is - // unordered, which is fine for a row that lives at most a day. - serviceVersionsJson: sql`( + const upserted = yield* inChunks(unknownRows, (chunk) => + tx + .insert(errorFingerprintCandidates) + .values( + chunk.map((row) => ({ + orgId: input.orgId, + fingerprintHash: row.fingerprintHash, + serviceName: row.serviceName, + exceptionType: row.exceptionType, + exceptionMessage: row.exceptionMessage, + errorLabel: row.errorLabel, + topFrame: row.topFrame, + serviceVersionsJson: mergeVersions([], row.serviceVersions), + occurrenceCount: row.count, + firstSeenAt: msToDate(row.firstSeenMs), + lastSeenAt: msToDate(row.lastSeenMs), + updatedAt: windowEnd, + })), + ) + .onConflictDoUpdate({ + target: [errorFingerprintCandidates.orgId, errorFingerprintCandidates.fingerprintHash], + set: { + serviceName: sql`excluded.service_name`, + exceptionType: sql`excluded.exception_type`, + exceptionMessage: sql`excluded.exception_message`, + errorLabel: sql`excluded.error_label`, + topFrame: sql`excluded.top_frame`, + // Union rather than overwrite: a candidate accumulates across ticks and + // the builds it was seen from seed the issue it is promoted into. The + // cap mirrors MAX_TRACKED_VERSIONS; which build is dropped past it is + // unordered, which is fine for a row that lives at most a day. + serviceVersionsJson: sql`( select coalesce(jsonb_agg(version), '[]'::jsonb) from ( select distinct value as version @@ -176,25 +195,31 @@ const promoteCandidates = ( limit ${sql.raw(String(MAX_TRACKED_VERSIONS))} ) as versions )`, - occurrenceCount: sql`${errorFingerprintCandidates.occurrenceCount} + excluded.occurrence_count`, - firstSeenAt: sql`least(${errorFingerprintCandidates.firstSeenAt}, excluded.first_seen_at)`, - lastSeenAt: sql`greatest(${errorFingerprintCandidates.lastSeenAt}, excluded.last_seen_at)`, - updatedAt: sql`excluded.updated_at`, - }, - }) - .returning() + occurrenceCount: sql`${errorFingerprintCandidates.occurrenceCount} + excluded.occurrence_count`, + firstSeenAt: sql`least(${errorFingerprintCandidates.firstSeenAt}, excluded.first_seen_at)`, + lastSeenAt: sql`greatest(${errorFingerprintCandidates.lastSeenAt}, excluded.last_seen_at)`, + updatedAt: sql`excluded.updated_at`, + }, + }) + .returning(), + ) const ready = upserted.filter((row) => row.occurrenceCount >= PROMOTION_MIN_OCCURRENCES) if (ready.length === 0) return [] - yield* tx.delete(errorFingerprintCandidates).where( - and( - eq(errorFingerprintCandidates.orgId, input.orgId), - inArray( - errorFingerprintCandidates.fingerprintHash, - ready.map((row) => row.fingerprintHash), + yield* Effect.forEach( + chunks(ready), + (chunk) => + tx.delete(errorFingerprintCandidates).where( + and( + eq(errorFingerprintCandidates.orgId, input.orgId), + inArray( + errorFingerprintCandidates.fingerprintHash, + chunk.map((row) => row.fingerprintHash), + ), + ), ), - ), + { discard: true }, ) return ready.map((row) => ({ @@ -376,8 +401,8 @@ const buildNotificationRows = (input: PersistErrorTickWindowInput, payload: Noti * commits. A transaction retry is safe because none of its generated IDs can * escape a rolled-back attempt. * - * Every write is set-based: the cost of a window is a fixed handful of - * statements rather than a few per fingerprint. That is what keeps a wide + * Every write is set-based: the cost of a window is a handful of statements + * per `BULK_CHUNK_ROWS` fingerprints rather than a few per fingerprint. That is what keeps a wide * catch-up window inside its lease — a per-row loop over Hyperdrive turns a * large window into a transaction that cannot finish before the lease TTL, * and a window that can never commit never advances the cursor either. @@ -415,18 +440,14 @@ export const persistErrorTickWindow = ( const rows = mergeScanRows(input.rows) const fingerprints = rows.map((row) => row.fingerprintHash) - const existingIssues = - fingerprints.length === 0 - ? [] - : yield* tx - .select() - .from(errorIssues) - .where( - and( - eq(errorIssues.orgId, input.orgId), - inArray(errorIssues.fingerprintHash, fingerprints), - ), - ) + const existingIssues = yield* inChunks(fingerprints, (chunk) => + tx + .select() + .from(errorIssues) + .where( + and(eq(errorIssues.orgId, input.orgId), inArray(errorIssues.fingerprintHash, chunk)), + ), + ) const issueByFingerprint = new Map(existingIssues.map((row) => [row.fingerprintHash, row])) // The cursor row lock makes this evaluator the only writer of error-kind @@ -507,69 +528,71 @@ export const persistErrorTickWindow = ( // One upsert covers new, ongoing, and regressed fingerprints. The SET // expressions all read the pre-update row, so the regression flip and the // counter accumulation stay consistent with each other. - const upserted = yield* tx - .insert(errorIssues) - .values( - applicable.map(({ row, prior }) => ({ - id: input.makeIssueId(), - orgId: input.orgId, - fingerprintHash: row.fingerprintHash, - fingerprintVersion: FINGERPRINT_VERSION, - serviceName: row.serviceName, - exceptionType: row.exceptionType, - exceptionMessage: row.exceptionMessage, - errorLabel: row.errorLabel, - topFrame: row.topFrame, - workflowState: "triage" as const, - priority: 3, - assignedActorId: null, - leaseHolderActorId: null, - leaseExpiresAt: null, - claimedAt: null, - notes: null, - firstSeenAt: msToDate(row.firstSeenMs), - lastSeenAt: msToDate(row.lastSeenMs), - occurrenceCount: row.count, - seenVersionsJson: mergeVersions( - prior?.seenVersionsJson ?? [], - row.serviceVersions, - ), - resolvedAt: null, - resolvedByActorId: null, - snoozeUntil: null, - archivedAt: null, - createdAt: windowEnd, - updatedAt: windowEnd, - })), - ) - .onConflictDoUpdate({ - target: [errorIssues.orgId, errorIssues.fingerprintHash], - set: { - fingerprintVersion: sql`excluded.fingerprint_version`, - serviceName: sql`excluded.service_name`, - exceptionType: sql`excluded.exception_type`, - exceptionMessage: sql`excluded.exception_message`, - errorLabel: sql`excluded.error_label`, - topFrame: sql`excluded.top_frame`, - firstSeenAt: sql`least(${errorIssues.firstSeenAt}, excluded.first_seen_at)`, - lastSeenAt: sql`greatest(${errorIssues.lastSeenAt}, excluded.last_seen_at)`, - occurrenceCount: sql`${errorIssues.occurrenceCount} + excluded.occurrence_count`, - // The merged set is computed against the pre-update row and arrives - // via the INSERT values, so this stays one statement. - seenVersionsJson: sql`excluded.seen_versions_json`, - // Reopening is NOT decided here. It used to be - // `case when workflow_state = 'done' then 'triage'`, which reopened a - // fixed issue on any occurrence at all — including one from a build - // that predates the fix. That needs the issue's resolved-build set, - // which SQL cannot reach from `excluded`, so `isRegression` decides in - // TypeScript and the targeted update below applies it. - updatedAt: sql`excluded.updated_at`, - }, - }) - .returning({ - id: errorIssues.id, - fingerprintHash: errorIssues.fingerprintHash, - }) + const upserted = yield* inChunks(applicable, (chunk) => + tx + .insert(errorIssues) + .values( + chunk.map(({ row, prior }) => ({ + id: input.makeIssueId(), + orgId: input.orgId, + fingerprintHash: row.fingerprintHash, + fingerprintVersion: FINGERPRINT_VERSION, + serviceName: row.serviceName, + exceptionType: row.exceptionType, + exceptionMessage: row.exceptionMessage, + errorLabel: row.errorLabel, + topFrame: row.topFrame, + workflowState: "triage" as const, + priority: 3, + assignedActorId: null, + leaseHolderActorId: null, + leaseExpiresAt: null, + claimedAt: null, + notes: null, + firstSeenAt: msToDate(row.firstSeenMs), + lastSeenAt: msToDate(row.lastSeenMs), + occurrenceCount: row.count, + seenVersionsJson: mergeVersions( + prior?.seenVersionsJson ?? [], + row.serviceVersions, + ), + resolvedAt: null, + resolvedByActorId: null, + snoozeUntil: null, + archivedAt: null, + createdAt: windowEnd, + updatedAt: windowEnd, + })), + ) + .onConflictDoUpdate({ + target: [errorIssues.orgId, errorIssues.fingerprintHash], + set: { + fingerprintVersion: sql`excluded.fingerprint_version`, + serviceName: sql`excluded.service_name`, + exceptionType: sql`excluded.exception_type`, + exceptionMessage: sql`excluded.exception_message`, + errorLabel: sql`excluded.error_label`, + topFrame: sql`excluded.top_frame`, + firstSeenAt: sql`least(${errorIssues.firstSeenAt}, excluded.first_seen_at)`, + lastSeenAt: sql`greatest(${errorIssues.lastSeenAt}, excluded.last_seen_at)`, + occurrenceCount: sql`${errorIssues.occurrenceCount} + excluded.occurrence_count`, + // The merged set is computed against the pre-update row and arrives + // via the INSERT values, so this stays one statement. + seenVersionsJson: sql`excluded.seen_versions_json`, + // Reopening is NOT decided here. It used to be + // `case when workflow_state = 'done' then 'triage'`, which reopened a + // fixed issue on any occurrence at all — including one from a build + // that predates the fix. That needs the issue's resolved-build set, + // which SQL cannot reach from `excluded`, so `isRegression` decides in + // TypeScript and the targeted update below applies it. + updatedAt: sql`excluded.updated_at`, + }, + }) + .returning({ + id: errorIssues.id, + fingerprintHash: errorIssues.fingerprintHash, + }), + ) const idByFingerprint = new Map(upserted.map((row) => [row.fingerprintHash, row.id])) @@ -631,34 +654,36 @@ export const persistErrorTickWindow = ( const regressedIds = observed .filter((entry) => entry.wasRegression) .map((entry) => entry.issueId) - if (regressedIds.length > 0) { - yield* tx - .update(errorIssues) - .set({ - workflowState: "regressed", - resolvedAt: null, - resolvedByActorId: null, - lastRegressedAt: windowEnd, - regressionCount: sql`${errorIssues.regressionCount} + 1`, - updatedAt: windowEnd, - }) - .where(and(eq(errorIssues.orgId, input.orgId), inArray(errorIssues.id, regressedIds))) - } + yield* Effect.forEach( + chunks(regressedIds), + (chunk) => + tx + .update(errorIssues) + .set({ + workflowState: "regressed", + resolvedAt: null, + resolvedByActorId: null, + lastRegressedAt: windowEnd, + regressionCount: sql`${errorIssues.regressionCount} + 1`, + updatedAt: windowEnd, + }) + .where(and(eq(errorIssues.orgId, input.orgId), inArray(errorIssues.id, chunk))), + { discard: true }, + ) } const observedIssueIds = observed.map((entry) => entry.issueId) - const states = - observedIssueIds.length === 0 - ? [] - : yield* tx - .select() - .from(errorIssueStates) - .where( - and( - eq(errorIssueStates.orgId, input.orgId), - inArray(errorIssueStates.issueId, observedIssueIds), - ), - ) + const states = yield* inChunks(observedIssueIds, (chunk) => + tx + .select() + .from(errorIssueStates) + .where( + and( + eq(errorIssueStates.orgId, input.orgId), + inArray(errorIssueStates.issueId, chunk), + ), + ), + ) const stateByIssue = new Map(states.map((row) => [row.issueId, row])) const refreshing: Array<{ @@ -684,15 +709,18 @@ export const persistErrorTickWindow = ( } } - if (refreshing.length > 0) { - const incidentValues = sql.join( - refreshing.map( - (item) => - sql`(${item.incidentId}::text, ${msToSqlTimestamp(item.row.lastSeenMs)}::timestamptz, ${item.row.count}::integer)`, - ), - sql`, `, - ) - yield* tx.execute(sql` + yield* Effect.forEach( + chunks(refreshing), + (chunk) => + Effect.gen(function* () { + const incidentValues = sql.join( + chunk.map( + (item) => + sql`(${item.incidentId}::text, ${msToSqlTimestamp(item.row.lastSeenMs)}::timestamptz, ${item.row.count}::integer)`, + ), + sql`, `, + ) + yield* tx.execute(sql` update ${errorIncidents} as t set last_triggered_at = greatest(t.last_triggered_at, v.last_seen), occurrence_count = t.occurrence_count + v.cnt, @@ -701,14 +729,14 @@ export const persistErrorTickWindow = ( where t.id = v.incident_id `) - const stateValues = sql.join( - refreshing.map( - (item) => - sql`(${item.issueId}::text, ${msToSqlTimestamp(item.row.lastSeenMs)}::timestamptz)`, - ), - sql`, `, - ) - yield* tx.execute(sql` + const stateValues = sql.join( + chunk.map( + (item) => + sql`(${item.issueId}::text, ${msToSqlTimestamp(item.row.lastSeenMs)}::timestamptz)`, + ), + sql`, `, + ) + yield* tx.execute(sql` update ${errorIssueStates} as t set last_observed_occurrence_at = greatest(t.last_observed_occurrence_at, v.last_seen), last_evaluated_at = ${windowEndSql}::timestamptz, @@ -716,7 +744,9 @@ export const persistErrorTickWindow = ( from (values ${stateValues}) as v(issue_id, last_seen) where t.org_id = ${input.orgId} and t.issue_id = v.issue_id `) - } + }), + { discard: true }, + ) const pendingTriages: PendingErrorTriage[] = [] let incidentsOpened = 0 @@ -725,32 +755,34 @@ export const persistErrorTickWindow = ( // The conditional upsert is the single-open-incident guard: a row whose // `open_incident_id` is already set fails the `setWhere` and is therefore // absent from RETURNING, so the claimed set is exactly what came back. - const claimed = yield* tx - .insert(errorIssueStates) - .values( - claiming.map((item) => ({ - orgId: input.orgId, - issueId: item.entry.issueId, - lastObservedOccurrenceAt: msToDate(item.entry.row.lastSeenMs), - lastEvaluatedAt: windowEnd, - openIncidentId: item.incidentId, - updatedAt: windowEnd, - })), - ) - .onConflictDoUpdate({ - target: [errorIssueStates.orgId, errorIssueStates.issueId], - set: { - lastObservedOccurrenceAt: sql`excluded.last_observed_occurrence_at`, - lastEvaluatedAt: sql`excluded.last_evaluated_at`, - openIncidentId: sql`excluded.open_incident_id`, - updatedAt: sql`excluded.updated_at`, - }, - setWhere: isNull(errorIssueStates.openIncidentId), - }) - .returning({ - issueId: errorIssueStates.issueId, - openIncidentId: errorIssueStates.openIncidentId, - }) + const claimed = yield* inChunks(claiming, (chunk) => + tx + .insert(errorIssueStates) + .values( + chunk.map((item) => ({ + orgId: input.orgId, + issueId: item.entry.issueId, + lastObservedOccurrenceAt: msToDate(item.entry.row.lastSeenMs), + lastEvaluatedAt: windowEnd, + openIncidentId: item.incidentId, + updatedAt: windowEnd, + })), + ) + .onConflictDoUpdate({ + target: [errorIssueStates.orgId, errorIssueStates.issueId], + set: { + lastObservedOccurrenceAt: sql`excluded.last_observed_occurrence_at`, + lastEvaluatedAt: sql`excluded.last_evaluated_at`, + openIncidentId: sql`excluded.open_incident_id`, + updatedAt: sql`excluded.updated_at`, + }, + setWhere: isNull(errorIssueStates.openIncidentId), + }) + .returning({ + issueId: errorIssueStates.issueId, + openIncidentId: errorIssueStates.openIncidentId, + }), + ) const claimedIncidentByIssue = new Map( claimed.map((row) => [row.issueId, row.openIncidentId]), @@ -760,20 +792,25 @@ export const persistErrorTickWindow = ( ) if (opened.length > 0) { - yield* tx.insert(errorIncidents).values( - opened.map((item) => ({ - id: item.incidentId, - orgId: input.orgId, - issueId: item.entry.issueId, - status: "open" as const, - reason: item.reason, - firstTriggeredAt: msToDate(item.entry.row.firstSeenMs), - lastTriggeredAt: msToDate(item.entry.row.lastSeenMs), - resolvedAt: null, - occurrenceCount: item.entry.row.count, - createdAt: windowEnd, - updatedAt: windowEnd, - })), + yield* Effect.forEach( + chunks(opened), + (chunk) => + tx.insert(errorIncidents).values( + chunk.map((item) => ({ + id: item.incidentId, + orgId: input.orgId, + issueId: item.entry.issueId, + status: "open" as const, + reason: item.reason, + firstTriggeredAt: msToDate(item.entry.row.firstSeenMs), + lastTriggeredAt: msToDate(item.entry.row.lastSeenMs), + resolvedAt: null, + occurrenceCount: item.entry.row.count, + createdAt: windowEnd, + updatedAt: windowEnd, + })), + ), + { discard: true }, ) incidentsOpened = opened.length @@ -817,56 +854,66 @@ export const persistErrorTickWindow = ( let incidentsResolved = 0 if (staleIncidents.length > 0) { - const staleIds = staleIncidents.map((incident) => incident.id) // `status = 'open'` in the predicate keeps the flip idempotent, and // RETURNING reports exactly which rows this transaction resolved. - const flipped = yield* tx - .update(errorIncidents) - .set({ status: "resolved", resolvedAt: windowEnd, updatedAt: windowEnd }) - .where( - and( - eq(errorIncidents.orgId, input.orgId), - inArray(errorIncidents.id, staleIds), - eq(errorIncidents.status, "open"), - ), - ) - .returning({ id: errorIncidents.id }) - const flippedIds = flipped.map((row) => row.id) - incidentsResolved = flippedIds.length - - if (flippedIds.length > 0) { - const resolvedIncidents = staleIncidents.filter((incident) => - flippedIds.includes(incident.id), - ) - // Match on `open_incident_id` too: a state pointing at some other - // incident must not be cleared by this one resolving. - yield* tx - .update(errorIssueStates) - .set({ openIncidentId: null, updatedAt: windowEnd }) + const flipped = yield* inChunks(staleIncidents, (chunk) => + tx + .update(errorIncidents) + .set({ status: "resolved", resolvedAt: windowEnd, updatedAt: windowEnd }) .where( and( - eq(errorIssueStates.orgId, input.orgId), + eq(errorIncidents.orgId, input.orgId), inArray( - errorIssueStates.issueId, - resolvedIncidents.map((incident) => incident.issueId), + errorIncidents.id, + chunk.map((incident) => incident.id), ), - inArray(errorIssueStates.openIncidentId, flippedIds), + eq(errorIncidents.status, "open"), ), ) + .returning({ id: errorIncidents.id }), + ) + const flippedIds = new Set(flipped.map((row) => row.id)) + incidentsResolved = flippedIds.size + + if (flippedIds.size > 0) { + const resolvedIncidents = staleIncidents.filter((incident) => flippedIds.has(incident.id)) + // Match on `open_incident_id` too: a state pointing at some other + // incident must not be cleared by this one resolving. An incident and + // its issue share a chunk, so chunking leaves the pairing intact. + yield* Effect.forEach( + chunks(resolvedIncidents), + (chunk) => + tx + .update(errorIssueStates) + .set({ openIncidentId: null, updatedAt: windowEnd }) + .where( + and( + eq(errorIssueStates.orgId, input.orgId), + inArray( + errorIssueStates.issueId, + chunk.map((incident) => incident.issueId), + ), + inArray( + errorIssueStates.openIncidentId, + chunk.map((incident) => incident.id), + ), + ), + ), + { discard: true }, + ) if (input.policy.enabled && input.policy.notifyOnResolve) { const staleIssueIds = [ ...new Set(resolvedIncidents.map((incident) => incident.issueId)), ] - const staleIssues = yield* tx - .select() - .from(errorIssues) - .where( - and( - eq(errorIssues.orgId, input.orgId), - inArray(errorIssues.id, staleIssueIds), + const staleIssues = yield* inChunks(staleIssueIds, (chunk) => + tx + .select() + .from(errorIssues) + .where( + and(eq(errorIssues.orgId, input.orgId), inArray(errorIssues.id, chunk)), ), - ) + ) const staleIssueById = new Map(staleIssues.map((issue) => [issue.id, issue])) for (const incident of resolvedIncidents) { const issue = staleIssueById.get(incident.issueId) @@ -888,18 +935,23 @@ export const persistErrorTickWindow = ( } } - if (events.length > 0) yield* tx.insert(errorIssueEvents).values(events) - if (notifications.length > 0) { - yield* tx - .insert(errorNotificationDeliveries) - .values(notifications) - .onConflictDoNothing({ - target: [ - errorNotificationDeliveries.deliveryKey, - errorNotificationDeliveries.destinationId, - ], - }) - } + yield* Effect.forEach(chunks(events), (chunk) => tx.insert(errorIssueEvents).values(chunk), { + discard: true, + }) + yield* Effect.forEach( + chunks(notifications), + (chunk) => + tx + .insert(errorNotificationDeliveries) + .values(chunk) + .onConflictDoNothing({ + target: [ + errorNotificationDeliveries.deliveryKey, + errorNotificationDeliveries.destinationId, + ], + }), + { discard: true }, + ) const advanced = yield* tx .update(errorTickStates)