diff --git a/package.json b/package.json index f2e1e154..e3f97b82 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 ea9282dc..add96185 100644 --- a/workers/release/src/index.ts +++ b/workers/release/src/index.ts @@ -89,11 +89,21 @@ 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) { + this.logger.debug(`Skipping release ${payload.release} for project ${projectId}: no valid commits or source maps`); + + return; + } + + await this.createReleaseIfMissing(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 +116,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 +130,54 @@ export default class ReleaseWorker extends Worker { } } + /** + * Create a release once and assign its first-registration sequence. + * + * 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 + */ + private async createReleaseIfMissing(projectId: string, release: string): Promise { + const existingRelease = await this.releasesCollection.findOne({ + projectId, + release, + }); + + if (existingRelease) { + return; + } + + const lastRelease = await this.releasesCollection.findOne({ + projectId, + releaseSequence: { $exists: true }, + }, { + sort: { releaseSequence: -1 }, + projection: { releaseSequence: 1 }, + }); + + const releaseSequence = (lastRelease?.releaseSequence || 0) + 1; + + 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; + } + } + /** * Check commtis for the content of all required data * @@ -208,31 +264,8 @@ 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. */ - if (!existedRelease) { - this.logger.info('trying insert new release'); - - try { - await this.releasesCollection.insertOne({ - projectId: projectId, - release: payload.release, - files: savedFilesWithoutContent, - } 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 24f5f46e..b464d06f 100644 --- a/workers/release/tests/index.test.ts +++ b/workers/release/tests/index.test.ts @@ -69,6 +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 mockBundle.build(); }); @@ -184,13 +194,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 +215,37 @@ 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 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 }) + .toArray(); - await expect(count).toEqual(1); + await expect(releases).toHaveLength(2); + await expect(releases.map(release => release.releaseSequence)).toEqual([1, 2]); }); test('should correctly handle release with multiple source maps in a single transaction', async () => {