Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
15 changes: 15 additions & 0 deletions docker-compose.dev.yml
Original file line number Diff line number Diff line change
Expand Up @@ -62,6 +62,7 @@ services:
#
# System workers:
# - archiver
# - release-validator
# - limiter
# - paymaster
#
Expand All @@ -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"
Expand Down
12 changes: 11 additions & 1 deletion docker-compose.prod.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand All @@ -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
#
Expand Down
4 changes: 3 additions & 1 deletion package.json
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand All @@ -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",
Comment thread
alisawavezen12 marked this conversation as resolved.
"run-email": "yarn worker hawk-worker-email",
"run-telegram": "yarn worker hawk-worker-telegram",
"run-limiter": "yarn worker hawk-worker-limiter",
Expand All @@ -57,7 +59,7 @@
"@babel/parser": "^7.26.9",
"@babel/traverse": "7.26.9",
"@hawk.so/nodejs": "^3.1.1",
"@hawk.so/types": "^0.5.9",
"@hawk.so/types": "^0.9.0",
"@types/amqplib": "^0.8.2",
"@types/jest": "^29.5.14",
"@types/mongodb": "^3.5.15",
Expand Down
68 changes: 68 additions & 0 deletions workers/grouper/src/check-and-mark-regression.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,68 @@
import { Db } from 'mongodb';
import type { ReleaseDBScheme } from '@hawk.so/types';

type ReleaseRecordPart = Pick<ReleaseDBScheme, '_id' | 'projectId' | '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.
*
* @param db - events database connection
* @param projectId - project identifier
* @param groupHash - original event group hash
* @param release - release in which the event occurred again
* @param resolvedInRelease - release in which the event was resolved
* @param regressionInRelease - regression from a previous resolution cycle
*/
export async function checkAndMarkRegression(
db: Db,
projectId: string,
groupHash: string,
release: string,
resolvedInRelease: string,
regressionInRelease?: string
): Promise<void> {
const releases = await db.collection<ReleaseRecordPart>('releases').find({
projectId,
release: {
$in: [resolvedInRelease, release, regressionInRelease].filter(Boolean),
},
})
.toArray();
const resolvedRelease = releases.find(item => item.release === resolvedInRelease);
const repetitionRelease = releases.find(item => item.release === release);
const previousRegressionRelease = regressionInRelease
? releases.find(item => item.release === regressionInRelease)
: undefined;

if (!resolvedRelease || !repetitionRelease) {
return;
}

const resolvedReleaseId = resolvedRelease._id.toHexString();
const isResolvedOrNewerRelease = repetitionRelease._id.toHexString() >= resolvedReleaseId;
/**
* A regression in or after the resolved release belongs to the current
* resolution cycle and must not be overwritten by later repetitions.
*/
const hasRegressionForCurrentCycle = previousRegressionRelease &&
previousRegressionRelease._id.toHexString() >= resolvedReleaseId;

if (!isResolvedOrNewerRelease || hasRegressionForCurrentCycle) {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

documentation needed

return;
}

await db.collection(`events:${projectId}`).updateOne({
groupHash,
resolvedInRelease,
...(regressionInRelease
? { regressionInRelease }
: { regressionInRelease: { $exists: false } }),
}, {
$set: {
regressionInRelease: release,
},
});
}
23 changes: 23 additions & 0 deletions workers/grouper/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@ import GrouperMetrics from './metrics/grouperMetrics';
import GrouperMemoryMonitor from './metrics/memoryMonitor';
import SlowHandleDiagnostics, { SlowHandleSession } from './metrics/slowHandleDiagnostics';
import { grouperDiagnosticsConfig, grouperMemoryConfig } from './metrics/config';
import { checkAndMarkRegression } from './check-and-mark-regression';

