From 9d6ede7bd1179bac828a15f7323c692e2230100e Mon Sep 17 00:00:00 2001 From: alisawavezen12 Date: Thu, 10 Sep 2026 15:50:32 +0300 Subject: [PATCH 1/6] feat(release-worker): assign sequence to new releases --- package.json | 4 +- workers/release/src/index.ts | 78 ++++++++++++++++++++++++++--- workers/release/tests/index.test.ts | 77 ++++++++++++++++++++++++++-- 3 files changed, 148 insertions(+), 11 deletions(-) diff --git a/package.json b/package.json index f2e1e154d..e3f97b82b 100644 --- a/package.json +++ b/package.json @@ -1,7 +1,7 @@ { "name": "hawk.workers", "private": true, - "version": "0.1.4", + "version": "0.1.5", "description": "Hawk workers", "repository": "git@github.com:codex-team/hawk.workers.git", "license": "BUSL-1.1", @@ -57,7 +57,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.7.1", "@types/amqplib": "^0.8.2", "@types/jest": "^29.5.14", "@types/mongodb": "^3.5.15", diff --git a/workers/release/src/index.ts b/workers/release/src/index.ts index ea9282dca..e43e905a4 100644 --- a/workers/release/src/index.ts +++ b/workers/release/src/index.ts @@ -9,7 +9,7 @@ import { Worker } from '../../../lib/worker'; import { DatabaseReadWriteError, NonCriticalError } from '../../../lib/workerErrors'; import * as pkg from '../package.json'; import { ReleaseWorkerTask, ReleaseWorkerAddReleasePayload, CommitDataUnparsed } from '../types'; -import { Collection, MongoClient, MongoError } from 'mongodb'; +import { Collection, FindAndModifyWriteOpResultObject, MongoClient, MongoError, UpdateQuery } from 'mongodb'; import { SourceMapDataExtended, SourceMapFileChunk, CommitData, SourcemapCollectedData, ReleaseDBScheme } from '@hawk.so/types'; /** @@ -18,6 +18,14 @@ import { SourceMapDataExtended, SourceMapFileChunk, CommitData, SourcemapCollect /* eslint-disable @typescript-eslint/no-magic-numbers */ const DB_DUPLICATE_KEY_ERROR = '11000'; +/** + * Counter used to assign the first-seen sequence to a project release. + */ +interface ReleaseSequenceCounter { + _id: string; + nextSequence: number; +} + /** * Worker to save releases */ @@ -43,6 +51,11 @@ export default class ReleaseWorker extends Worker { */ private releasesCollection: Collection; + /** + * Collection with per-project release sequence counters. + */ + private releaseSequencesCollection: Collection; + /** * Mongo client for events database, used for transactions */ @@ -56,6 +69,7 @@ export default class ReleaseWorker extends Worker { await this.db.connect(); this.db.createGridFsBucket(this.dbCollectionName); this.releasesCollection = this.db.getConnection().collection(this.dbCollectionName); + this.releaseSequencesCollection = this.db.getConnection().collection('releaseSequenceCounters'); await super.start(); } @@ -89,11 +103,19 @@ export default class ReleaseWorker extends Worker { this.logger.info(`saveRelease: save release for project: ${projectId}, release: ${payload.release}`); try { const commits = payload.commits; + const validCommits = !!commits && this.areCommitsValid(commits); + const hasFiles = Array.isArray(payload.files) && payload.files.length > 0; + + if (!validCommits && !hasFiles) { + return; + } + + await this.ensureRelease(projectId, payload.release); /** * Save commits */ - if (commits && this.areCommitsValid(commits)) { + if (validCommits) { const commitsWithParsedDate: CommitData[] = commits.map(commit => ({ ...commit, date: new Date(commit.date), @@ -106,13 +128,11 @@ export default class ReleaseWorker extends Worker { $set: { commits: commitsWithParsedDate, }, - }, { - upsert: true, }); } // save source maps - if (payload.files) { + if (hasFiles) { await this.saveSourceMap(projectId, payload); } } catch (err) { @@ -122,6 +142,52 @@ export default class ReleaseWorker extends Worker { } } + /** + * Ensure that a release gets one stable first-seen sequence. + * + * Sequence gaps are acceptable when concurrent requests race: only the + * request that wins the unique release insert owns the assigned sequence. + * + * @param projectId - project id to bind the corresponding release. + * @param release - release name + */ + private async ensureRelease(projectId: string, release: string): Promise { + const existingRelease = await this.releasesCollection.findOne({ + projectId, + release, + }); + + if (existingRelease) { + return; + } + + const counter = await this.releaseSequencesCollection.findOneAndUpdate({ + _id: projectId, + }, { + $inc: { + nextSequence: 1, + }, + } as unknown as UpdateQuery, { + upsert: true, + returnOriginal: false, + }) as FindAndModifyWriteOpResultObject; + + const releaseSequence = counter.value.nextSequence; + + try { + await this.releasesCollection.insertOne({ + projectId, + release, + releaseSequence, + commits: [], + } as unknown as ReleaseDBScheme); + } catch (error) { + if (error.code?.toString() !== DB_DUPLICATE_KEY_ERROR) { + throw error; + } + } + } + /** * Check commtis for the content of all required data * @@ -220,7 +286,7 @@ export default class ReleaseWorker extends Worker { projectId: projectId, release: payload.release, files: savedFilesWithoutContent, - } as ReleaseDBScheme); + } as unknown as ReleaseDBScheme); this.logger.info('inserted new release'); } catch (err) { if ((err as MongoError).code.toString() === DB_DUPLICATE_KEY_ERROR) { diff --git a/workers/release/tests/index.test.ts b/workers/release/tests/index.test.ts index 24f5f46ef..bab3c628a 100644 --- a/workers/release/tests/index.test.ts +++ b/workers/release/tests/index.test.ts @@ -69,6 +69,10 @@ describe('Release Worker', () => { }); db = connection.db(); collection = await db.collection('releases'); + await collection.createIndex({ projectId: 1, release: 1 }, { + name: 'projectId_release_unique_idx', + unique: true, + }); await mockBundle.build(); }); @@ -78,6 +82,7 @@ describe('Release Worker', () => { */ afterAll(async () => { await collection.deleteMany({}); + await db.collection('releaseSequenceCounters').deleteMany({}); await db.collection('releases.chunks').deleteMany({}); await db.collection('releases.files').deleteMany({}); @@ -88,6 +93,7 @@ describe('Release Worker', () => { beforeEach(async () => { await collection.deleteMany({}); + await db.collection('releaseSequenceCounters').deleteMany({}); await db.collection('releases.chunks').deleteMany({}); await db.collection('releases.files').deleteMany({}); }); @@ -184,13 +190,18 @@ describe('Release Worker', () => { await expect(release).toMatchObject(parsedReleasePayload); }); - test('should update a release if it is already exists', async () => { + test('should keep the same release sequence on repeated uploads', async () => { await worker.handle({ projectId, type: 'add-release', payload: releasePayload, }); + const firstRelease = await collection.findOne({ + projectId, + release: releasePayload.release, + }); + await worker.handle({ projectId, type: 'add-release', @@ -200,9 +211,69 @@ describe('Release Worker', () => { }, }); - const count = await collection.countDocuments(); + const secondRelease = await collection.findOne({ + projectId, + release: releasePayload.release, + }); + + await expect(secondRelease.releaseSequence).toEqual(firstRelease.releaseSequence); + await expect(await collection.countDocuments()).toEqual(1); + }); + + test('should assign different sequences to concurrently created releases', async () => { + await Promise.all([ + worker.handle({ + projectId, + type: 'add-release', + payload: releasePayload, + }), + worker.handle({ + projectId, + type: 'add-release', + payload: { + ...releasePayload, + release: 'Dapper Dragon 2', + }, + }), + ]); + + const releases = await collection.find({ projectId }) + .sort({ releaseSequence: 1 }) + .toArray(); + + await expect(releases).toHaveLength(2); + await expect(releases.map(release => release.releaseSequence)).toEqual([1, 2]); + }); + + test('should use one sequence when commits and source maps create the same release concurrently', async () => { + const map = await mockBundle.getSourceMap(); + + await Promise.all([ + worker.handle({ + projectId, + type: 'add-release', + payload: releasePayload, + }), + worker.handle({ + projectId, + type: 'add-release', + payload: { + ...releasePayload, + files: [ { + name: 'main.js.map', + payload: map, + } ], + }, + }), + ]); + + const releases = await collection.find({ + projectId, + release: releasePayload.release, + }).toArray(); - await expect(count).toEqual(1); + await expect(releases).toHaveLength(1); + await expect(releases[0].releaseSequence).toEqual(1); }); test('should correctly handle release with multiple source maps in a single transaction', async () => { From 0b35f82ed2d35dc8d76cdd0da6c8367f0af853a2 Mon Sep 17 00:00:00 2001 From: alisawavezen12 Date: Thu, 10 Sep 2026 16:32:34 +0300 Subject: [PATCH 2/6] Simplify Release Sequence Assignment --- workers/release/src/index.ts | 81 ++++++----------------------- workers/release/tests/index.test.ts | 64 ++++++----------------- 2 files changed, 30 insertions(+), 115 deletions(-) diff --git a/workers/release/src/index.ts b/workers/release/src/index.ts index e43e905a4..8d6483211 100644 --- a/workers/release/src/index.ts +++ b/workers/release/src/index.ts @@ -9,23 +9,9 @@ import { Worker } from '../../../lib/worker'; import { DatabaseReadWriteError, NonCriticalError } from '../../../lib/workerErrors'; import * as pkg from '../package.json'; import { ReleaseWorkerTask, ReleaseWorkerAddReleasePayload, CommitDataUnparsed } from '../types'; -import { Collection, FindAndModifyWriteOpResultObject, MongoClient, MongoError, UpdateQuery } from 'mongodb'; +import { Collection, MongoClient } from 'mongodb'; import { SourceMapDataExtended, SourceMapFileChunk, CommitData, SourcemapCollectedData, ReleaseDBScheme } from '@hawk.so/types'; -/** - * Error code of MongoDB key duplication error - */ -/* eslint-disable @typescript-eslint/no-magic-numbers */ -const DB_DUPLICATE_KEY_ERROR = '11000'; - -/** - * Counter used to assign the first-seen sequence to a project release. - */ -interface ReleaseSequenceCounter { - _id: string; - nextSequence: number; -} - /** * Worker to save releases */ @@ -51,11 +37,6 @@ export default class ReleaseWorker extends Worker { */ private releasesCollection: Collection; - /** - * Collection with per-project release sequence counters. - */ - private releaseSequencesCollection: Collection; - /** * Mongo client for events database, used for transactions */ @@ -69,7 +50,6 @@ export default class ReleaseWorker extends Worker { await this.db.connect(); this.db.createGridFsBucket(this.dbCollectionName); this.releasesCollection = this.db.getConnection().collection(this.dbCollectionName); - this.releaseSequencesCollection = this.db.getConnection().collection('releaseSequenceCounters'); await super.start(); } @@ -145,8 +125,7 @@ export default class ReleaseWorker extends Worker { /** * Ensure that a release gets one stable first-seen sequence. * - * Sequence gaps are acceptable when concurrent requests race: only the - * request that wins the unique release insert owns the assigned sequence. + * Release processing is intentionally sequential in the current deployment. * * @param projectId - project id to bind the corresponding release. * @param release - release name @@ -161,31 +140,22 @@ export default class ReleaseWorker extends Worker { return; } - const counter = await this.releaseSequencesCollection.findOneAndUpdate({ - _id: projectId, + const lastRelease = await this.releasesCollection.findOne({ + projectId, + releaseSequence: { $exists: true }, }, { - $inc: { - nextSequence: 1, - }, - } as unknown as UpdateQuery, { - upsert: true, - returnOriginal: false, - }) as FindAndModifyWriteOpResultObject; + sort: { releaseSequence: -1 }, + projection: { releaseSequence: 1 }, + }); - const releaseSequence = counter.value.nextSequence; + const releaseSequence = (lastRelease?.releaseSequence || 0) + 1; - try { - await this.releasesCollection.insertOne({ - projectId, - release, - releaseSequence, - commits: [], - } as unknown as ReleaseDBScheme); - } catch (error) { - if (error.code?.toString() !== DB_DUPLICATE_KEY_ERROR) { - throw error; - } - } + await this.releasesCollection.insertOne({ + projectId, + release, + releaseSequence, + commits: [], + } as unknown as ReleaseDBScheme); } /** @@ -278,27 +248,6 @@ export default class ReleaseWorker extends Worker { * or * - update previous record with adding new saved maps */ - if (!existedRelease) { - this.logger.info('trying insert new release'); - - try { - await this.releasesCollection.insertOne({ - projectId: projectId, - release: payload.release, - files: savedFilesWithoutContent, - } as unknown as ReleaseDBScheme); - this.logger.info('inserted new release'); - } catch (err) { - if ((err as MongoError).code.toString() === DB_DUPLICATE_KEY_ERROR) { - this.logger.warn(`Duplicate key on insert, retrying update after small delay`); - /* eslint-disable @typescript-eslint/no-magic-numbers */ - await new Promise(resolve => setTimeout(resolve, 200)); - } else { - throw err; - } - } - } - await this.releasesCollection.findOneAndUpdate({ projectId: projectId, release: payload.release, diff --git a/workers/release/tests/index.test.ts b/workers/release/tests/index.test.ts index bab3c628a..e4474ac4c 100644 --- a/workers/release/tests/index.test.ts +++ b/workers/release/tests/index.test.ts @@ -82,7 +82,6 @@ describe('Release Worker', () => { */ afterAll(async () => { await collection.deleteMany({}); - await db.collection('releaseSequenceCounters').deleteMany({}); await db.collection('releases.chunks').deleteMany({}); await db.collection('releases.files').deleteMany({}); @@ -93,7 +92,6 @@ describe('Release Worker', () => { beforeEach(async () => { await collection.deleteMany({}); - await db.collection('releaseSequenceCounters').deleteMany({}); await db.collection('releases.chunks').deleteMany({}); await db.collection('releases.files').deleteMany({}); }); @@ -220,22 +218,21 @@ describe('Release Worker', () => { await expect(await collection.countDocuments()).toEqual(1); }); - test('should assign different sequences to concurrently created releases', async () => { - await Promise.all([ - worker.handle({ - projectId, - type: 'add-release', - payload: releasePayload, - }), - worker.handle({ - projectId, - type: 'add-release', - payload: { - ...releasePayload, - release: 'Dapper Dragon 2', - }, - }), - ]); + test('should assign increasing sequences to new releases', async () => { + await worker.handle({ + projectId, + type: 'add-release', + payload: releasePayload, + }); + + await worker.handle({ + projectId, + type: 'add-release', + payload: { + ...releasePayload, + release: 'Dapper Dragon 2', + }, + }); const releases = await collection.find({ projectId }) .sort({ releaseSequence: 1 }) @@ -245,37 +242,6 @@ describe('Release Worker', () => { await expect(releases.map(release => release.releaseSequence)).toEqual([1, 2]); }); - test('should use one sequence when commits and source maps create the same release concurrently', async () => { - const map = await mockBundle.getSourceMap(); - - await Promise.all([ - worker.handle({ - projectId, - type: 'add-release', - payload: releasePayload, - }), - worker.handle({ - projectId, - type: 'add-release', - payload: { - ...releasePayload, - files: [ { - name: 'main.js.map', - payload: map, - } ], - }, - }), - ]); - - const releases = await collection.find({ - projectId, - release: releasePayload.release, - }).toArray(); - - await expect(releases).toHaveLength(1); - await expect(releases[0].releaseSequence).toEqual(1); - }); - test('should correctly handle release with multiple source maps in a single transaction', async () => { const map = await mockBundle.getSourceMap(); From cb1538def98ba134e59912920070a4abe8096b85 Mon Sep 17 00:00:00 2001 From: alisawavezen12 Date: Thu, 10 Sep 2026 16:38:30 +0300 Subject: [PATCH 3/6] lint --- workers/release/tests/index.test.ts | 14 ++++++++++---- 1 file changed, 10 insertions(+), 4 deletions(-) diff --git a/workers/release/tests/index.test.ts b/workers/release/tests/index.test.ts index e4474ac4c..b464d06f7 100644 --- a/workers/release/tests/index.test.ts +++ b/workers/release/tests/index.test.ts @@ -69,10 +69,16 @@ describe('Release Worker', () => { }); db = connection.db(); collection = await db.collection('releases'); - await collection.createIndex({ projectId: 1, release: 1 }, { - name: 'projectId_release_unique_idx', - unique: true, - }); + await collection.createIndex( + { + projectId: 1, + release: 1, + }, + { + name: 'projectId_release_unique_idx', + unique: true, + } + ); await mockBundle.build(); }); From 7fe5034ed5c92f155134eff3bec3471c582a33f6 Mon Sep 17 00:00:00 2001 From: alisawavezen12 Date: Sun, 13 Sep 2026 13:33:12 +0300 Subject: [PATCH 4/6] Create releases only when data is available Add debug logging for skipped releases and rename the release creation helper. --- workers/release/src/index.ts | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/workers/release/src/index.ts b/workers/release/src/index.ts index 8d6483211..dc07942d9 100644 --- a/workers/release/src/index.ts +++ b/workers/release/src/index.ts @@ -87,10 +87,12 @@ export default class ReleaseWorker extends Worker { const hasFiles = Array.isArray(payload.files) && payload.files.length > 0; if (!validCommits && !hasFiles) { + this.logger.debug(`Skipping release ${payload.release} for project ${projectId}: no valid commits or source maps`); + return; } - await this.ensureRelease(projectId, payload.release); + await this.createReleaseIfMissing(projectId, payload.release); /** * Save commits @@ -130,7 +132,7 @@ export default class ReleaseWorker extends Worker { * @param projectId - project id to bind the corresponding release. * @param release - release name */ - private async ensureRelease(projectId: string, release: string): Promise { + private async createReleaseIfMissing(projectId: string, release: string): Promise { const existingRelease = await this.releasesCollection.findOne({ projectId, release, From d72ac6943377fc1310dbbd72faea2f4d67012d81 Mon Sep 17 00:00:00 2001 From: alisawavezen12 Date: Sun, 13 Sep 2026 13:35:08 +0300 Subject: [PATCH 5/6] Clarify Release Sequence Semantics --- workers/release/src/index.ts | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/workers/release/src/index.ts b/workers/release/src/index.ts index dc07942d9..3b6630efc 100644 --- a/workers/release/src/index.ts +++ b/workers/release/src/index.ts @@ -125,9 +125,11 @@ export default class ReleaseWorker extends Worker { } /** - * Ensure that a release gets one stable first-seen sequence. + * Create a release once and assign its first-registration sequence. * - * Release processing is intentionally sequential in the current deployment. + * The sequence starts at 1 for each project and is based on the order in + * which new releases are received by this sequential worker. Repeated + * commits or source-map uploads keep the existing sequence. * * @param projectId - project id to bind the corresponding release. * @param release - release name @@ -246,9 +248,7 @@ export default class ReleaseWorker extends Worker { try { /** - * - insert new record with saved maps - * or - * - update previous record with adding new saved maps + * Add new source maps to the release created by saveRelease. */ await this.releasesCollection.findOneAndUpdate({ projectId: projectId, From 17dccbe77a581e481627505ae087a1096163e0c5 Mon Sep 17 00:00:00 2001 From: alisawavezen12 Date: Sun, 13 Sep 2026 13:41:45 +0300 Subject: [PATCH 6/6] Handle duplicate release creation gracefully --- workers/release/src/index.ts | 30 +++++++++++++++++++++++------- 1 file changed, 23 insertions(+), 7 deletions(-) diff --git a/workers/release/src/index.ts b/workers/release/src/index.ts index 3b6630efc..add96185b 100644 --- a/workers/release/src/index.ts +++ b/workers/release/src/index.ts @@ -9,9 +9,15 @@ import { Worker } from '../../../lib/worker'; import { DatabaseReadWriteError, NonCriticalError } from '../../../lib/workerErrors'; import * as pkg from '../package.json'; import { ReleaseWorkerTask, ReleaseWorkerAddReleasePayload, CommitDataUnparsed } from '../types'; -import { Collection, MongoClient } from 'mongodb'; +import { Collection, MongoClient, MongoError } from 'mongodb'; import { SourceMapDataExtended, SourceMapFileChunk, CommitData, SourcemapCollectedData, ReleaseDBScheme } from '@hawk.so/types'; +/** + * Error code of MongoDB key duplication error + */ +/* eslint-disable @typescript-eslint/no-magic-numbers */ +const DB_DUPLICATE_KEY_ERROR = '11000'; + /** * Worker to save releases */ @@ -154,12 +160,22 @@ export default class ReleaseWorker extends Worker { const releaseSequence = (lastRelease?.releaseSequence || 0) + 1; - await this.releasesCollection.insertOne({ - projectId, - release, - releaseSequence, - commits: [], - } as unknown as ReleaseDBScheme); + try { + await this.releasesCollection.insertOne({ + projectId, + release, + releaseSequence, + commits: [], + } as unknown as ReleaseDBScheme); + } catch (error) { + if ((error as MongoError).code?.toString() === DB_DUPLICATE_KEY_ERROR) { + this.logger.debug(`Release ${release} for project ${projectId} was created by another worker`); + + return; + } + + throw error; + } } /**