From 8f58256f24bd120375fd4a0301f96782e2d55ad6 Mon Sep 17 00:00:00 2001 From: alisawavezen12 Date: Thu, 17 Sep 2026 17:53:23 +0300 Subject: [PATCH 01/10] Add release validator worker --- package.json | 2 + workers/release-validator/README.md | 15 + workers/release-validator/package.json | 10 + workers/release-validator/src/index.ts | 46 +++ workers/release-validator/src/types.ts | 31 ++ .../src/utils/build-event-release-map.ts | 38 +++ .../src/utils/group-releases-by-project.ts | 19 ++ .../src/validate-releases.ts | 155 +++++++++ workers/release-validator/tests/index.test.ts | 15 + .../utils/build-event-release-map.test.ts | 88 ++++++ .../utils/group-releases-by-project.test.ts | 53 ++++ .../tests/validate-releases.test.ts | 298 ++++++++++++++++++ 12 files changed, 770 insertions(+) create mode 100644 workers/release-validator/README.md create mode 100644 workers/release-validator/package.json create mode 100644 workers/release-validator/src/index.ts create mode 100644 workers/release-validator/src/types.ts create mode 100644 workers/release-validator/src/utils/build-event-release-map.ts create mode 100644 workers/release-validator/src/utils/group-releases-by-project.ts create mode 100644 workers/release-validator/src/validate-releases.ts create mode 100644 workers/release-validator/tests/index.test.ts create mode 100644 workers/release-validator/tests/utils/build-event-release-map.test.ts create mode 100644 workers/release-validator/tests/utils/group-releases-by-project.test.ts create mode 100644 workers/release-validator/tests/validate-releases.test.ts diff --git a/package.json b/package.json index f2e1e154..d9d168b8 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", 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/types.ts b/workers/release-validator/src/types.ts new file mode 100644 index 00000000..ecd76364 --- /dev/null +++ b/workers/release-validator/src/types.ts @@ -0,0 +1,31 @@ +import { ObjectId } from 'mongodb'; + +/** + * Release data used during validation. + */ +export interface ReleaseRecord { + _id: ObjectId; + projectId: string; + release: string; + fixChecked?: boolean; +} + +/** + * Original event data used during validation. + */ +export interface EventRecord { + _id: ObjectId; + groupHash: string; + payload?: { + release?: string; + }; + resolvedInRelease?: string | null; +} + +/** + * Repetition data used during validation. + */ +export interface RepetitionRecord { + groupHash: string; + release?: string; +} 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..2563eaf7 --- /dev/null +++ b/workers/release-validator/src/utils/build-event-release-map.ts @@ -0,0 +1,38 @@ +import { EventRecord, RepetitionRecord } from '../types'; + +/** + * Build a map of releases in which each event occurred. + * + * @param events - original events + * @param repetitions - event repetitions + */ +export function buildEventReleaseMap( + events: EventRecord[], + repetitions: RepetitionRecord[] +): 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..6b79665b --- /dev/null +++ b/workers/release-validator/src/utils/group-releases-by-project.ts @@ -0,0 +1,19 @@ +import { ReleaseRecord } from '../types'; + +/** + * Group chronologically ordered releases by project. + * + * @param releases - releases ordered from oldest to newest + */ +export function groupReleasesByProject(releases: ReleaseRecord[]): 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..af126462 --- /dev/null +++ b/workers/release-validator/src/validate-releases.ts @@ -0,0 +1,155 @@ +import { Db, ObjectID } from 'mongodb'; +import { HOURS_IN_DAY, MINUTES_IN_HOUR, MS_IN_SEC, SECONDS_IN_MINUTE } from '../../../lib/utils/consts'; +import { EventRecord, ReleaseRecord, RepetitionRecord } from './types'; +import { buildEventReleaseMap } from './utils/build-event-release-map'; +import { groupReleasesByProject } from './utils/group-releases-by-project'; + +const RELEASE_OBSERVATION_PERIOD_SECONDS = HOURS_IN_DAY * MINUTES_IN_HOUR * SECONDS_IN_MINUTE; +const RELEASE_MAX_AGE_DAYS = 30; +const RELEASE_MAX_AGE_SECONDS = RELEASE_MAX_AGE_DAYS * RELEASE_OBSERVATION_PERIOD_SECONDS; + +/** + * 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); + 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, + }, + projectId: { + $type: 'string', + $ne: '', + }, + release: { + $type: 'string', + $ne: '', + }, + fixChecked: { $ne: true }, + }) + .sort({ _id: 1 }) + .toArray(); +} + +/** + * 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: ReleaseRecord[]): Promise { + const releasesCollection = db.collection('releases'); + const eventsCollection = db.collection(`events:${projectId}`); + const repetitionsCollection = db.collection(`repetitions:${projectId}`); + const allProjectReleases = await releasesCollection + .find({ + projectId, + release: { + $type: 'string', + $ne: '', + }, + }) + .sort({ _id: 1 }) + .toArray(); + const releasesByName = new Map(); + + for (const release of allProjectReleases) { + releasesByName.set(release.release, release); + } + + const events = await eventsCollection.find({ + groupHash: { + $type: 'string', + $ne: '', + }, + 'payload.release': { + $type: 'string', + $ne: '', + }, + $or: [ + { resolvedInRelease: { $exists: false } }, + { resolvedInRelease: null }, + ], + }).toArray(); + const eventGroupHashes = events.map(event => event.groupHash); + const repetitions = await repetitionsCollection.find({ + groupHash: { $in: eventGroupHashes }, + release: { + $type: 'string', + $ne: '', + }, + }).toArray(); + const eventReleases = buildEventReleaseMap(events, repetitions); + + for (const event of events) { + const originalReleaseName = event.payload.release; + const originalRelease = releasesByName.get(originalReleaseName); + + if (!originalRelease) { + continue; + } + + const releasesWithEvent = eventReleases.get(event.groupHash) || new Set(); + + for (const release of releasesToCheck) { + const releaseId = release._id.toHexString(); + const isNewerThanOriginal = releaseId > originalRelease._id.toHexString(); + const occurredInRelease = releasesWithEvent.has(release.release); + const occurredInNewerRelease = allProjectReleases.some(projectRelease => { + return projectRelease._id.toHexString() > releaseId && releasesWithEvent.has(projectRelease.release); + }); + + if (!isNewerThanOriginal || occurredInRelease || occurredInNewerRelease) { + continue; + } + + await eventsCollection.updateOne({ + _id: event._id, + $or: [ + { resolvedInRelease: { $exists: false } }, + { resolvedInRelease: null }, + ], + }, { + $set: { + resolvedInRelease: release.release, + }, + }); + + break; + } + } + + 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 { + 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/index.test.ts b/workers/release-validator/tests/index.test.ts new file mode 100644 index 00000000..96be5476 --- /dev/null +++ b/workers/release-validator/tests/index.test.ts @@ -0,0 +1,15 @@ +import '../../../env-test'; +import ReleaseValidatorWorker from '../src'; + +jest.mock('amqplib'); + +/** + * Release Validator worker smoke tests. + */ +describe('ReleaseValidatorWorker', () => { + test('should use the release validator queue', () => { + const worker = new ReleaseValidatorWorker(); + + expect(worker.type).toBe('cron-tasks/release-validator'); + }); +}); 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..9735eb96 --- /dev/null +++ b/workers/release-validator/tests/utils/build-event-release-map.test.ts @@ -0,0 +1,88 @@ +import { ObjectID } from 'mongodb'; +import { EventRecord, RepetitionRecord } from '../../src/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): EventRecord { + return { + _id: new ObjectID(), + groupHash, + payload: release ? { release } : {}, + }; +} + +/** + * Create an event repetition. + * + * @param groupHash - event group hash + * @param release - release in which the event occurred + */ +function createRepetition(groupHash: string, release?: string): RepetitionRecord { + return { + groupHash, + 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..5a22c503 --- /dev/null +++ b/workers/release-validator/tests/utils/group-releases-by-project.test.ts @@ -0,0 +1,53 @@ +import { ObjectID } from 'mongodb'; +import { ReleaseRecord } from '../../src/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): ReleaseRecord { + return { + _id: ObjectID.createFromTime(createdAtSeconds), + projectId, + release, + }; +} + +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: ReleaseRecord[] = []; + + // 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..b781daf8 --- /dev/null +++ b/workers/release-validator/tests/validate-releases.test.ts @@ -0,0 +1,298 @@ +import '../../../env-test'; +import { Collection, Db, MongoClient, ObjectID } from 'mongodb'; +import { validateReleases } from '../src/validate-releases'; +import { EventRecord, ReleaseRecord, RepetitionRecord } from '../src/types'; + +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): ReleaseRecord { + return { + _id: releaseId(hoursBeforeNow), + projectId: PROJECT_ID, + release, + 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): EventRecord { + const event: EventRecord = { + _id: new ObjectID(), + groupHash, + payload: { release }, + }; + + 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): RepetitionRecord { + return { + groupHash, + release, + }; +} + +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, + }, + 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: {}, + }, + 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', + }); + + 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 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'); + }); +}); From 7d752cbcf19afe70b73003125c0ad601fa0eae42 Mon Sep 17 00:00:00 2001 From: alisawavezen12 Date: Fri, 18 Sep 2026 13:26:42 +0300 Subject: [PATCH 02/10] Track event regressions by release --- package.json | 2 +- workers/grouper/src/index.ts | 22 ++++++ workers/grouper/src/mark-regression.ts | 57 ++++++++++++++ workers/grouper/tests/index.test.ts | 101 +++++++++++++++++++++++++ yarn.lock | 8 +- 5 files changed, 185 insertions(+), 5 deletions(-) create mode 100644 workers/grouper/src/mark-regression.ts diff --git a/package.json b/package.json index d9d168b8..f01b42db 100644 --- a/package.json +++ b/package.json @@ -59,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.8.0", "@types/amqplib": "^0.8.2", "@types/jest": "^29.5.14", "@types/mongodb": "^3.5.15", diff --git a/workers/grouper/src/index.ts b/workers/grouper/src/index.ts index 6750d768..73fa0c33 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 { markRegression } from './mark-regression'; /** * eslint does not count decorators as a variable usage @@ -343,10 +344,31 @@ 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 && !existedEvent.regressionInRelease) { + try { + await markRegression( + this.eventsDb.getConnection(), + task.projectId, + uniqueEventHash, + task.payload.release, + existedEvent.resolvedInRelease + ); + } catch (error) { + this.logger.error( + `[markRegression] 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/src/mark-regression.ts b/workers/grouper/src/mark-regression.ts new file mode 100644 index 00000000..bc5aa743 --- /dev/null +++ b/workers/grouper/src/mark-regression.ts @@ -0,0 +1,57 @@ +import { Db, ObjectID } from 'mongodb'; + +interface ReleaseRecord { + _id: ObjectID; + projectId: string; + release: string; +} + +/** + * Mark a resolved event as regressed 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 + */ +export async function markRegression( + db: Db, + projectId: string, + groupHash: string, + release: string, + resolvedInRelease: string +): Promise { + const releases = await db.collection('releases').find({ + projectId, + release: { + $in: [resolvedInRelease, release], + }, + }) + .toArray(); + const resolvedRelease = releases.find(item => item.release === resolvedInRelease); + const repetitionRelease = releases.find(item => item.release === release); + + if (!resolvedRelease || !repetitionRelease) { + return; + } + + const isResolvedOrNewerRelease = repetitionRelease._id.toHexString() >= resolvedRelease._id.toHexString(); + + if (!isResolvedOrNewerRelease) { + return; + } + + await db.collection(`events:${projectId}`).updateOne({ + groupHash, + resolvedInRelease, + regressionInRelease: { $exists: false }, + }, { + $set: { + regressionInRelease: release, + }, + }); +} diff --git a/workers/grouper/tests/index.test.ts b/workers/grouper/tests/index.test.ts index e4f2b3bd..5e494100 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,106 @@ 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(); + }); + + 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 a resolved event as regressed in the resolved 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 in an older release', 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 not overwrite the first regression release', async () => { + 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/yarn.lock b/yarn.lock index 3172b2bc..924cb365 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.8.0": + version "0.8.0" + resolved "https://registry.yarnpkg.com/@hawk.so/types/-/types-0.8.0.tgz#4d682f0c8df1e857d08a2af1a5ecb0ec8a5fcbcd" + integrity sha512-nfpu40G6Gj7woGktAronNAXfTA8k940R5gUqk2NIpNH9hb8oRB/HJRselRdeVxTnBb+3XT7TBC52TviaEis+0w== dependencies: bson "^7.0.0" From 8ac4dab986f3b4b7bcfb6006fad45146eafe3a24 Mon Sep 17 00:00:00 2001 From: alisawavezen12 Date: Fri, 18 Sep 2026 18:07:16 +0300 Subject: [PATCH 03/10] chore --- package.json | 2 +- workers/release-validator/src/types.ts | 31 ------------- .../src/utils/build-event-release-map.ts | 6 +-- .../src/utils/group-releases-by-project.ts | 6 +-- .../src/validate-releases.ts | 16 +++---- .../utils/build-event-release-map.test.ts | 19 +++++--- .../utils/group-releases-by-project.test.ts | 7 +-- .../tests/validate-releases.test.ts | 45 +++++++++++++------ yarn.lock | 8 ++-- 9 files changed, 69 insertions(+), 71 deletions(-) delete mode 100644 workers/release-validator/src/types.ts diff --git a/package.json b/package.json index f01b42db..32953a73 100644 --- a/package.json +++ b/package.json @@ -59,7 +59,7 @@ "@babel/parser": "^7.26.9", "@babel/traverse": "7.26.9", "@hawk.so/nodejs": "^3.1.1", - "@hawk.so/types": "^0.8.0", + "@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/release-validator/src/types.ts b/workers/release-validator/src/types.ts deleted file mode 100644 index ecd76364..00000000 --- a/workers/release-validator/src/types.ts +++ /dev/null @@ -1,31 +0,0 @@ -import { ObjectId } from 'mongodb'; - -/** - * Release data used during validation. - */ -export interface ReleaseRecord { - _id: ObjectId; - projectId: string; - release: string; - fixChecked?: boolean; -} - -/** - * Original event data used during validation. - */ -export interface EventRecord { - _id: ObjectId; - groupHash: string; - payload?: { - release?: string; - }; - resolvedInRelease?: string | null; -} - -/** - * Repetition data used during validation. - */ -export interface RepetitionRecord { - groupHash: string; - release?: string; -} diff --git a/workers/release-validator/src/utils/build-event-release-map.ts b/workers/release-validator/src/utils/build-event-release-map.ts index 2563eaf7..26c3022f 100644 --- a/workers/release-validator/src/utils/build-event-release-map.ts +++ b/workers/release-validator/src/utils/build-event-release-map.ts @@ -1,4 +1,4 @@ -import { EventRecord, RepetitionRecord } from '../types'; +import type { GroupedEventDBScheme, RepetitionDBScheme } from '@hawk.so/types'; /** * Build a map of releases in which each event occurred. @@ -7,8 +7,8 @@ import { EventRecord, RepetitionRecord } from '../types'; * @param repetitions - event repetitions */ export function buildEventReleaseMap( - events: EventRecord[], - repetitions: RepetitionRecord[] + events: GroupedEventDBScheme[], + repetitions: RepetitionDBScheme[] ): Map> { const releasesByGroupHash = new Map>(); diff --git a/workers/release-validator/src/utils/group-releases-by-project.ts b/workers/release-validator/src/utils/group-releases-by-project.ts index 6b79665b..efd30ae3 100644 --- a/workers/release-validator/src/utils/group-releases-by-project.ts +++ b/workers/release-validator/src/utils/group-releases-by-project.ts @@ -1,12 +1,12 @@ -import { ReleaseRecord } from '../types'; +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: ReleaseRecord[]): Map { - const releasesByProject = new Map(); +export function groupReleasesByProject(releases: ReleaseDBScheme[]): Map { + const releasesByProject = new Map(); for (const release of releases) { const projectReleases = releasesByProject.get(release.projectId) || []; diff --git a/workers/release-validator/src/validate-releases.ts b/workers/release-validator/src/validate-releases.ts index af126462..80524528 100644 --- a/workers/release-validator/src/validate-releases.ts +++ b/workers/release-validator/src/validate-releases.ts @@ -1,6 +1,6 @@ import { 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 { EventRecord, ReleaseRecord, RepetitionRecord } from './types'; import { buildEventReleaseMap } from './utils/build-event-release-map'; import { groupReleasesByProject } from './utils/group-releases-by-project'; @@ -14,12 +14,12 @@ const RELEASE_MAX_AGE_SECONDS = RELEASE_MAX_AGE_DAYS * RELEASE_OBSERVATION_PERIO * @param db - events database connection * @param now - current time */ -async function findReleasesToCheck(db: Db, now: Date): Promise { +async function findReleasesToCheck(db: Db, now: Date): Promise { const nowSeconds = Math.floor(now.getTime() / MS_IN_SEC); const oldestReleaseId = ObjectID.createFromTime(nowSeconds - RELEASE_MAX_AGE_SECONDS); const newestReleaseId = ObjectID.createFromTime(nowSeconds - RELEASE_OBSERVATION_PERIOD_SECONDS); - return db.collection('releases') + return db.collection('releases') .find({ _id: { $gte: oldestReleaseId, @@ -46,10 +46,10 @@ async function findReleasesToCheck(db: Db, now: Date): Promise * @param projectId - project identifier * @param releasesToCheck - ready releases ordered from oldest to newest */ -async function validateProject(db: Db, projectId: string, releasesToCheck: ReleaseRecord[]): Promise { - const releasesCollection = db.collection('releases'); - const eventsCollection = db.collection(`events:${projectId}`); - const repetitionsCollection = db.collection(`repetitions:${projectId}`); +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}`); const allProjectReleases = await releasesCollection .find({ projectId, @@ -60,7 +60,7 @@ async function validateProject(db: Db, projectId: string, releasesToCheck: Relea }) .sort({ _id: 1 }) .toArray(); - const releasesByName = new Map(); + const releasesByName = new Map(); for (const release of allProjectReleases) { releasesByName.set(release.release, release); 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 index 9735eb96..a9d6025f 100644 --- a/workers/release-validator/tests/utils/build-event-release-map.test.ts +++ b/workers/release-validator/tests/utils/build-event-release-map.test.ts @@ -1,5 +1,5 @@ import { ObjectID } from 'mongodb'; -import { EventRecord, RepetitionRecord } from '../../src/types'; +import type { GroupedEventDBScheme, RepetitionDBScheme } from '@hawk.so/types'; import { buildEventReleaseMap } from '../../src/utils/build-event-release-map'; /** @@ -8,11 +8,19 @@ import { buildEventReleaseMap } from '../../src/utils/build-event-release-map'; * @param groupHash - event group hash * @param release - release in which the event first occurred */ -function createEvent(groupHash: string, release?: string): EventRecord { +function createEvent(groupHash: string, release?: string): GroupedEventDBScheme { return { _id: new ObjectID(), groupHash, - payload: release ? { release } : {}, + payload: { + title: groupHash, + ...(release ? { release } : {}), + }, + totalCount: 1, + catcherType: 'errors/default', + usersAffected: 0, + visitedBy: [], + timestamp: 1, }; } @@ -22,10 +30,11 @@ function createEvent(groupHash: string, release?: string): EventRecord { * @param groupHash - event group hash * @param release - release in which the event occurred */ -function createRepetition(groupHash: string, release?: string): RepetitionRecord { +function createRepetition(groupHash: string, release?: string): RepetitionDBScheme { return { groupHash, - release, + timestamp: 1, + ...(release ? { 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 index 5a22c503..0e53fe78 100644 --- a/workers/release-validator/tests/utils/group-releases-by-project.test.ts +++ b/workers/release-validator/tests/utils/group-releases-by-project.test.ts @@ -1,5 +1,5 @@ import { ObjectID } from 'mongodb'; -import { ReleaseRecord } from '../../src/types'; +import type { ReleaseDBScheme } from '@hawk.so/types'; import { groupReleasesByProject } from '../../src/utils/group-releases-by-project'; /** @@ -9,11 +9,12 @@ import { groupReleasesByProject } from '../../src/utils/group-releases-by-projec * @param release - release name * @param createdAtSeconds - release creation time */ -function createRelease(projectId: string, release: string, createdAtSeconds: number): ReleaseRecord { +function createRelease(projectId: string, release: string, createdAtSeconds: number): ReleaseDBScheme { return { _id: ObjectID.createFromTime(createdAtSeconds), projectId, release, + commits: [], }; } @@ -42,7 +43,7 @@ describe('groupReleasesByProject', () => { test('should return an empty map for an empty release list', () => { // Arrange - const releases: ReleaseRecord[] = []; + const releases: ReleaseDBScheme[] = []; // Act const result = groupReleasesByProject(releases); diff --git a/workers/release-validator/tests/validate-releases.test.ts b/workers/release-validator/tests/validate-releases.test.ts index b781daf8..327d8d8a 100644 --- a/workers/release-validator/tests/validate-releases.test.ts +++ b/workers/release-validator/tests/validate-releases.test.ts @@ -1,7 +1,7 @@ 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'; -import { EventRecord, ReleaseRecord, RepetitionRecord } from '../src/types'; const PROJECT_ID = 'release-validator-project'; const HOUR_IN_SECONDS = 60 * 60; @@ -24,11 +24,12 @@ function releaseId(hoursBeforeNow: number): ObjectID { * @param hoursBeforeNow - release age in hours * @param fixChecked - whether the release has already been checked */ -function createRelease(release: string, hoursBeforeNow: number, fixChecked = false): ReleaseRecord { +function createRelease(release: string, hoursBeforeNow: number, fixChecked = false): ReleaseDBScheme { return { _id: releaseId(hoursBeforeNow), projectId: PROJECT_ID, release, + commits: [], fixChecked, }; } @@ -40,11 +41,19 @@ function createRelease(release: string, hoursBeforeNow: number, fixChecked = fal * @param release - release in which the event first occurred * @param resolvedInRelease - existing resolved release */ -function createEvent(groupHash: string, release: string, resolvedInRelease?: string): EventRecord { - const event: EventRecord = { +function createEvent(groupHash: string, release: string, resolvedInRelease?: string): GroupedEventDBScheme { + const event: GroupedEventDBScheme = { _id: new ObjectID(), groupHash, - payload: { release }, + payload: { + title: groupHash, + release, + }, + totalCount: 1, + catcherType: 'errors/default', + usersAffected: 0, + visitedBy: [], + timestamp: NOW_SECONDS, }; if (resolvedInRelease !== undefined) { @@ -60,19 +69,20 @@ function createEvent(groupHash: string, release: string, resolvedInRelease?: str * @param groupHash - event group hash * @param release - release in which the event occurred */ -function createRepetition(groupHash: string, release: string): RepetitionRecord { +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; + let releases: Collection; + let events: Collection; + let repetitions: Collection; beforeAll(async () => { connection = await MongoClient.connect(process.env.MONGO_EVENTS_DATABASE_URI, { @@ -80,9 +90,9 @@ describe('validateReleases', () => { useUnifiedTopology: true, }); db = connection.db(); - releases = db.collection('releases'); - events = db.collection(`events:${PROJECT_ID}`); - repetitions = db.collection(`repetitions:${PROJECT_ID}`); + releases = db.collection('releases'); + events = db.collection(`events:${PROJECT_ID}`); + repetitions = db.collection(`repetitions:${PROJECT_ID}`); }); beforeEach(async () => { @@ -211,6 +221,7 @@ describe('validateReleases', () => { _id: releaseId(72), projectId: PROJECT_ID, release: 123 as unknown as string, + commits: [], }, createRelease('d', 48), ]); @@ -232,7 +243,14 @@ describe('validateReleases', () => { { _id: new ObjectID(), groupHash: 'broken-event', - payload: {}, + payload: { + title: 'broken-event', + }, + totalCount: 1, + catcherType: 'errors/default', + usersAffected: 0, + visitedBy: [], + timestamp: NOW_SECONDS, }, createEvent('valid-event', 'a'), ]); @@ -251,6 +269,7 @@ describe('validateReleases', () => { await events.insertOne(createEvent('error-10', 'a')); await repetitions.insertOne({ groupHash: 'error-10', + timestamp: NOW_SECONDS, }); await expect(validateReleases(db, NOW)).resolves.toBeUndefined(); diff --git a/yarn.lock b/yarn.lock index 924cb365..31792312 100644 --- a/yarn.lock +++ b/yarn.lock @@ -402,10 +402,10 @@ dependencies: "@types/mongodb" "^3.5.34" -"@hawk.so/types@^0.8.0": - version "0.8.0" - resolved "https://registry.yarnpkg.com/@hawk.so/types/-/types-0.8.0.tgz#4d682f0c8df1e857d08a2af1a5ecb0ec8a5fcbcd" - integrity sha512-nfpu40G6Gj7woGktAronNAXfTA8k940R5gUqk2NIpNH9hb8oRB/HJRselRdeVxTnBb+3XT7TBC52TviaEis+0w== +"@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" From d8164f75b86a09347b562e20876df68f9a26dc5d Mon Sep 17 00:00:00 2001 From: alisawavezen12 Date: Fri, 18 Sep 2026 18:23:10 +0300 Subject: [PATCH 04/10] Support repeated resolution and regression cycles --- workers/grouper/src/index.ts | 5 +- workers/grouper/src/mark-regression.ts | 20 +++++-- workers/grouper/tests/index.test.ts | 50 ++++++++++++++++- .../src/validate-releases.ts | 54 +++++++++++++++---- .../tests/validate-releases.test.ts | 38 +++++++++++++ 5 files changed, 150 insertions(+), 17 deletions(-) diff --git a/workers/grouper/src/index.ts b/workers/grouper/src/index.ts index 73fa0c33..64862f39 100644 --- a/workers/grouper/src/index.ts +++ b/workers/grouper/src/index.ts @@ -352,14 +352,15 @@ export default class GrouperWorker extends Worker { return this.saveRepetition(task.projectId, newRepetition); }); - if (task.payload.release && existedEvent.resolvedInRelease && !existedEvent.regressionInRelease) { + if (task.payload.release && existedEvent.resolvedInRelease) { try { await markRegression( this.eventsDb.getConnection(), task.projectId, uniqueEventHash, task.payload.release, - existedEvent.resolvedInRelease + existedEvent.resolvedInRelease, + existedEvent.regressionInRelease ); } catch (error) { this.logger.error( diff --git a/workers/grouper/src/mark-regression.ts b/workers/grouper/src/mark-regression.ts index bc5aa743..80c25803 100644 --- a/workers/grouper/src/mark-regression.ts +++ b/workers/grouper/src/mark-regression.ts @@ -17,38 +17,48 @@ interface ReleaseRecord { * @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 markRegression( db: Db, projectId: string, groupHash: string, release: string, - resolvedInRelease: string + resolvedInRelease: string, + regressionInRelease?: string ): Promise { const releases = await db.collection('releases').find({ projectId, release: { - $in: [resolvedInRelease, 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 isResolvedOrNewerRelease = repetitionRelease._id.toHexString() >= resolvedRelease._id.toHexString(); + const resolvedReleaseId = resolvedRelease._id.toHexString(); + const isResolvedOrNewerRelease = repetitionRelease._id.toHexString() >= resolvedReleaseId; + const hasRegressionForCurrentCycle = previousRegressionRelease && + previousRegressionRelease._id.toHexString() >= resolvedReleaseId; - if (!isResolvedOrNewerRelease) { + if (!isResolvedOrNewerRelease || hasRegressionForCurrentCycle) { return; } await db.collection(`events:${projectId}`).updateOne({ groupHash, resolvedInRelease, - regressionInRelease: { $exists: false }, + ...(regressionInRelease + ? { regressionInRelease } + : { regressionInRelease: { $exists: false } }), }, { $set: { regressionInRelease: release, diff --git a/workers/grouper/tests/index.test.ts b/workers/grouper/tests/index.test.ts index 5e494100..fe88f1d7 100644 --- a/workers/grouper/tests/index.test.ts +++ b/workers/grouper/tests/index.test.ts @@ -478,7 +478,55 @@ describe('GrouperWorker', () => { expect((await eventsCollection.findOne({})).regressionInRelease).toBeUndefined(); }); - test('Should not overwrite the first regression release', async () => { + 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 regression from the current resolution cycle', 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: { diff --git a/workers/release-validator/src/validate-releases.ts b/workers/release-validator/src/validate-releases.ts index 80524528..2a81f343 100644 --- a/workers/release-validator/src/validate-releases.ts +++ b/workers/release-validator/src/validate-releases.ts @@ -78,6 +78,16 @@ async function validateProject(db: Db, projectId: string, releasesToCheck: Relea $or: [ { resolvedInRelease: { $exists: false } }, { resolvedInRelease: null }, + { + resolvedInRelease: { + $type: 'string', + $ne: '', + }, + regressionInRelease: { + $type: 'string', + $ne: '', + }, + }, ], }).toArray(); const eventGroupHashes = events.map(event => event.groupHash); @@ -91,10 +101,27 @@ async function validateProject(db: Db, projectId: string, releasesToCheck: Relea const eventReleases = buildEventReleaseMap(events, repetitions); for (const event of events) { - const originalReleaseName = event.payload.release; - const originalRelease = releasesByName.get(originalReleaseName); + const originalRelease = releasesByName.get(event.payload.release); + let lastOccurrenceRelease = originalRelease; + + if (event.resolvedInRelease && event.regressionInRelease) { + const resolvedRelease = releasesByName.get(event.resolvedInRelease); + const regressionRelease = releasesByName.get(event.regressionInRelease); + + if (!resolvedRelease || !regressionRelease) { + continue; + } - if (!originalRelease) { + const isCurrentlyRegressed = regressionRelease._id.toHexString() >= resolvedRelease._id.toHexString(); + + if (!isCurrentlyRegressed) { + continue; + } + + lastOccurrenceRelease = regressionRelease; + } + + if (!lastOccurrenceRelease) { continue; } @@ -102,22 +129,31 @@ async function validateProject(db: Db, projectId: string, releasesToCheck: Relea for (const release of releasesToCheck) { const releaseId = release._id.toHexString(); - const isNewerThanOriginal = releaseId > originalRelease._id.toHexString(); + const isNewerThanLastOccurrence = releaseId > lastOccurrenceRelease._id.toHexString(); const occurredInRelease = releasesWithEvent.has(release.release); const occurredInNewerRelease = allProjectReleases.some(projectRelease => { return projectRelease._id.toHexString() > releaseId && releasesWithEvent.has(projectRelease.release); }); - if (!isNewerThanOriginal || occurredInRelease || occurredInNewerRelease) { + if (!isNewerThanLastOccurrence || occurredInRelease || occurredInNewerRelease) { continue; } + const eventState = event.regressionInRelease + ? { + resolvedInRelease: event.resolvedInRelease, + regressionInRelease: event.regressionInRelease, + } + : { + $or: [ + { resolvedInRelease: { $exists: false } }, + { resolvedInRelease: null }, + ], + }; + await eventsCollection.updateOne({ _id: event._id, - $or: [ - { resolvedInRelease: { $exists: false } }, - { resolvedInRelease: null }, - ], + ...eventState, }, { $set: { resolvedInRelease: release.release, diff --git a/workers/release-validator/tests/validate-releases.test.ts b/workers/release-validator/tests/validate-releases.test.ts index 327d8d8a..332c12a2 100644 --- a/workers/release-validator/tests/validate-releases.test.ts +++ b/workers/release-validator/tests/validate-releases.test.ts @@ -290,6 +290,44 @@ describe('validateReleases', () => { 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), From 64555302dbc069f0c34e5843daa62e76ffcf3fcd Mon Sep 17 00:00:00 2001 From: alisawavezen12 Date: Sun, 20 Sep 2026 12:29:20 +0300 Subject: [PATCH 05/10] validate Events Batch --- .../src/utils/build-event-release-map.ts | 3 +- .../src/validate-releases.ts | 245 ++++++++++++++---- workers/task-manager/src/index.ts | 4 +- .../types/project-task-manager-config.ts | 19 ++ 4 files changed, 220 insertions(+), 51 deletions(-) create mode 100644 workers/task-manager/types/project-task-manager-config.ts diff --git a/workers/release-validator/src/utils/build-event-release-map.ts b/workers/release-validator/src/utils/build-event-release-map.ts index 26c3022f..2c0a0eb5 100644 --- a/workers/release-validator/src/utils/build-event-release-map.ts +++ b/workers/release-validator/src/utils/build-event-release-map.ts @@ -1,7 +1,8 @@ import type { GroupedEventDBScheme, RepetitionDBScheme } from '@hawk.so/types'; /** - * Build a map of releases in which each event occurred. + * Build a lookup set of releases in which each event occurred. + * Release ordering is handled separately by the validation flow. * * @param events - original events * @param repetitions - event repetitions diff --git a/workers/release-validator/src/validate-releases.ts b/workers/release-validator/src/validate-releases.ts index 2a81f343..e79b1349 100644 --- a/workers/release-validator/src/validate-releases.ts +++ b/workers/release-validator/src/validate-releases.ts @@ -1,4 +1,4 @@ -import { Db, ObjectID } from 'mongodb'; +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'; @@ -8,6 +8,15 @@ const RELEASE_OBSERVATION_PERIOD_SECONDS = HOURS_IN_DAY * MINUTES_IN_HOUR * SECO const RELEASE_MAX_AGE_DAYS = 30; 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. * @@ -16,6 +25,11 @@ const RELEASE_MAX_AGE_SECONDS = RELEASE_MAX_AGE_DAYS * RELEASE_OBSERVATION_PERIO */ 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); @@ -40,56 +54,27 @@ async function findReleasesToCheck(db: Db, now: Date): Promise { - const releasesCollection = db.collection('releases'); - const eventsCollection = db.collection(`events:${projectId}`); - const repetitionsCollection = db.collection(`repetitions:${projectId}`); - const allProjectReleases = await releasesCollection - .find({ - projectId, - release: { - $type: 'string', - $ne: '', - }, - }) - .sort({ _id: 1 }) - .toArray(); - const releasesByName = new Map(); - - for (const release of allProjectReleases) { - releasesByName.set(release.release, release); - } - - const events = await 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: '', - }, - }, - ], - }).toArray(); +async function validateEventsBatch( + events: GroupedEventDBScheme[], + eventsCollection: Collection, + repetitionsCollection: Collection, + releasesToCheck: ReleaseDBScheme[], + allProjectReleases: ReleaseDBScheme[], + releasesByName: Map +): Promise { const eventGroupHashes = events.map(event => event.groupHash); const repetitions = await repetitionsCollection.find({ groupHash: { $in: eventGroupHashes }, @@ -97,23 +82,48 @@ async function validateProject(db: Db, projectId: string, releasesToCheck: Relea $type: 'string', $ne: '', }, + }, { + projection: { + _id: 0, + groupHash: 1, + release: 1, + }, }).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; + /** + * For a regressed event, continue validation from the regression release + * instead of its original release. This enables repeated resolve cycles. + */ 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; } @@ -121,6 +131,10 @@ async function validateProject(db: Db, projectId: string, releasesToCheck: Relea lastOccurrenceRelease = regressionRelease; } + /** + * Skip events whose original or latest occurrence release is missing from + * the project release history. + */ if (!lastOccurrenceRelease) { continue; } @@ -131,6 +145,11 @@ async function validateProject(db: Db, projectId: string, releasesToCheck: Relea 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); }); @@ -139,6 +158,10 @@ async function validateProject(db: Db, projectId: string, releasesToCheck: Relea continue; } + /** + * Keep the state transition atomic. A concurrent Grouper update must make + * this conditional update miss instead of overwriting newer state. + */ const eventState = event.regressionInRelease ? { resolvedInRelease: event.resolvedInRelease, @@ -160,10 +183,132 @@ async function validateProject(db: Db, projectId: string, releasesToCheck: Relea }, }); + /** + * 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 complete release history. Releases outside the validation window + * still define occurrence order and can block an incorrect resolution. + */ + const allProjectReleases = await releasesCollection + .find({ + projectId, + release: { + $type: 'string', + $ne: '', + }, + }) + .sort({ _id: 1 }) + .toArray(); + + /** + * Resolve release names to their Mongo records so ObjectIds can be used as + * the chronological source of truth instead of comparing version strings. + */ + const releasesByName = new Map(); + + for (const release of allProjectReleases) { + releasesByName.set(release.release, release); + } + + /** + * Stream eligible events from MongoDB instead of materializing the complete + * project result. 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), @@ -182,6 +327,10 @@ async function validateProject(db: Db, projectId: string, releasesToCheck: Relea * @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); 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; + }; +}; From e77e4a3709dcdb19e82032e0aca003473c2bee54 Mon Sep 17 00:00:00 2001 From: alisawavezen12 Date: Mon, 21 Sep 2026 12:15:22 +0300 Subject: [PATCH 06/10] Add release validator to Docker Compose --- docker-compose.dev.yml | 15 +++++++++++++++ docker-compose.prod.yml | 12 +++++++++++- 2 files changed, 26 insertions(+), 1 deletion(-) 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 # From 60df83740e9c5754b02ca5089c59d786499d2401 Mon Sep 17 00:00:00 2001 From: alisawavezen12 Date: Fri, 25 Sep 2026 07:16:34 +0300 Subject: [PATCH 07/10] chore --- ...ession.ts => check-and-mark-regression.ts} | 19 ++++++++++--------- workers/grouper/src/index.ts | 6 +++--- workers/grouper/tests/index.test.ts | 2 +- 3 files changed, 14 insertions(+), 13 deletions(-) rename workers/grouper/src/{mark-regression.ts => check-and-mark-regression.ts} (77%) diff --git a/workers/grouper/src/mark-regression.ts b/workers/grouper/src/check-and-mark-regression.ts similarity index 77% rename from workers/grouper/src/mark-regression.ts rename to workers/grouper/src/check-and-mark-regression.ts index 80c25803..e5084a49 100644 --- a/workers/grouper/src/mark-regression.ts +++ b/workers/grouper/src/check-and-mark-regression.ts @@ -1,13 +1,10 @@ -import { Db, ObjectID } from 'mongodb'; +import { Db } from 'mongodb'; +import type { ReleaseDBScheme } from '@hawk.so/types'; -interface ReleaseRecord { - _id: ObjectID; - projectId: string; - release: string; -} +type ReleaseRecordPart = Pick; /** - * Mark a resolved event as regressed in the resolved or a newer release. + * 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. @@ -19,7 +16,7 @@ interface ReleaseRecord { * @param resolvedInRelease - release in which the event was resolved * @param regressionInRelease - regression from a previous resolution cycle */ -export async function markRegression( +export async function checkAndMarkRegression( db: Db, projectId: string, groupHash: string, @@ -27,7 +24,7 @@ export async function markRegression( resolvedInRelease: string, regressionInRelease?: string ): Promise { - const releases = await db.collection('releases').find({ + const releases = await db.collection('releases').find({ projectId, release: { $in: [resolvedInRelease, release, regressionInRelease].filter(Boolean), @@ -46,6 +43,10 @@ export async function markRegression( 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; diff --git a/workers/grouper/src/index.ts b/workers/grouper/src/index.ts index 64862f39..e6839c39 100644 --- a/workers/grouper/src/index.ts +++ b/workers/grouper/src/index.ts @@ -24,7 +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 { markRegression } from './mark-regression'; +import { checkAndMarkRegression } from './check-and-mark-regression'; /** * eslint does not count decorators as a variable usage @@ -354,7 +354,7 @@ export default class GrouperWorker extends Worker { if (task.payload.release && existedEvent.resolvedInRelease) { try { - await markRegression( + await checkAndMarkRegression( this.eventsDb.getConnection(), task.projectId, uniqueEventHash, @@ -364,7 +364,7 @@ export default class GrouperWorker extends Worker { ); } catch (error) { this.logger.error( - `[markRegression] project=${task.projectId} groupHash=${uniqueEventHash} release=${task.payload.release}`, + `[checkAndMarkRegression] project=${task.projectId} groupHash=${uniqueEventHash} release=${task.payload.release}`, error ); } diff --git a/workers/grouper/tests/index.test.ts b/workers/grouper/tests/index.test.ts index fe88f1d7..219d3099 100644 --- a/workers/grouper/tests/index.test.ts +++ b/workers/grouper/tests/index.test.ts @@ -435,7 +435,7 @@ describe('GrouperWorker', () => { expect((await eventsCollection.findOne({})).regressionInRelease).toBe('release-c'); }); - test('Should mark a resolved event as regressed in the resolved release', async () => { + 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, From 1ee7b0171f11f3969cdac010062f7f044e74a41d Mon Sep 17 00:00:00 2001 From: alisawavezen12 Date: Fri, 25 Sep 2026 07:25:49 +0300 Subject: [PATCH 08/10] chore --- workers/grouper/tests/index.test.ts | 222 +++++++++--------- .../src/utils/build-event-release-map.ts | 12 +- .../src/validate-releases.ts | 51 ++-- 3 files changed, 157 insertions(+), 128 deletions(-) diff --git a/workers/grouper/tests/index.test.ts b/workers/grouper/tests/index.test.ts index 219d3099..6e2421c3 100644 --- a/workers/grouper/tests/index.test.ts +++ b/workers/grouper/tests/index.test.ts @@ -410,134 +410,136 @@ describe('GrouperWorker', () => { expect(savedRepetition.release).toBeUndefined(); }); - 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' })); + 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', + }, + }); - expect((await eventsCollection.findOne({})).regressionInRelease).toBe('release-c'); - }); + await worker.handle(generateTask({ release: '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', - }, + expect((await eventsCollection.findOne({})).regressionInRelease).toBe('release-c'); }); - await worker.handle(generateTask({ release: 'release-b' })); - - expect((await eventsCollection.findOne({})).regressionInRelease).toBe('release-b'); - }); - - test('Should not mark regression in an older release', async () => { - await connection.db().collection('releases').insertMany([ - { + 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-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' })); + await eventsCollection.updateOne({}, { + $set: { + resolvedInRelease: 'release-b', + }, + }); + + await worker.handle(generateTask({ release: 'release-b' })); + + expect((await eventsCollection.findOne({})).regressionInRelease).toBe('release-b'); }); - await worker.handle(generateTask({ release: 'release-a' })); + 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', + }, + }); - expect((await eventsCollection.findOne({})).regressionInRelease).toBeUndefined(); - }); + await worker.handle(generateTask({ release: 'release-a' })); - 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', - }, + expect((await eventsCollection.findOne({})).regressionInRelease).toBeUndefined(); }); - await worker.handle(generateTask({ release: 'release-d' })); + 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', + }, + }); - expect((await eventsCollection.findOne({})).regressionInRelease).toBe('release-d'); - }); + await worker.handle(generateTask({ release: 'release-d' })); - test('Should not overwrite the regression from the current resolution cycle', 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', - }, + expect((await eventsCollection.findOne({})).regressionInRelease).toBe('release-d'); }); - await worker.handle(generateTask({ release: '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', + }, + }); - expect((await eventsCollection.findOne({})).regressionInRelease).toBe('release-c'); + await worker.handle(generateTask({ release: 'release-d' })); + + expect((await eventsCollection.findOne({})).regressionInRelease).toBe('release-c'); + }); }); test('Should stringify delta', async () => { diff --git a/workers/release-validator/src/utils/build-event-release-map.ts b/workers/release-validator/src/utils/build-event-release-map.ts index 2c0a0eb5..13db2382 100644 --- a/workers/release-validator/src/utils/build-event-release-map.ts +++ b/workers/release-validator/src/utils/build-event-release-map.ts @@ -1,15 +1,21 @@ 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 - * @param repetitions - event repetitions + * @param events - original events from the current validation batch + * @param repetitions - unique event and release pairs from MongoDB */ export function buildEventReleaseMap( events: GroupedEventDBScheme[], - repetitions: RepetitionDBScheme[] + repetitions: RepetitionRelease[] ): Map> { const releasesByGroupHash = new Map>(); diff --git a/workers/release-validator/src/validate-releases.ts b/workers/release-validator/src/validate-releases.ts index e79b1349..476a523a 100644 --- a/workers/release-validator/src/validate-releases.ts +++ b/workers/release-validator/src/validate-releases.ts @@ -4,8 +4,20 @@ import { HOURS_IN_DAY, MINUTES_IN_HOUR, MS_IN_SEC, SECONDS_IN_MINUTE } from '../ import { buildEventReleaseMap } from './utils/build-event-release-map'; import { groupReleasesByProject } from './utils/group-releases-by-project'; +/** + * 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; /** @@ -39,10 +51,6 @@ async function findReleasesToCheck(db: Db, now: Date): Promise ): Promise { const eventGroupHashes = events.map(event => event.groupHash); - const repetitions = await repetitionsCollection.find({ - groupHash: { $in: eventGroupHashes }, - release: { - $type: 'string', - $ne: '', + const repetitions = await repetitionsCollection.aggregate>([ + { + $match: { + groupHash: { $in: eventGroupHashes }, + release: { + $type: 'string', + $ne: '', + }, + }, }, - }, { - projection: { - _id: 0, - groupHash: 1, - release: 1, + { + $group: { + _id: { + groupHash: '$groupHash', + release: '$release', + }, + }, + }, + { + $project: { + _id: 0, + groupHash: '$_id.groupHash', + release: '$_id.release', + }, }, - }).toArray(); + ]).toArray(); const eventReleases = buildEventReleaseMap(events, repetitions); /** From aa79701c581692976d666d12593ba52768450d33 Mon Sep 17 00:00:00 2001 From: alisawavezen12 Date: Fri, 25 Sep 2026 07:43:00 +0300 Subject: [PATCH 09/10] chore --- .../src/validate-releases.ts | 32 +++++++++++++------ 1 file changed, 22 insertions(+), 10 deletions(-) diff --git a/workers/release-validator/src/validate-releases.ts b/workers/release-validator/src/validate-releases.ts index 476a523a..f0b05454 100644 --- a/workers/release-validator/src/validate-releases.ts +++ b/workers/release-validator/src/validate-releases.ts @@ -4,6 +4,8 @@ import { HOURS_IN_DAY, MINUTES_IN_HOUR, MS_IN_SEC, SECONDS_IN_MINUTE } from '../ 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. @@ -80,8 +82,8 @@ async function validateEventsBatch( eventsCollection: Collection, repetitionsCollection: Collection, releasesToCheck: ReleaseDBScheme[], - allProjectReleases: ReleaseDBScheme[], - releasesByName: Map + allProjectReleases: ReleaseHistoryEntry[], + releasesByName: Map ): Promise { const eventGroupHashes = events.map(event => event.groupHash); const repetitions = await repetitionsCollection.aggregate>([ @@ -226,8 +228,10 @@ async function validateProject(db: Db, projectId: string, releasesToCheck: Relea const repetitionsCollection = db.collection(`repetitions:${projectId}`); /** - * Load the complete release history. Releases outside the validation window - * still define occurrence order and can block an incorrect resolution. + * 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({ @@ -237,23 +241,31 @@ async function validateProject(db: Db, projectId: string, releasesToCheck: Relea $ne: '', }, }) + .project({ + _id: 1, + release: 1, + }) .sort({ _id: 1 }) .toArray(); /** - * Resolve release names to their Mongo records so ObjectIds can be used as - * the chronological source of truth instead of comparing version strings. + * 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(); + const releasesByName = new Map(); for (const release of allProjectReleases) { releasesByName.set(release.release, release); } /** - * Stream eligible events from MongoDB instead of materializing the complete - * project result. The cursor and application batch use the same limit so the - * worker holds at most one bounded portion of event documents in memory. + * 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: { From ecca00deae6d711c2b49412062361e756ec3fdb4 Mon Sep 17 00:00:00 2001 From: alisawavezen12 Date: Fri, 25 Sep 2026 07:53:33 +0300 Subject: [PATCH 10/10] chore --- .../release-validator/src/validate-releases.ts | 17 +++++++++++++---- workers/release-validator/tests/index.test.ts | 15 --------------- 2 files changed, 13 insertions(+), 19 deletions(-) delete mode 100644 workers/release-validator/tests/index.test.ts diff --git a/workers/release-validator/src/validate-releases.ts b/workers/release-validator/src/validate-releases.ts index f0b05454..b6d58289 100644 --- a/workers/release-validator/src/validate-releases.ts +++ b/workers/release-validator/src/validate-releases.ts @@ -126,8 +126,9 @@ async function validateEventsBatch( let lastOccurrenceRelease = originalRelease; /** - * For a regressed event, continue validation from the regression release - * instead of its original release. This enables repeated resolve cycles. + * 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); @@ -177,13 +178,21 @@ async function validateEventsBatch( 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; } /** - * Keep the state transition atomic. A concurrent Grouper update must make - * this conditional update miss instead of overwriting newer state. + * 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 ? { diff --git a/workers/release-validator/tests/index.test.ts b/workers/release-validator/tests/index.test.ts deleted file mode 100644 index 96be5476..00000000 --- a/workers/release-validator/tests/index.test.ts +++ /dev/null @@ -1,15 +0,0 @@ -import '../../../env-test'; -import ReleaseValidatorWorker from '../src'; - -jest.mock('amqplib'); - -/** - * Release Validator worker smoke tests. - */ -describe('ReleaseValidatorWorker', () => { - test('should use the release validator queue', () => { - const worker = new ReleaseValidatorWorker(); - - expect(worker.type).toBe('cron-tasks/release-validator'); - }); -});