/**
* eslint does not count decorators as a variable usage
Expand Down Expand Up @@ -343,10 +344,32 @@ export default class GrouperWorker extends Worker {
timestamp: task.timestamp,
} as RepetitionDBScheme;

if (task.payload.release) {
newRepetition.release = task.payload.release;
}

repetitionId = await session.measureStep('saveRepetition', () => {
return this.saveRepetition(task.projectId, newRepetition);
});

if (task.payload.release && existedEvent.resolvedInRelease) {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

we need to check if task.payload.release is newer than existedEvent.resolvedInRelease

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

we do this later in the markRegression function.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

then we need to rename markRegression to checkForRegression

try {
await checkAndMarkRegression(
this.eventsDb.getConnection(),
task.projectId,
uniqueEventHash,
task.payload.release,
existedEvent.resolvedInRelease,
existedEvent.regressionInRelease
);
} catch (error) {
this.logger.error(
`[checkAndMarkRegression] project=${task.projectId} groupHash=${uniqueEventHash} release=${task.payload.release}`,
error
);
Comment thread
alisawavezen12 marked this conversation as resolved.
}
}

/**
* Clear the large event payload references to allow garbage collection
* This prevents memory leaks from retaining full event objects after delta is computed
Expand Down
151 changes: 151 additions & 0 deletions workers/grouper/tests/index.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 () => {
Expand Down Expand Up @@ -391,6 +392,156 @@ describe('GrouperWorker', () => {
}).toArray()).length).toBe(2);
});

test('Should save repetition release as a separate field', async () => {
await worker.handle(generateTask({ release: 'release-a' }));
await worker.handle(generateTask({ release: 'release-b' }));

const savedRepetition = await repetitionsCollection.findOne({});

expect(savedRepetition.release).toBe('release-b');
});

test('Should not save repetition release when event has no release', async () => {
await worker.handle(generateTask());
await worker.handle(generateTask());

const savedRepetition = await repetitionsCollection.findOne({});

expect(savedRepetition.release).toBeUndefined();
});

describe('Regression marking', () => {
test('Should mark a resolved event as regressed in a newer repetition release', async () => {
await connection.db().collection('releases').insertMany([
{
_id: mongodb.ObjectID.createFromTime(1),
projectId: projectIdMock,
release: 'release-b',
},
{
_id: mongodb.ObjectID.createFromTime(2),
projectId: projectIdMock,
release: 'release-c',
},
]);
await worker.handle(generateTask({ release: 'release-a' }));
await eventsCollection.updateOne({}, {
$set: {
resolvedInRelease: 'release-b',
},
});

await worker.handle(generateTask({ release: 'release-c' }));

expect((await eventsCollection.findOne({})).regressionInRelease).toBe('release-c');
});

test('Should mark as regressed if we later encounter this event with a release that is considered a resolving release', async () => {
await connection.db().collection('releases').insertOne({
_id: mongodb.ObjectID.createFromTime(1),
projectId: projectIdMock,
release: 'release-b',
});
await worker.handle(generateTask({ release: 'release-a' }));
await eventsCollection.updateOne({}, {
$set: {
resolvedInRelease: 'release-b',
},
});

await worker.handle(generateTask({ release: 'release-b' }));

expect((await eventsCollection.findOne({})).regressionInRelease).toBe('release-b');
});

test('Should not mark regression if we encounter an event from one of the old releases', async () => {
await connection.db().collection('releases').insertMany([
{
_id: mongodb.ObjectID.createFromTime(1),
projectId: projectIdMock,
release: 'release-a',
},
{
_id: mongodb.ObjectID.createFromTime(2),
projectId: projectIdMock,
release: 'release-b',
},
]);
await worker.handle(generateTask({ release: 'release-b' }));
await eventsCollection.updateOne({}, {
$set: {
resolvedInRelease: 'release-b',
},
});

await worker.handle(generateTask({ release: 'release-a' }));

expect((await eventsCollection.findOne({})).regressionInRelease).toBeUndefined();
});

test('Should replace an old regression after a newer resolution', async () => {
await connection.db().collection('releases').insertMany([
{
_id: mongodb.ObjectID.createFromTime(1),
projectId: projectIdMock,
release: 'release-b',
},
{
_id: mongodb.ObjectID.createFromTime(2),
projectId: projectIdMock,
release: 'release-c',
},
{
_id: mongodb.ObjectID.createFromTime(3),
projectId: projectIdMock,
release: 'release-d',
},
]);
await worker.handle(generateTask({ release: 'release-a' }));
await eventsCollection.updateOne({}, {
$set: {
resolvedInRelease: 'release-c',
regressionInRelease: 'release-b',
},
});

await worker.handle(generateTask({ release: 'release-d' }));

expect((await eventsCollection.findOne({})).regressionInRelease).toBe('release-d');
});

test('Should not overwrite the first marked regression if we continue receiving this event in newer releases', async () => {
await connection.db().collection('releases').insertMany([
{
_id: mongodb.ObjectID.createFromTime(1),
projectId: projectIdMock,
release: 'release-b',
},
{
_id: mongodb.ObjectID.createFromTime(2),
projectId: projectIdMock,
release: 'release-c',
},
{
_id: mongodb.ObjectID.createFromTime(3),
projectId: projectIdMock,
release: 'release-d',
},
]);
await worker.handle(generateTask({ release: 'release-a' }));
await eventsCollection.updateOne({}, {
$set: {
resolvedInRelease: 'release-b',
regressionInRelease: 'release-c',
},
});

await worker.handle(generateTask({ release: 'release-d' }));

expect((await eventsCollection.findOne({})).regressionInRelease).toBe('release-c');
});
});

test('Should stringify delta', async () => {
const generatedTask = generateTask();

Expand Down
15 changes: 15 additions & 0 deletions workers/release-validator/README.md
Original file line number Diff line number Diff line change
@@ -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
```
10 changes: 10 additions & 0 deletions workers/release-validator/package.json
Original file line number Diff line number Diff line change
@@ -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"
}
Loading
Loading