From 8f58256f24bd120375fd4a0301f96782e2d55ad6 Mon Sep 17 00:00:00 2001 From: alisawavezen12 Date: Thu, 17 Sep 2026 17:53:23 +0300 Subject: [PATCH 01/16] 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/16] 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/16] 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/16] 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/16] 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/16] 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/16] 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/16] 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/16] 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/16] 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'); - }); -}); From ef64cf1a36a499e2b659859f4fe0076b5c58eb33 Mon Sep 17 00:00:00 2001 From: alisawavezen12 Date: Sun, 4 Oct 2026 16:19:06 +0300 Subject: [PATCH 11/16] Align release validation with retention policy Use MAX_DAYS_NUMBER to define the candidate release window and fall back to event timestamps when original release records are archived. --- workers/release-validator/README.md | 2 +- .../src/validate-releases.ts | 40 ++++++----- .../tests/validate-releases.test.ts | 71 +++++++++++++------ 3 files changed, 74 insertions(+), 39 deletions(-) diff --git a/workers/release-validator/README.md b/workers/release-validator/README.md index 333d4fb6..a64e9665 100644 --- a/workers/release-validator/README.md +++ b/workers/release-validator/README.md @@ -4,7 +4,7 @@ Checks releases after a 24-hour observation period and marks original events tha 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. +The worker only uses the fields required for matching releases and event groups. Records without these fields are ignored. Candidate releases must be older than 24 hours and no older than the `MAX_DAYS_NUMBER` retention period, while all available project releases are still used to compare event history. If an event's original release was archived, the original event timestamp is used as its chronological boundary. Queue: `cron-tasks/release-validator` diff --git a/workers/release-validator/src/validate-releases.ts b/workers/release-validator/src/validate-releases.ts index b6d58289..8d5515f7 100644 --- a/workers/release-validator/src/validate-releases.ts +++ b/workers/release-validator/src/validate-releases.ts @@ -12,16 +12,6 @@ type ReleaseHistoryEntry = Pick; */ const RELEASE_OBSERVATION_PERIOD_SECONDS = HOURS_IN_DAY * MINUTES_IN_HOUR * SECONDS_IN_MINUTE; -/** - * Maximum age of a release eligible for validation, in days. - */ -const RELEASE_MAX_AGE_DAYS = 30; - -/** - * Maximum candidate release age expressed in seconds for ObjectId boundaries. - */ -const RELEASE_MAX_AGE_SECONDS = RELEASE_MAX_AGE_DAYS * RELEASE_OBSERVATION_PERIOD_SECONDS; - /** * Maximum number of events processed in one validation batch. * @@ -39,12 +29,18 @@ const EVENTS_BATCH_SIZE = 500; */ async function findReleasesToCheck(db: Db, now: Date): Promise { const nowSeconds = Math.floor(now.getTime() / MS_IN_SEC); + const releaseMaxAgeDays = Number(process.env.MAX_DAYS_NUMBER); + + if (!Number.isFinite(releaseMaxAgeDays) || releaseMaxAgeDays <= 0) { + throw new Error('MAX_DAYS_NUMBER must be a positive number'); + } /** - * Limit candidates to the rollout window: releases must be old enough to - * observe for 24 hours, but recent enough to contain release-aware repetitions. + * Use the archiver retention period as the candidate window so releases are + * validated only while their records are expected to remain available. */ - const oldestReleaseId = ObjectID.createFromTime(nowSeconds - RELEASE_MAX_AGE_SECONDS); + const releaseMaxAgeSeconds = releaseMaxAgeDays * RELEASE_OBSERVATION_PERIOD_SECONDS; + const oldestReleaseId = ObjectID.createFromTime(nowSeconds - releaseMaxAgeSeconds); const newestReleaseId = ObjectID.createFromTime(nowSeconds - RELEASE_OBSERVATION_PERIOD_SECONDS); return db.collection('releases') @@ -156,18 +152,25 @@ async function validateEventsBatch( } /** - * Skip events whose original or latest occurrence release is missing from - * the project release history. + * The archiver can remove the original release while its event is still + * active. In that case, use the original event occurrence time as the + * chronological boundary instead of skipping the event forever. */ - if (!lastOccurrenceRelease) { - continue; + let lastOccurrenceId = lastOccurrenceRelease?._id; + + if (!lastOccurrenceId) { + if (!Number.isFinite(event.timestamp)) { + continue; + } + + lastOccurrenceId = ObjectID.createFromTime(Math.floor(event.timestamp)); } const releasesWithEvent = eventReleases.get(event.groupHash) || new Set(); for (const release of releasesToCheck) { const releaseId = release._id.toHexString(); - const isNewerThanLastOccurrence = releaseId > lastOccurrenceRelease._id.toHexString(); + const isNewerThanLastOccurrence = releaseId > lastOccurrenceId.toHexString(); const occurredInRelease = releasesWithEvent.has(release.release); /** @@ -304,6 +307,7 @@ async function validateProject(db: Db, projectId: string, releasesToCheck: Relea _id: 1, groupHash: 1, 'payload.release': 1, + timestamp: 1, resolvedInRelease: 1, regressionInRelease: 1, }, diff --git a/workers/release-validator/tests/validate-releases.test.ts b/workers/release-validator/tests/validate-releases.test.ts index 332c12a2..a616511d 100644 --- a/workers/release-validator/tests/validate-releases.test.ts +++ b/workers/release-validator/tests/validate-releases.test.ts @@ -191,16 +191,31 @@ describe('validateReleases', () => { expect((await releases.findOne({ release: 'b' })).fixChecked).toBe(false); }); - test('should skip an event when its original release is unknown', async () => { + test('should use the original event timestamp when its release record was archived', async () => { await releases.insertOne(createRelease('b', 48)); - await events.insertOne(createEvent('error-7', 'unknown')); + const event = createEvent('error-7', 'archived-release'); + + event.timestamp = NOW_SECONDS - 72 * HOUR_IN_SECONDS; + await events.insertOne(event); await validateReleases(db, NOW); - expect((await events.findOne({ groupHash: 'error-7' })).resolvedInRelease).toBeUndefined(); + expect((await events.findOne({ groupHash: 'error-7' })).resolvedInRelease).toBe('b'); expect((await releases.findOne({ release: 'b' })).fixChecked).toBe(true); }); + test('should not resolve an event before its timestamp when its release record is missing', async () => { + await releases.insertOne(createRelease('b', 48)); + const event = createEvent('error-after-release', 'missing-release'); + + event.timestamp = NOW_SECONDS - 24 * HOUR_IN_SECONDS; + await events.insertOne(event); + + await validateReleases(db, NOW); + + expect((await events.findOne({ groupHash: 'error-after-release' })).resolvedInRelease).toBeUndefined(); + }); + test('should not overwrite an existing resolved release on repeated validation', async () => { await releases.insertMany([ createRelease('a', 72, true), @@ -328,28 +343,44 @@ describe('validateReleases', () => { expect((await events.findOne({ groupHash: 'error-active-regression' })).resolvedInRelease).toBe('b'); }); - test('should not select a release older than 30 days as a candidate', async () => { - await releases.insertMany([ - createRelease('a', 960, true), - createRelease('b', 744), - ]); - await events.insertOne(createEvent('error-12', 'a')); + test('should not select a release older than the configured retention period as a candidate', async () => { + const originalMaxDaysNumber = process.env.MAX_DAYS_NUMBER; - await validateReleases(db, NOW); + process.env.MAX_DAYS_NUMBER = '10'; - expect((await events.findOne({ groupHash: 'error-12' })).resolvedInRelease).toBeUndefined(); - expect((await releases.findOne({ release: 'b' })).fixChecked).toBe(false); + try { + await releases.insertMany([ + createRelease('a', 480, true), + createRelease('b', 264), + ]); + 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); + } finally { + process.env.MAX_DAYS_NUMBER = originalMaxDaysNumber; + } }); - 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')); + test('should use a release older than the retention period as event history', async () => { + const originalMaxDaysNumber = process.env.MAX_DAYS_NUMBER; - await validateReleases(db, NOW); + process.env.MAX_DAYS_NUMBER = '10'; + + try { + await releases.insertMany([ + createRelease('a', 480, 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'); + expect((await events.findOne({ groupHash: 'error-13' })).resolvedInRelease).toBe('b'); + } finally { + process.env.MAX_DAYS_NUMBER = originalMaxDaysNumber; + } }); }); From 15817a5b823e695c90625bccde957b3caa2995ac Mon Sep 17 00:00:00 2001 From: alisawavezen12 Date: Sun, 4 Oct 2026 16:23:22 +0300 Subject: [PATCH 12/16] Isolate release validation failures by project --- .../src/validate-releases.ts | 9 ++- .../tests/validate-releases.test.ts | 58 +++++++++++++++++++ 2 files changed, 66 insertions(+), 1 deletion(-) diff --git a/workers/release-validator/src/validate-releases.ts b/workers/release-validator/src/validate-releases.ts index 8d5515f7..6c95c10a 100644 --- a/workers/release-validator/src/validate-releases.ts +++ b/workers/release-validator/src/validate-releases.ts @@ -1,11 +1,14 @@ 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 createLogger from '../../../lib/logger'; import { buildEventReleaseMap } from './utils/build-event-release-map'; import { groupReleasesByProject } from './utils/group-releases-by-project'; type ReleaseHistoryEntry = Pick; +const logger = createLogger(); + /** * Time allowed for repetitions to arrive before a release is checked for * resolved events. @@ -381,6 +384,10 @@ export async function validateReleases(db: Db, now = new Date()): Promise const releasesByProject = groupReleasesByProject(releasesToCheck); for (const [projectId, projectReleases] of releasesByProject) { - await validateProject(db, projectId, projectReleases); + try { + await validateProject(db, projectId, projectReleases); + } catch (error) { + logger.error(`Failed to validate releases for project ${projectId}`, error); + } } } diff --git a/workers/release-validator/tests/validate-releases.test.ts b/workers/release-validator/tests/validate-releases.test.ts index a616511d..bac09105 100644 --- a/workers/release-validator/tests/validate-releases.test.ts +++ b/workers/release-validator/tests/validate-releases.test.ts @@ -343,6 +343,64 @@ describe('validateReleases', () => { expect((await events.findOne({ groupHash: 'error-active-regression' })).resolvedInRelease).toBe('b'); }); + test('should continue validating other projects when one project fails', async () => { + const failedProjectId = 'failed-release-validator-project'; + + await releases.insertMany([ + createRelease('a', 72, true), + { + ...createRelease('b', 49), + projectId: failedProjectId, + }, + createRelease('c', 48), + ]); + await db.createCollection(`events:${failedProjectId}`, { + validator: { + $jsonSchema: { + properties: { + resolvedInRelease: { + bsonType: 'int', + }, + }, + }, + }, + validationLevel: 'strict', + validationAction: 'error', + }); + await db.collection(`events:${failedProjectId}`).insertOne({ + _id: new ObjectID(), + groupHash: 'failed-project-event', + payload: { + title: 'failed-project-event', + release: 'a', + }, + totalCount: 1, + catcherType: 'errors/default', + usersAffected: 0, + visitedBy: [], + timestamp: NOW_SECONDS - 96 * HOUR_IN_SECONDS, + }); + await events.insertOne(createEvent('successful-project-event', 'a')); + + try { + await expect(validateReleases(db, NOW)).resolves.toBeUndefined(); + + expect((await releases.findOne({ + projectId: failedProjectId, + release: 'b', + })).fixChecked).toBeUndefined(); + expect((await events.findOne({ groupHash: 'successful-project-event' })).resolvedInRelease).toBe('c'); + expect((await releases.findOne({ + projectId: PROJECT_ID, + release: 'c', + })).fixChecked).toBe(true); + } finally { + await db.dropCollection(`events:${failedProjectId}`).catch(() => undefined); + await db.dropCollection(`repetitions:${failedProjectId}`).catch(() => undefined); + await releases.deleteMany({ projectId: failedProjectId }); + } + }); + test('should not select a release older than the configured retention period as a candidate', async () => { const originalMaxDaysNumber = process.env.MAX_DAYS_NUMBER; From e2760e88905aa29ba0e9854facda3522420f71de Mon Sep 17 00:00:00 2001 From: alisawavezen12 Date: Sun, 4 Oct 2026 17:05:08 +0300 Subject: [PATCH 13/16] Clear regression when resolving event again --- .../grouper/src/check-and-mark-regression.ts | 21 +++---------- workers/grouper/src/index.ts | 5 ++- workers/grouper/tests/index.test.ts | 31 ------------------- .../src/validate-releases.ts | 3 ++ .../tests/validate-releases.test.ts | 2 +- 5 files changed, 10 insertions(+), 52 deletions(-) diff --git a/workers/grouper/src/check-and-mark-regression.ts b/workers/grouper/src/check-and-mark-regression.ts index e5084a49..53c59689 100644 --- a/workers/grouper/src/check-and-mark-regression.ts +++ b/workers/grouper/src/check-and-mark-regression.ts @@ -14,28 +14,23 @@ type ReleaseRecordPart = Pick; * @param groupHash - original event group hash * @param release - release in which the event occurred again * @param resolvedInRelease - release in which the event was resolved - * @param regressionInRelease - regression from a previous resolution cycle */ export async function checkAndMarkRegression( db: Db, projectId: string, groupHash: string, release: string, - resolvedInRelease: string, - regressionInRelease?: string + resolvedInRelease: string ): Promise { const releases = await db.collection('releases').find({ projectId, release: { - $in: [resolvedInRelease, release, regressionInRelease].filter(Boolean), + $in: [resolvedInRelease, release], }, }) .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; @@ -43,23 +38,15 @@ export async function checkAndMarkRegression( const resolvedReleaseId = resolvedRelease._id.toHexString(); const isResolvedOrNewerRelease = repetitionRelease._id.toHexString() >= resolvedReleaseId; - /** - * A regression in or after the resolved release belongs to the current - * resolution cycle and must not be overwritten by later repetitions. - */ - const hasRegressionForCurrentCycle = previousRegressionRelease && - previousRegressionRelease._id.toHexString() >= resolvedReleaseId; - if (!isResolvedOrNewerRelease || hasRegressionForCurrentCycle) { + if (!isResolvedOrNewerRelease) { return; } await db.collection(`events:${projectId}`).updateOne({ groupHash, resolvedInRelease, - ...(regressionInRelease - ? { regressionInRelease } - : { regressionInRelease: { $exists: false } }), + regressionInRelease: { $exists: false }, }, { $set: { regressionInRelease: release, diff --git a/workers/grouper/src/index.ts b/workers/grouper/src/index.ts index e6839c39..60e9cb25 100644 --- a/workers/grouper/src/index.ts +++ b/workers/grouper/src/index.ts @@ -352,15 +352,14 @@ export default class GrouperWorker extends Worker { return this.saveRepetition(task.projectId, newRepetition); }); - if (task.payload.release && existedEvent.resolvedInRelease) { + if (task.payload.release && existedEvent.resolvedInRelease && !existedEvent.regressionInRelease) { try { await checkAndMarkRegression( this.eventsDb.getConnection(), task.projectId, uniqueEventHash, task.payload.release, - existedEvent.resolvedInRelease, - existedEvent.regressionInRelease + existedEvent.resolvedInRelease ); } catch (error) { this.logger.error( diff --git a/workers/grouper/tests/index.test.ts b/workers/grouper/tests/index.test.ts index 6e2421c3..c71d03ea 100644 --- a/workers/grouper/tests/index.test.ts +++ b/workers/grouper/tests/index.test.ts @@ -479,37 +479,6 @@ describe('GrouperWorker', () => { expect((await eventsCollection.findOne({})).regressionInRelease).toBeUndefined(); }); - test('Should replace an old regression after a newer resolution', async () => { - await connection.db().collection('releases').insertMany([ - { - _id: mongodb.ObjectID.createFromTime(1), - projectId: projectIdMock, - release: 'release-b', - }, - { - _id: mongodb.ObjectID.createFromTime(2), - projectId: projectIdMock, - release: 'release-c', - }, - { - _id: mongodb.ObjectID.createFromTime(3), - projectId: projectIdMock, - release: 'release-d', - }, - ]); - await worker.handle(generateTask({ release: 'release-a' })); - await eventsCollection.updateOne({}, { - $set: { - resolvedInRelease: 'release-c', - regressionInRelease: 'release-b', - }, - }); - - await worker.handle(generateTask({ release: 'release-d' })); - - expect((await eventsCollection.findOne({})).regressionInRelease).toBe('release-d'); - }); - test('Should not overwrite the first marked regression if we continue receiving this event in newer releases', async () => { await connection.db().collection('releases').insertMany([ { diff --git a/workers/release-validator/src/validate-releases.ts b/workers/release-validator/src/validate-releases.ts index 6c95c10a..9faf338e 100644 --- a/workers/release-validator/src/validate-releases.ts +++ b/workers/release-validator/src/validate-releases.ts @@ -219,6 +219,9 @@ async function validateEventsBatch( $set: { resolvedInRelease: release.release, }, + ...(event.regressionInRelease + ? { $unset: { regressionInRelease: '' } } + : {}), }); /** diff --git a/workers/release-validator/tests/validate-releases.test.ts b/workers/release-validator/tests/validate-releases.test.ts index bac09105..b04b8a3a 100644 --- a/workers/release-validator/tests/validate-releases.test.ts +++ b/workers/release-validator/tests/validate-releases.test.ts @@ -323,7 +323,7 @@ describe('validateReleases', () => { const updatedEvent = await events.findOne({ groupHash: 'error-cycle' }); expect(updatedEvent.resolvedInRelease).toBe('d'); - expect(updatedEvent.regressionInRelease).toBe('c'); + expect(updatedEvent.regressionInRelease).toBeUndefined(); }); test('should not resolve an event again before a release newer than its regression', async () => { From a82acd316dce5380a5a9746e7846aa2f02be1427 Mon Sep 17 00:00:00 2001 From: alisawavezen12 Date: Sun, 4 Oct 2026 17:11:24 +0300 Subject: [PATCH 14/16] Track latest event occurrence before validation --- .../src/validate-releases.ts | 33 +++++++++---------- 1 file changed, 16 insertions(+), 17 deletions(-) diff --git a/workers/release-validator/src/validate-releases.ts b/workers/release-validator/src/validate-releases.ts index 9faf338e..3e860de7 100644 --- a/workers/release-validator/src/validate-releases.ts +++ b/workers/release-validator/src/validate-releases.ts @@ -171,25 +171,24 @@ async function validateEventsBatch( const releasesWithEvent = eventReleases.get(event.groupHash) || new Set(); - for (const release of releasesToCheck) { - const releaseId = release._id.toHexString(); - const isNewerThanLastOccurrence = releaseId > lastOccurrenceId.toHexString(); - const occurredInRelease = releasesWithEvent.has(release.release); + /** + * Find the latest known release in which the event occurred once per event. + * This includes releases younger than 24 hours, which must still prevent an + * older candidate from resolving the event. + */ + for (const projectRelease of allProjectReleases) { + if ( + projectRelease._id.toHexString() > lastOccurrenceId.toHexString() && + releasesWithEvent.has(projectRelease.release) + ) { + lastOccurrenceId = projectRelease._id; + } + } - /** - * 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); - }); + for (const release of releasesToCheck) { + const isNewerThanLastOccurrence = release._id.toHexString() > lastOccurrenceId.toHexString(); - /** - * 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) { + if (!isNewerThanLastOccurrence) { continue; } From 59ef3e4986697337e45dfa02bceeef723bb9534f Mon Sep 17 00:00:00 2001 From: alisawavezen12 Date: Sun, 4 Oct 2026 17:33:19 +0300 Subject: [PATCH 15/16] tests fix --- .gitignore | 2 ++ package.json | 2 +- workers/release-validator/jest.config.js | 18 ++++++++++++++++++ workers/release-validator/jest.setup.js | 1 + .../tests/validate-releases.test.ts | 2 +- 5 files changed, 23 insertions(+), 2 deletions(-) create mode 100644 workers/release-validator/jest.config.js create mode 100644 workers/release-validator/jest.setup.js diff --git a/.gitignore b/.gitignore index b4b45798..2507214f 100644 --- a/.gitignore +++ b/.gitignore @@ -13,6 +13,8 @@ globalConfig.json !jest.setup.mongo-repl-set.js !jest.setup.redis-mock.js !jest-mongodb-config.js +!workers/release-validator/jest.config.js +!workers/release-validator/jest.setup.js !migrate-mongo-config.js !/env.js !convertors/**/*.js diff --git a/package.json b/package.json index 32953a73..b5bb807e 100644 --- a/package.json +++ b/package.json @@ -25,7 +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:release-validator": "jest --config workers/release-validator/jest.config.js", "test:slack": "jest workers/slack", "test:loop": "jest workers/loop", "test:limiter": "jest workers/limiter --runInBand", diff --git a/workers/release-validator/jest.config.js b/workers/release-validator/jest.config.js new file mode 100644 index 00000000..2b8182ed --- /dev/null +++ b/workers/release-validator/jest.config.js @@ -0,0 +1,18 @@ +const baseConfig = require('../../jest.config'); + +module.exports = { + ...baseConfig, + rootDir: '../..', + setupFiles: [ + '/jest.setup.js', + '/workers/release-validator/jest.setup.js', + ], + setupFilesAfterEnv: [ '/jest.setup.mongo-repl-set.js' ], + globalTeardown: '/jest.global-teardown.js', + roots: [ + '/workers/release-validator', + '/lib', + ], + testMatch: [ '/workers/release-validator/**/*.test.ts' ], + moduleFileExtensions: ['ts', 'tsx', 'js', 'json', 'node'], +}; diff --git a/workers/release-validator/jest.setup.js b/workers/release-validator/jest.setup.js new file mode 100644 index 00000000..d0c9bf1e --- /dev/null +++ b/workers/release-validator/jest.setup.js @@ -0,0 +1 @@ +process.env.MAX_DAYS_NUMBER = process.env.MAX_DAYS_NUMBER || '30'; diff --git a/workers/release-validator/tests/validate-releases.test.ts b/workers/release-validator/tests/validate-releases.test.ts index b04b8a3a..adeffe92 100644 --- a/workers/release-validator/tests/validate-releases.test.ts +++ b/workers/release-validator/tests/validate-releases.test.ts @@ -388,7 +388,7 @@ describe('validateReleases', () => { expect((await releases.findOne({ projectId: failedProjectId, release: 'b', - })).fixChecked).toBeUndefined(); + })).fixChecked).toBe(false); expect((await events.findOne({ groupHash: 'successful-project-event' })).resolvedInRelease).toBe('c'); expect((await releases.findOne({ projectId: PROJECT_ID, From 8b4e8d071e79969b5b7040e6bfcf1b96b50331d4 Mon Sep 17 00:00:00 2001 From: alisawavezen12 Date: Sun, 4 Oct 2026 17:42:08 +0300 Subject: [PATCH 16/16] chore --- workers/release-validator/tests/validate-releases.test.ts | 2 ++ 1 file changed, 2 insertions(+) diff --git a/workers/release-validator/tests/validate-releases.test.ts b/workers/release-validator/tests/validate-releases.test.ts index adeffe92..57551d72 100644 --- a/workers/release-validator/tests/validate-releases.test.ts +++ b/workers/release-validator/tests/validate-releases.test.ts @@ -5,6 +5,7 @@ import { validateReleases } from '../src/validate-releases'; const PROJECT_ID = 'release-validator-project'; const HOUR_IN_SECONDS = 60 * 60; +const DEFAULT_MAX_DAYS_NUMBER = '30'; const NOW_SECONDS = Math.floor(new Date('2026-09-17T12:00:00.000Z').getTime() / 1000); const NOW = new Date(NOW_SECONDS * 1000); @@ -96,6 +97,7 @@ describe('validateReleases', () => { }); beforeEach(async () => { + process.env.MAX_DAYS_NUMBER = DEFAULT_MAX_DAYS_NUMBER; await releases.deleteMany({ projectId: PROJECT_ID }); await events.deleteMany({}); await repetitions.deleteMany({});