diff --git a/docker-compose.dev.yml b/docker-compose.dev.yml index 29766b23..cc131d23 100644 --- a/docker-compose.dev.yml +++ b/docker-compose.dev.yml @@ -62,6 +62,7 @@ services: # # System workers: # - archiver + # - release-validator # - limiter # - paymaster # @@ -79,6 +80,20 @@ services: - ./:/usr/src/app - workers-deps:/usr/src/app/node_modules + hawk-worker-release-validator: + build: + dockerfile: "dev.Dockerfile" + context: . + env_file: + - .env + environment: + - SIMULTANEOUS_TASKS=1 + restart: unless-stopped + entrypoint: yarn run-release-validator + volumes: + - ./:/usr/src/app + - workers-deps:/usr/src/app/node_modules + hawk-worker-limiter: build: dockerfile: "dev.Dockerfile" diff --git a/docker-compose.prod.yml b/docker-compose.prod.yml index 88a6e0fa..a702aefc 100644 --- a/docker-compose.prod.yml +++ b/docker-compose.prod.yml @@ -43,7 +43,7 @@ services: entrypoint: /usr/local/bin/node runner.js hawk-worker-grouper # - # System workers: archiver + # System workers: archiver, release-validator # hawk-worker-archiver: @@ -55,6 +55,16 @@ services: restart: unless-stopped entrypoint: /usr/local/bin/node runner.js hawk-worker-archiver + hawk-worker-release-validator: + image: "codexteamuser/hawk-workers:prod" + network_mode: host + env_file: + - .env + environment: + - SIMULTANEOUS_TASKS=1 + restart: unless-stopped + entrypoint: /usr/local/bin/node runner.js hawk-worker-release-validator + # # Notification workers: notifier, email, telegram # diff --git a/package.json b/package.json index f2e1e154..32953a73 100644 --- a/package.json +++ b/package.json @@ -25,6 +25,7 @@ "test:sentry": "jest workers/sentry --config workers/sentry/jest.config.js", "test:javascript": "jest workers/javascript", "test:release": "jest workers/release", + "test:release-validator": "jest workers/release-validator", "test:slack": "jest workers/slack", "test:loop": "jest workers/loop", "test:limiter": "jest workers/limiter --runInBand", @@ -47,6 +48,7 @@ "run-paymaster": "yarn worker hawk-worker-paymaster", "run-notifier": "yarn worker hawk-worker-notifier", "run-release": "yarn worker hawk-worker-release", + "run-release-validator": "yarn worker hawk-worker-release-validator", "run-email": "yarn worker hawk-worker-email", "run-telegram": "yarn worker hawk-worker-telegram", "run-limiter": "yarn worker hawk-worker-limiter", @@ -57,7 +59,7 @@ "@babel/parser": "^7.26.9", "@babel/traverse": "7.26.9", "@hawk.so/nodejs": "^3.1.1", - "@hawk.so/types": "^0.5.9", + "@hawk.so/types": "^0.9.0", "@types/amqplib": "^0.8.2", "@types/jest": "^29.5.14", "@types/mongodb": "^3.5.15", diff --git a/workers/grouper/src/check-and-mark-regression.ts b/workers/grouper/src/check-and-mark-regression.ts new file mode 100644 index 00000000..e5084a49 --- /dev/null +++ b/workers/grouper/src/check-and-mark-regression.ts @@ -0,0 +1,68 @@ +import { Db } from 'mongodb'; +import type { ReleaseDBScheme } from '@hawk.so/types'; + +type ReleaseRecordPart = Pick; + +/** + * Mark an event as regressed if it reoccurs in the resolved or a newer release. + * + * The update is atomic: only the first repetition after resolution sets the + * regression release, and later repetitions do not overwrite it. + * + * @param db - events database connection + * @param projectId - project identifier + * @param groupHash - original event group hash + * @param release - release in which the event occurred again + * @param resolvedInRelease - release in which the event was resolved + * @param regressionInRelease - regression from a previous resolution cycle + */ +export async function checkAndMarkRegression( + db: Db, + projectId: string, + groupHash: string, + release: string, + resolvedInRelease: string, + regressionInRelease?: string +): Promise { + const releases = await db.collection('releases').find({ + projectId, + release: { + $in: [resolvedInRelease, release, regressionInRelease].filter(Boolean), + }, + }) + .toArray(); + const resolvedRelease = releases.find(item => item.release === resolvedInRelease); + const repetitionRelease = releases.find(item => item.release === release); + const previousRegressionRelease = regressionInRelease + ? releases.find(item => item.release === regressionInRelease) + : undefined; + + if (!resolvedRelease || !repetitionRelease) { + return; + } + + const resolvedReleaseId = resolvedRelease._id.toHexString(); + const isResolvedOrNewerRelease = repetitionRelease._id.toHexString() >= resolvedReleaseId; + /** + * A regression in or after the resolved release belongs to the current + * resolution cycle and must not be overwritten by later repetitions. + */ + const hasRegressionForCurrentCycle = previousRegressionRelease && + previousRegressionRelease._id.toHexString() >= resolvedReleaseId; + + if (!isResolvedOrNewerRelease || hasRegressionForCurrentCycle) { + return; + } + + await db.collection(`events:${projectId}`).updateOne({ + groupHash, + resolvedInRelease, + ...(regressionInRelease + ? { regressionInRelease } + : { regressionInRelease: { $exists: false } }), + }, { + $set: { + regressionInRelease: release, + }, + }); +} diff --git a/workers/grouper/src/index.ts b/workers/grouper/src/index.ts index 6750d768..e6839c39 100644 --- a/workers/grouper/src/index.ts +++ b/workers/grouper/src/index.ts @@ -24,6 +24,7 @@ import GrouperMetrics from './metrics/grouperMetrics'; import GrouperMemoryMonitor from './metrics/memoryMonitor'; import SlowHandleDiagnostics, { SlowHandleSession } from './metrics/slowHandleDiagnostics'; import { grouperDiagnosticsConfig, grouperMemoryConfig } from './metrics/config'; +import { checkAndMarkRegression } from './check-and-mark-regression'; /** * eslint does not count decorators as a variable usage @@ -343,10 +344,32 @@ export default class GrouperWorker extends Worker { timestamp: task.timestamp, } as RepetitionDBScheme; + if (task.payload.release) { + newRepetition.release = task.payload.release; + } + repetitionId = await session.measureStep('saveRepetition', () => { return this.saveRepetition(task.projectId, newRepetition); }); + if (task.payload.release && existedEvent.resolvedInRelease) { + try { + await checkAndMarkRegression( + this.eventsDb.getConnection(), + task.projectId, + uniqueEventHash, + task.payload.release, + existedEvent.resolvedInRelease, + existedEvent.regressionInRelease + ); + } catch (error) { + this.logger.error( + `[checkAndMarkRegression] project=${task.projectId} groupHash=${uniqueEventHash} release=${task.payload.release}`, + error + ); + } + } + /** * Clear the large event payload references to allow garbage collection * This prevents memory leaks from retaining full event objects after delta is computed diff --git a/workers/grouper/tests/index.test.ts b/workers/grouper/tests/index.test.ts index e4f2b3bd..6e2421c3 100644 --- a/workers/grouper/tests/index.test.ts +++ b/workers/grouper/tests/index.test.ts @@ -158,6 +158,7 @@ describe('GrouperWorker', () => { await eventsCollection.deleteMany({}); await dailyEventsCollection.deleteMany({}); await repetitionsCollection.deleteMany({}); + await connection.db().collection('releases').deleteMany({ projectId: projectIdMock }); }); afterEach(async () => { @@ -391,6 +392,156 @@ describe('GrouperWorker', () => { }).toArray()).length).toBe(2); }); + test('Should save repetition release as a separate field', async () => { + await worker.handle(generateTask({ release: 'release-a' })); + await worker.handle(generateTask({ release: 'release-b' })); + + const savedRepetition = await repetitionsCollection.findOne({}); + + expect(savedRepetition.release).toBe('release-b'); + }); + + test('Should not save repetition release when event has no release', async () => { + await worker.handle(generateTask()); + await worker.handle(generateTask()); + + const savedRepetition = await repetitionsCollection.findOne({}); + + expect(savedRepetition.release).toBeUndefined(); + }); + + describe('Regression marking', () => { + test('Should mark a resolved event as regressed in a newer repetition release', async () => { + await connection.db().collection('releases').insertMany([ + { + _id: mongodb.ObjectID.createFromTime(1), + projectId: projectIdMock, + release: 'release-b', + }, + { + _id: mongodb.ObjectID.createFromTime(2), + projectId: projectIdMock, + release: 'release-c', + }, + ]); + await worker.handle(generateTask({ release: 'release-a' })); + await eventsCollection.updateOne({}, { + $set: { + resolvedInRelease: 'release-b', + }, + }); + + await worker.handle(generateTask({ release: 'release-c' })); + + expect((await eventsCollection.findOne({})).regressionInRelease).toBe('release-c'); + }); + + test('Should mark as regressed if we later encounter this event with a release that is considered a resolving release', async () => { + await connection.db().collection('releases').insertOne({ + _id: mongodb.ObjectID.createFromTime(1), + projectId: projectIdMock, + release: 'release-b', + }); + await worker.handle(generateTask({ release: 'release-a' })); + await eventsCollection.updateOne({}, { + $set: { + resolvedInRelease: 'release-b', + }, + }); + + await worker.handle(generateTask({ release: 'release-b' })); + + expect((await eventsCollection.findOne({})).regressionInRelease).toBe('release-b'); + }); + + test('Should not mark regression if we encounter an event from one of the old releases', async () => { + await connection.db().collection('releases').insertMany([ + { + _id: mongodb.ObjectID.createFromTime(1), + projectId: projectIdMock, + release: 'release-a', + }, + { + _id: mongodb.ObjectID.createFromTime(2), + projectId: projectIdMock, + release: 'release-b', + }, + ]); + await worker.handle(generateTask({ release: 'release-b' })); + await eventsCollection.updateOne({}, { + $set: { + resolvedInRelease: 'release-b', + }, + }); + + await worker.handle(generateTask({ release: 'release-a' })); + + expect((await eventsCollection.findOne({})).regressionInRelease).toBeUndefined(); + }); + + test('Should replace an old regression after a newer resolution', async () => { + await connection.db().collection('releases').insertMany([ + { + _id: mongodb.ObjectID.createFromTime(1), + projectId: projectIdMock, + release: 'release-b', + }, + { + _id: mongodb.ObjectID.createFromTime(2), + projectId: projectIdMock, + release: 'release-c', + }, + { + _id: mongodb.ObjectID.createFromTime(3), + projectId: projectIdMock, + release: 'release-d', + }, + ]); + await worker.handle(generateTask({ release: 'release-a' })); + await eventsCollection.updateOne({}, { + $set: { + resolvedInRelease: 'release-c', + regressionInRelease: 'release-b', + }, + }); + + await worker.handle(generateTask({ release: 'release-d' })); + + expect((await eventsCollection.findOne({})).regressionInRelease).toBe('release-d'); + }); + + test('Should not overwrite the first marked regression if we continue receiving this event in newer releases', async () => { + await connection.db().collection('releases').insertMany([ + { + _id: mongodb.ObjectID.createFromTime(1), + projectId: projectIdMock, + release: 'release-b', + }, + { + _id: mongodb.ObjectID.createFromTime(2), + projectId: projectIdMock, + release: 'release-c', + }, + { + _id: mongodb.ObjectID.createFromTime(3), + projectId: projectIdMock, + release: 'release-d', + }, + ]); + await worker.handle(generateTask({ release: 'release-a' })); + await eventsCollection.updateOne({}, { + $set: { + resolvedInRelease: 'release-b', + regressionInRelease: 'release-c', + }, + }); + + await worker.handle(generateTask({ release: 'release-d' })); + + expect((await eventsCollection.findOne({})).regressionInRelease).toBe('release-c'); + }); + }); + test('Should stringify delta', async () => { const generatedTask = generateTask(); diff --git a/workers/release-validator/README.md b/workers/release-validator/README.md new file mode 100644 index 00000000..333d4fb6 --- /dev/null +++ b/workers/release-validator/README.md @@ -0,0 +1,15 @@ +# Release Validator Worker + +Checks releases after a 24-hour observation period and marks original events that no longer occur as resolved. + +The worker processes releases from oldest to newest, stores the first release without an event in `resolvedInRelease`, and marks successfully processed releases with `fixChecked: true`. + +The worker only uses the fields required for matching releases and event groups. Records without these fields are ignored. Candidate releases must be between 24 hours and 30 days old, while all project releases are still used to compare event history. + +Queue: `cron-tasks/release-validator` + +Run locally: + +```sh +yarn run-release-validator +``` diff --git a/workers/release-validator/package.json b/workers/release-validator/package.json new file mode 100644 index 00000000..3cd8c88f --- /dev/null +++ b/workers/release-validator/package.json @@ -0,0 +1,10 @@ +{ + "name": "hawk-worker-release-validator", + "version": "0.0.1", + "description": "Detects events fixed by a release", + "main": "src/index.ts", + "author": "CodeX", + "license": "UNLICENSED", + "private": true, + "workerType": "cron-tasks/release-validator" +} diff --git a/workers/release-validator/src/index.ts b/workers/release-validator/src/index.ts new file mode 100644 index 00000000..ac268f6a --- /dev/null +++ b/workers/release-validator/src/index.ts @@ -0,0 +1,46 @@ +import { DatabaseController } from '../../../lib/db/controller'; +import { Worker } from '../../../lib/worker'; +import * as pkg from '../package.json'; +import { validateReleases } from './validate-releases'; + +/** + * Worker that detects events fixed by a release. + */ +export default class ReleaseValidatorWorker extends Worker { + /** + * Worker type. + */ + public readonly type: string = pkg.workerType; + + /** + * Events database controller. + */ + private eventsDb = new DatabaseController(process.env.MONGO_EVENTS_DATABASE_URI); + + /** + * Connect to the events database and start consuming tasks. + */ + public async start(): Promise { + await this.eventsDb.connect(); + await super.start(); + } + + /** + * Stop consuming tasks and close the database connection. + */ + public async finish(): Promise { + await super.finish(); + await this.eventsDb.close(); + } + + /** + * Handle a scheduled release validation task. + */ + public async handle(): Promise { + this.logger.info('Release validation started'); + + await validateReleases(this.eventsDb.getConnection()); + + this.logger.info('Release validation finished'); + } +} diff --git a/workers/release-validator/src/utils/build-event-release-map.ts b/workers/release-validator/src/utils/build-event-release-map.ts new file mode 100644 index 00000000..13db2382 --- /dev/null +++ b/workers/release-validator/src/utils/build-event-release-map.ts @@ -0,0 +1,45 @@ +import type { GroupedEventDBScheme, RepetitionDBScheme } from '@hawk.so/types'; + +type RepetitionRelease = Pick; + +/** + * Build a lookup set of releases in which each event occurred. + * + * Repetitions are deduplicated by event and release in MongoDB before this + * function is called. Therefore, memory usage depends on the number of unique + * releases per event rather than the potentially much larger repetition count. + * Release ordering is handled separately by the validation flow. + * + * @param events - original events from the current validation batch + * @param repetitions - unique event and release pairs from MongoDB + */ +export function buildEventReleaseMap( + events: GroupedEventDBScheme[], + repetitions: RepetitionRelease[] +): Map> { + const releasesByGroupHash = new Map>(); + + for (const event of events) { + const releases = new Set(); + + if (event.payload?.release) { + releases.add(event.payload.release); + } + + releasesByGroupHash.set(event.groupHash, releases); + } + + for (const repetition of repetitions) { + if (!repetition.release) { + continue; + } + + const releases = releasesByGroupHash.get(repetition.groupHash); + + if (releases) { + releases.add(repetition.release); + } + } + + return releasesByGroupHash; +} diff --git a/workers/release-validator/src/utils/group-releases-by-project.ts b/workers/release-validator/src/utils/group-releases-by-project.ts new file mode 100644 index 00000000..efd30ae3 --- /dev/null +++ b/workers/release-validator/src/utils/group-releases-by-project.ts @@ -0,0 +1,19 @@ +import type { ReleaseDBScheme } from '@hawk.so/types'; + +/** + * Group chronologically ordered releases by project. + * + * @param releases - releases ordered from oldest to newest + */ +export function groupReleasesByProject(releases: ReleaseDBScheme[]): Map { + const releasesByProject = new Map(); + + for (const release of releases) { + const projectReleases = releasesByProject.get(release.projectId) || []; + + projectReleases.push(release); + releasesByProject.set(release.projectId, projectReleases); + } + + return releasesByProject; +} diff --git a/workers/release-validator/src/validate-releases.ts b/workers/release-validator/src/validate-releases.ts new file mode 100644 index 00000000..b6d58289 --- /dev/null +++ b/workers/release-validator/src/validate-releases.ts @@ -0,0 +1,382 @@ +import { Collection, Db, ObjectID } from 'mongodb'; +import type { GroupedEventDBScheme, ReleaseDBScheme, RepetitionDBScheme } from '@hawk.so/types'; +import { HOURS_IN_DAY, MINUTES_IN_HOUR, MS_IN_SEC, SECONDS_IN_MINUTE } from '../../../lib/utils/consts'; +import { buildEventReleaseMap } from './utils/build-event-release-map'; +import { groupReleasesByProject } from './utils/group-releases-by-project'; + +type ReleaseHistoryEntry = Pick; + +/** + * Time allowed for repetitions to arrive before a release is checked for + * resolved events. + */ +const RELEASE_OBSERVATION_PERIOD_SECONDS = HOURS_IN_DAY * MINUTES_IN_HOUR * SECONDS_IN_MINUTE; + +/** + * Maximum age of a release eligible for validation, in days. + */ +const RELEASE_MAX_AGE_DAYS = 30; + +/** + * Maximum candidate release age expressed in seconds for ObjectId boundaries. + */ +const RELEASE_MAX_AGE_SECONDS = RELEASE_MAX_AGE_DAYS * RELEASE_OBSERVATION_PERIOD_SECONDS; + +/** + * Maximum number of events processed in one validation batch. + * + * A bounded batch keeps the repetitions `$in` query below MongoDB document + * limits and prevents large projects from loading all events and repetitions + * into worker memory at once. + */ +const EVENTS_BATCH_SIZE = 500; + +/** + * Find unchecked releases older than the observation period. + * + * @param db - events database connection + * @param now - current time + */ +async function findReleasesToCheck(db: Db, now: Date): Promise { + const nowSeconds = Math.floor(now.getTime() / MS_IN_SEC); + + /** + * Limit candidates to the rollout window: releases must be old enough to + * observe for 24 hours, but recent enough to contain release-aware repetitions. + */ + const oldestReleaseId = ObjectID.createFromTime(nowSeconds - RELEASE_MAX_AGE_SECONDS); + const newestReleaseId = ObjectID.createFromTime(nowSeconds - RELEASE_OBSERVATION_PERIOD_SECONDS); + + return db.collection('releases') + .find({ + _id: { + $gte: oldestReleaseId, + $lt: newestReleaseId, + }, + release: { + $type: 'string', + $ne: '', + }, + fixChecked: { $ne: true }, + }) + .sort({ _id: 1 }) + .toArray(); +} + +/** + * Validate a bounded portion of project events. + * + * Only the repetition fields required to build the occurrence lookup are + * projected. This avoids transferring serialized deltas and payloads that are + * not used by release validation. + * + * @param events - current event portion + * @param eventsCollection - project events collection + * @param repetitionsCollection - project repetitions collection + * @param releasesToCheck - ready releases ordered from oldest to newest + * @param allProjectReleases - complete project release history + * @param releasesByName - project releases indexed by release name + */ +async function validateEventsBatch( + events: GroupedEventDBScheme[], + eventsCollection: Collection, + repetitionsCollection: Collection, + releasesToCheck: ReleaseDBScheme[], + allProjectReleases: ReleaseHistoryEntry[], + releasesByName: Map +): Promise { + const eventGroupHashes = events.map(event => event.groupHash); + const repetitions = await repetitionsCollection.aggregate>([ + { + $match: { + groupHash: { $in: eventGroupHashes }, + release: { + $type: 'string', + $ne: '', + }, + }, + }, + { + $group: { + _id: { + groupHash: '$groupHash', + release: '$release', + }, + }, + }, + { + $project: { + _id: 0, + groupHash: '$_id.groupHash', + release: '$_id.release', + }, + }, + ]).toArray(); + const eventReleases = buildEventReleaseMap(events, repetitions); + + /** + * Resolve each event in the first eligible release after its latest known + * occurrence. A later occurrence blocks resolution in an older release. + */ + for (const event of events) { + /** + * The original event release is the initial occurrence boundary. + */ + const originalRelease = releasesByName.get(event.payload.release); + let lastOccurrenceRelease = originalRelease; + + /** + * An event can be resolved, reappear in a later release, and then stop + * occurring again. In that case, search for the next resolving release + * after the regression rather than after the event's original occurrence. + */ + if (event.resolvedInRelease && event.regressionInRelease) { + const resolvedRelease = releasesByName.get(event.resolvedInRelease); + const regressionRelease = releasesByName.get(event.regressionInRelease); + + /** + * Without both release records their chronological relation is unknown, + * so the event cannot be safely resolved again. + */ + if (!resolvedRelease || !regressionRelease) { + continue; + } + + const isCurrentlyRegressed = regressionRelease._id.toHexString() >= resolvedRelease._id.toHexString(); + + /** + * A regression older than the current resolution belongs to a previous + * cycle and does not make the event eligible for another resolution. + */ + if (!isCurrentlyRegressed) { + continue; + } + + lastOccurrenceRelease = regressionRelease; + } + + /** + * Skip events whose original or latest occurrence release is missing from + * the project release history. + */ + if (!lastOccurrenceRelease) { + continue; + } + + const releasesWithEvent = eventReleases.get(event.groupHash) || new Set(); + + for (const release of releasesToCheck) { + const releaseId = release._id.toHexString(); + const isNewerThanLastOccurrence = releaseId > lastOccurrenceRelease._id.toHexString(); + const occurredInRelease = releasesWithEvent.has(release.release); + + /** + * A repetition in any later release, including one younger than 24 hours, + * proves that this candidate did not fix the event. + */ + const occurredInNewerRelease = allProjectReleases.some(projectRelease => { + return projectRelease._id.toHexString() > releaseId && releasesWithEvent.has(projectRelease.release); + }); + + /** + * A candidate resolves the event only if it was deployed after the latest + * occurrence and the event appears neither in that candidate nor in any + * newer release. Otherwise, continue with the next candidate. + */ + if (!isNewerThanLastOccurrence || occurredInRelease || occurredInNewerRelease) { + continue; + } + + /** + * Resolve only the state that was evaluated in this batch: an unresolved + * event must still be unresolved, while a regressed event must retain the + * same resolution and regression releases. Including that state in the + * update also makes the transition atomic, so a concurrent Grouper update + * cannot be overwritten. + */ + const eventState = event.regressionInRelease + ? { + resolvedInRelease: event.resolvedInRelease, + regressionInRelease: event.regressionInRelease, + } + : { + $or: [ + { resolvedInRelease: { $exists: false } }, + { resolvedInRelease: null }, + ], + }; + + await eventsCollection.updateOne({ + _id: event._id, + ...eventState, + }, { + $set: { + resolvedInRelease: release.release, + }, + }); + + /** + * Releases are ordered from oldest to newest, so the first matching + * candidate is the release in which the event became likely fixed. + */ + break; + } + } +} + +/** + * Validate all ready releases of one project. + * + * @param db - events database connection + * @param projectId - project identifier + * @param releasesToCheck - ready releases ordered from oldest to newest + */ +async function validateProject(db: Db, projectId: string, releasesToCheck: ReleaseDBScheme[]): Promise { + const releasesCollection = db.collection('releases'); + const eventsCollection = db.collection(`events:${projectId}`); + const repetitionsCollection = db.collection(`repetitions:${projectId}`); + + /** + * Load the project's release timeline once. Releases outside the validation + * window are required because an old original release or a newer occurrence + * can change whether a candidate release resolved an event. Only identifiers + * and names are loaded; source maps and commit data are not used here. + */ + const allProjectReleases = await releasesCollection + .find({ + projectId, + release: { + $type: 'string', + $ne: '', + }, + }) + .project({ + _id: 1, + release: 1, + }) + .sort({ _id: 1 }) + .toArray(); + + /** + * Index the timeline by release name for event-field lookups. The map key is + * only a lookup key; chronology is still determined by each value's ObjectId. + */ + const releasesByName = new Map(); + + for (const release of allProjectReleases) { + releasesByName.set(release.release, release); + } + + /** + * Find events whose resolution state needs evaluation. An event must have a + * group hash and an original release, and must either be unresolved or have + * both resolution and regression releases for a subsequent resolution cycle. + * + * Stream the result instead of materializing all project events. The cursor + * and application batch use the same limit so the worker holds at most one + * bounded portion of event documents in memory. + */ + const eventsCursor = eventsCollection.find({ + groupHash: { + $type: 'string', + $ne: '', + }, + 'payload.release': { + $type: 'string', + $ne: '', + }, + $or: [ + { resolvedInRelease: { $exists: false } }, + { resolvedInRelease: null }, + { + resolvedInRelease: { + $type: 'string', + $ne: '', + }, + regressionInRelease: { + $type: 'string', + $ne: '', + }, + }, + ], + }, { + projection: { + _id: 1, + groupHash: 1, + 'payload.release': 1, + resolvedInRelease: 1, + regressionInRelease: 1, + }, + }).batchSize(EVENTS_BATCH_SIZE); + let eventsBatch: GroupedEventDBScheme[] = []; + + while (await eventsCursor.hasNext()) { + const event = await eventsCursor.next(); + + if (!event) { + break; + } + + eventsBatch.push(event); + + if (eventsBatch.length < EVENTS_BATCH_SIZE) { + continue; + } + + await validateEventsBatch( + eventsBatch, + eventsCollection, + repetitionsCollection, + releasesToCheck, + allProjectReleases, + releasesByName + ); + eventsBatch = []; + } + + /** + * Process the final partial portion left after the cursor is exhausted. + */ + if (eventsBatch.length > 0) { + await validateEventsBatch( + eventsBatch, + eventsCollection, + repetitionsCollection, + releasesToCheck, + allProjectReleases, + releasesByName + ); + } + + /** + * Mark releases only after every project batch finishes successfully. + */ + await releasesCollection.updateMany({ + _id: { + $in: releasesToCheck.map(release => release._id), + }, + }, { + $set: { + fixChecked: true, + }, + }); +} + +/** + * Validate all releases whose observation period has elapsed. + * + * @param db - events database connection + * @param now - current time + */ +export async function validateReleases(db: Db, now = new Date()): Promise { + /** + * Candidate releases are globally ordered and then grouped so each project's + * release history, events, and repetitions are loaded only once per run. + */ + const releasesToCheck = await findReleasesToCheck(db, now); + const releasesByProject = groupReleasesByProject(releasesToCheck); + + for (const [projectId, projectReleases] of releasesByProject) { + await validateProject(db, projectId, projectReleases); + } +} diff --git a/workers/release-validator/tests/utils/build-event-release-map.test.ts b/workers/release-validator/tests/utils/build-event-release-map.test.ts new file mode 100644 index 00000000..a9d6025f --- /dev/null +++ b/workers/release-validator/tests/utils/build-event-release-map.test.ts @@ -0,0 +1,97 @@ +import { ObjectID } from 'mongodb'; +import type { GroupedEventDBScheme, RepetitionDBScheme } from '@hawk.so/types'; +import { buildEventReleaseMap } from '../../src/utils/build-event-release-map'; + +/** + * Create an original event. + * + * @param groupHash - event group hash + * @param release - release in which the event first occurred + */ +function createEvent(groupHash: string, release?: string): GroupedEventDBScheme { + return { + _id: new ObjectID(), + groupHash, + payload: { + title: groupHash, + ...(release ? { release } : {}), + }, + totalCount: 1, + catcherType: 'errors/default', + usersAffected: 0, + visitedBy: [], + timestamp: 1, + }; +} + +/** + * Create an event repetition. + * + * @param groupHash - event group hash + * @param release - release in which the event occurred + */ +function createRepetition(groupHash: string, release?: string): RepetitionDBScheme { + return { + groupHash, + timestamp: 1, + ...(release ? { release } : {}), + }; +} + +describe('buildEventReleaseMap', () => { + test('should build the release map from original events and repetitions', () => { + // Arrange + const events = [ + createEvent('error-1', 'a'), + createEvent('error-2', 'a'), + createEvent('error-3', 'a'), + ]; + const repetitions = [ + createRepetition('error-1', 'b'), + createRepetition('error-1', 'c'), + createRepetition('error-2', 'b'), + createRepetition('error-3', 'b'), + createRepetition('error-3', 'c'), + createRepetition('error-3', 'd'), + createRepetition('error-3', 'e'), + ]; + + // Act + const result = buildEventReleaseMap(events, repetitions); + + // Assert + expect(Array.from(result.entries()).map(([groupHash, releases]) => { + return [groupHash, Array.from(releases)]; + })).toEqual([ + ['error-1', ['a', 'b', 'c'] ], + ['error-2', ['a', 'b'] ], + ['error-3', ['a', 'b', 'c', 'd', 'e'] ], + ]); + }); + + test('should ignore duplicate releases and repetitions without matching events', () => { + // Arrange + const events = [ + createEvent('error-1', 'a'), + createEvent('error-without-release'), + ]; + const repetitions = [ + createRepetition('error-1', 'a'), + createRepetition('error-1', 'b'), + createRepetition('error-1', 'b'), + createRepetition('error-1'), + createRepetition('unknown-error', 'c'), + ]; + + // Act + const result = buildEventReleaseMap(events, repetitions); + + // Assert + expect(Array.from(result.entries()).map(([groupHash, releases]) => { + return [groupHash, Array.from(releases)]; + })).toEqual([ + ['error-1', ['a', 'b'] ], + ['error-without-release', [] ], + ]); + }); +}); diff --git a/workers/release-validator/tests/utils/group-releases-by-project.test.ts b/workers/release-validator/tests/utils/group-releases-by-project.test.ts new file mode 100644 index 00000000..0e53fe78 --- /dev/null +++ b/workers/release-validator/tests/utils/group-releases-by-project.test.ts @@ -0,0 +1,54 @@ +import { ObjectID } from 'mongodb'; +import type { ReleaseDBScheme } from '@hawk.so/types'; +import { groupReleasesByProject } from '../../src/utils/group-releases-by-project'; + +/** + * Create a release record with a predictable id. + * + * @param projectId - project identifier + * @param release - release name + * @param createdAtSeconds - release creation time + */ +function createRelease(projectId: string, release: string, createdAtSeconds: number): ReleaseDBScheme { + return { + _id: ObjectID.createFromTime(createdAtSeconds), + projectId, + release, + commits: [], + }; +} + +describe('groupReleasesByProject', () => { + test('should group releases by project and preserve their input order', () => { + // Arrange + const releases = [ + createRelease('project-a', 'a', 1), + createRelease('project-b', 'x', 2), + createRelease('project-a', 'b', 3), + createRelease('project-b', 'y', 4), + createRelease('project-a', 'c', 5), + ]; + + // Act + const result = groupReleasesByProject(releases); + + // Assert + expect(Array.from(result.entries()).map(([projectId, projectReleases]) => { + return [projectId, projectReleases.map(release => release.release)]; + })).toEqual([ + ['project-a', ['a', 'b', 'c'] ], + ['project-b', ['x', 'y'] ], + ]); + }); + + test('should return an empty map for an empty release list', () => { + // Arrange + const releases: ReleaseDBScheme[] = []; + + // Act + const result = groupReleasesByProject(releases); + + // Assert + expect(Array.from(result.entries())).toEqual([]); + }); +}); diff --git a/workers/release-validator/tests/validate-releases.test.ts b/workers/release-validator/tests/validate-releases.test.ts new file mode 100644 index 00000000..332c12a2 --- /dev/null +++ b/workers/release-validator/tests/validate-releases.test.ts @@ -0,0 +1,355 @@ +import '../../../env-test'; +import { Collection, Db, MongoClient, ObjectID } from 'mongodb'; +import type { GroupedEventDBScheme, ReleaseDBScheme, RepetitionDBScheme } from '@hawk.so/types'; +import { validateReleases } from '../src/validate-releases'; + +const PROJECT_ID = 'release-validator-project'; +const HOUR_IN_SECONDS = 60 * 60; +const NOW_SECONDS = Math.floor(new Date('2026-09-17T12:00:00.000Z').getTime() / 1000); +const NOW = new Date(NOW_SECONDS * 1000); + +/** + * Create a release ObjectId with a predictable creation time. + * + * @param hoursBeforeNow - release age in hours + */ +function releaseId(hoursBeforeNow: number): ObjectID { + return ObjectID.createFromTime(NOW_SECONDS - hoursBeforeNow * HOUR_IN_SECONDS); +} + +/** + * Create a release record for the test project. + * + * @param release - release name + * @param hoursBeforeNow - release age in hours + * @param fixChecked - whether the release has already been checked + */ +function createRelease(release: string, hoursBeforeNow: number, fixChecked = false): ReleaseDBScheme { + return { + _id: releaseId(hoursBeforeNow), + projectId: PROJECT_ID, + release, + commits: [], + fixChecked, + }; +} + +/** + * Create an original event. + * + * @param groupHash - event group hash + * @param release - release in which the event first occurred + * @param resolvedInRelease - existing resolved release + */ +function createEvent(groupHash: string, release: string, resolvedInRelease?: string): GroupedEventDBScheme { + const event: GroupedEventDBScheme = { + _id: new ObjectID(), + groupHash, + payload: { + title: groupHash, + release, + }, + totalCount: 1, + catcherType: 'errors/default', + usersAffected: 0, + visitedBy: [], + timestamp: NOW_SECONDS, + }; + + if (resolvedInRelease !== undefined) { + event.resolvedInRelease = resolvedInRelease; + } + + return event; +} + +/** + * Create an event repetition. + * + * @param groupHash - event group hash + * @param release - release in which the event occurred + */ +function createRepetition(groupHash: string, release: string): RepetitionDBScheme { + return { + groupHash, + release, + timestamp: NOW_SECONDS, + }; +} + +describe('validateReleases', () => { + let connection: MongoClient; + let db: Db; + let releases: Collection; + let events: Collection; + let repetitions: Collection; + + beforeAll(async () => { + connection = await MongoClient.connect(process.env.MONGO_EVENTS_DATABASE_URI, { + useNewUrlParser: true, + useUnifiedTopology: true, + }); + db = connection.db(); + releases = db.collection('releases'); + events = db.collection(`events:${PROJECT_ID}`); + repetitions = db.collection(`repetitions:${PROJECT_ID}`); + }); + + beforeEach(async () => { + await releases.deleteMany({ projectId: PROJECT_ID }); + await events.deleteMany({}); + await repetitions.deleteMany({}); + }); + + afterAll(async () => { + await releases.deleteMany({ projectId: PROJECT_ID }); + await events.drop().catch(() => undefined); + await repetitions.drop().catch(() => undefined); + await connection.close(); + }); + + test('should not resolve an event that first appeared in the checked release', async () => { + await releases.insertMany([ + createRelease('a', 72, true), + createRelease('b', 48), + ]); + await events.insertOne(createEvent('error-1', 'b')); + + await validateReleases(db, NOW); + + expect((await events.findOne({ groupHash: 'error-1' })).resolvedInRelease).toBeUndefined(); + expect((await releases.findOne({ release: 'b' })).fixChecked).toBe(true); + }); + + test('should not resolve an event that occurred in the checked release', async () => { + await releases.insertMany([ + createRelease('a', 72, true), + createRelease('b', 48), + ]); + await events.insertOne(createEvent('error-2', 'a')); + await repetitions.insertOne(createRepetition('error-2', 'b')); + + await validateReleases(db, NOW); + + expect((await events.findOne({ groupHash: 'error-2' })).resolvedInRelease).toBeUndefined(); + }); + + test('should not resolve an event that reappeared in a newer release within 24 hours', async () => { + await releases.insertMany([ + createRelease('a', 72, true), + createRelease('b', 48), + createRelease('c', 2), + ]); + await events.insertOne(createEvent('error-3', 'a')); + await repetitions.insertOne(createRepetition('error-3', 'c')); + + await validateReleases(db, NOW); + + expect((await events.findOne({ groupHash: 'error-3' })).resolvedInRelease).toBeUndefined(); + expect((await releases.findOne({ release: 'b' })).fixChecked).toBe(true); + expect((await releases.findOne({ release: 'c' })).fixChecked).toBe(false); + }); + + test('should resolve an event in the checked release when it never appears again', async () => { + await releases.insertMany([ + createRelease('a', 72, true), + createRelease('b', 48), + createRelease('c', 2), + ]); + await events.insertOne(createEvent('error-4', 'a')); + + await validateReleases(db, NOW); + + expect((await events.findOne({ groupHash: 'error-4' })).resolvedInRelease).toBe('b'); + }); + + test('should resolve an event in the first checked release after its last occurrence', async () => { + await releases.insertMany([ + createRelease('a', 120, true), + createRelease('b', 96), + createRelease('c', 72), + createRelease('d', 48), + ]); + await events.insertOne(createEvent('error-5', 'a')); + await repetitions.insertOne(createRepetition('error-5', 'c')); + + await validateReleases(db, NOW); + + expect((await events.findOne({ groupHash: 'error-5' })).resolvedInRelease).toBe('d'); + }); + + test('should not check a release until its 24-hour observation period has elapsed', async () => { + await releases.insertMany([ + createRelease('a', 48, true), + createRelease('b', 23), + ]); + await events.insertOne(createEvent('error-6', 'a')); + + await validateReleases(db, NOW); + + expect((await events.findOne({ groupHash: 'error-6' })).resolvedInRelease).toBeUndefined(); + expect((await releases.findOne({ release: 'b' })).fixChecked).toBe(false); + }); + + test('should skip an event when its original release is unknown', async () => { + await releases.insertOne(createRelease('b', 48)); + await events.insertOne(createEvent('error-7', 'unknown')); + + await validateReleases(db, NOW); + + expect((await events.findOne({ groupHash: 'error-7' })).resolvedInRelease).toBeUndefined(); + expect((await releases.findOne({ release: 'b' })).fixChecked).toBe(true); + }); + + test('should not overwrite an existing resolved release on repeated validation', async () => { + await releases.insertMany([ + createRelease('a', 72, true), + createRelease('b', 48), + ]); + await events.insertOne(createEvent('error-8', 'a', 'previous-release')); + + await validateReleases(db, NOW); + await validateReleases(db, NOW); + + expect((await events.findOne({ groupHash: 'error-8' })).resolvedInRelease).toBe('previous-release'); + }); + + test('should skip a malformed release and continue with the next valid release', async () => { + await releases.insertMany([ + createRelease('a', 96, true), + { + _id: releaseId(72), + projectId: PROJECT_ID, + release: 123 as unknown as string, + commits: [], + }, + createRelease('d', 48), + ]); + await events.insertOne(createEvent('error-9', 'a')); + + await expect(validateReleases(db, NOW)).resolves.toBeUndefined(); + + expect((await events.findOne({ groupHash: 'error-9' })).resolvedInRelease).toBe('d'); + expect((await releases.findOne({ release: 123 as unknown as string })).fixChecked).toBeUndefined(); + expect((await releases.findOne({ release: 'd' })).fixChecked).toBe(true); + }); + + test('should skip a malformed event without stopping valid events', async () => { + await releases.insertMany([ + createRelease('a', 72, true), + createRelease('b', 48), + ]); + await events.insertMany([ + { + _id: new ObjectID(), + groupHash: 'broken-event', + payload: { + title: 'broken-event', + }, + totalCount: 1, + catcherType: 'errors/default', + usersAffected: 0, + visitedBy: [], + timestamp: NOW_SECONDS, + }, + createEvent('valid-event', 'a'), + ]); + + await expect(validateReleases(db, NOW)).resolves.toBeUndefined(); + + expect((await events.findOne({ groupHash: 'broken-event' })).resolvedInRelease).toBeUndefined(); + expect((await events.findOne({ groupHash: 'valid-event' })).resolvedInRelease).toBe('b'); + }); + + test('should ignore a repetition without a release', async () => { + await releases.insertMany([ + createRelease('a', 72, true), + createRelease('b', 48), + ]); + await events.insertOne(createEvent('error-10', 'a')); + await repetitions.insertOne({ + groupHash: 'error-10', + timestamp: NOW_SECONDS, + }); + + await expect(validateReleases(db, NOW)).resolves.toBeUndefined(); + + expect((await events.findOne({ groupHash: 'error-10' })).resolvedInRelease).toBe('b'); + }); + + test('should ignore a repetition from an unknown release', async () => { + await releases.insertMany([ + createRelease('a', 72, true), + createRelease('b', 48), + ]); + await events.insertOne(createEvent('error-11', 'a')); + await repetitions.insertOne(createRepetition('error-11', 'unknown')); + + await expect(validateReleases(db, NOW)).resolves.toBeUndefined(); + + expect((await events.findOne({ groupHash: 'error-11' })).resolvedInRelease).toBe('b'); + }); + + test('should resolve an event again after its regression stops occurring', async () => { + await releases.insertMany([ + createRelease('a', 120, true), + createRelease('b', 96, true), + createRelease('c', 72, true), + createRelease('d', 48), + ]); + const event = createEvent('error-cycle', 'a', 'b'); + + event.regressionInRelease = 'c'; + await events.insertOne(event); + await repetitions.insertOne(createRepetition('error-cycle', 'c')); + + await validateReleases(db, NOW); + + const updatedEvent = await events.findOne({ groupHash: 'error-cycle' }); + + expect(updatedEvent.resolvedInRelease).toBe('d'); + expect(updatedEvent.regressionInRelease).toBe('c'); + }); + + test('should not resolve an event again before a release newer than its regression', async () => { + await releases.insertMany([ + createRelease('a', 96, true), + createRelease('b', 72, true), + createRelease('c', 48), + ]); + const event = createEvent('error-active-regression', 'a', 'b'); + + event.regressionInRelease = 'c'; + await events.insertOne(event); + await repetitions.insertOne(createRepetition('error-active-regression', 'c')); + + await validateReleases(db, NOW); + + expect((await events.findOne({ groupHash: 'error-active-regression' })).resolvedInRelease).toBe('b'); + }); + + test('should not select a release older than 30 days as a candidate', async () => { + await releases.insertMany([ + createRelease('a', 960, true), + createRelease('b', 744), + ]); + await events.insertOne(createEvent('error-12', 'a')); + + await validateReleases(db, NOW); + + expect((await events.findOne({ groupHash: 'error-12' })).resolvedInRelease).toBeUndefined(); + expect((await releases.findOne({ release: 'b' })).fixChecked).toBe(false); + }); + + test('should use a release older than 30 days as event history', async () => { + await releases.insertMany([ + createRelease('a', 960, true), + createRelease('b', 48), + ]); + await events.insertOne(createEvent('error-13', 'a')); + + await validateReleases(db, NOW); + + expect((await events.findOne({ groupHash: 'error-13' })).resolvedInRelease).toBe('b'); + }); +}); diff --git a/workers/task-manager/src/index.ts b/workers/task-manager/src/index.ts index d6dec032..cda3f338 100644 --- a/workers/task-manager/src/index.ts +++ b/workers/task-manager/src/index.ts @@ -6,9 +6,9 @@ import * as pkg from '../package.json'; import type { TaskManagerWorkerTask } from '../types/task-manager-worker-task'; import type { ProjectDBScheme, - GroupedEventDBScheme, - ProjectTaskManagerConfig + GroupedEventDBScheme } from '@hawk.so/types'; +import type { ProjectTaskManagerConfig } from '../types/project-task-manager-config'; import type { TaskManagerItem } from '@hawk.so/types/src/base/event/taskManagerItem'; import HawkCatcher from '@hawk.so/nodejs'; import { decodeUnsafeFields } from '../../../lib/utils/unsafeFields'; diff --git a/workers/task-manager/types/project-task-manager-config.ts b/workers/task-manager/types/project-task-manager-config.ts new file mode 100644 index 00000000..a255450f --- /dev/null +++ b/workers/task-manager/types/project-task-manager-config.ts @@ -0,0 +1,19 @@ +import type { ProjectTaskManagerConfig as ProjectTaskManagerConfigType } from '@hawk.so/types'; + +interface LegacyDelegatedUser { + accessToken: string; + accessTokenExpiresAt: Date | null; + refreshToken: string; + refreshTokenExpiresAt: Date | null; + status: 'active' | 'revoked' | 'missing'; +} + +/** + * Project task manager config with the legacy delegated user data still used + * by the API and task-manager worker. + */ +export type ProjectTaskManagerConfig = ProjectTaskManagerConfigType & { + config: ProjectTaskManagerConfigType['config'] & { + delegatedUser?: LegacyDelegatedUser; + }; +}; diff --git a/yarn.lock b/yarn.lock index 3172b2bc..31792312 100644 --- a/yarn.lock +++ b/yarn.lock @@ -402,10 +402,10 @@ dependencies: "@types/mongodb" "^3.5.34" -"@hawk.so/types@^0.5.9": - version "0.5.9" - resolved "https://registry.yarnpkg.com/@hawk.so/types/-/types-0.5.9.tgz#817e8b26283d0367371125f055f2e37a274797bc" - integrity sha512-86aE0Bdzvy8C+Dqd1iZpnDho44zLGX/t92SGuAv2Q52gjSJ7SHQdpGDWtM91FXncfT5uzAizl9jYMuE6Qrtm0Q== +"@hawk.so/types@^0.9.0": + version "0.9.0" + resolved "https://registry.yarnpkg.com/@hawk.so/types/-/types-0.9.0.tgz#19b68065e9bafa0f4ee14f22a71cbb946f502cbc" + integrity sha512-zY9yw83Dzw3RcbTtMR9AKj8yNEPwD1sZZCeLlw5onZr1+GdjfnR1AhLvoPeZXBbRPRLx3MzyeuM85lAl0do+DA== dependencies: bson "^7.0.0"