From 8149ae098c765e4088d38f154e3bdee7319a209b Mon Sep 17 00:00:00 2001 From: Peter Date: Wed, 29 Jul 2026 18:54:47 +0300 Subject: [PATCH 1/5] fix(grouper): fix some ts errors (#581) --- workers/grouper/src/data-filter.ts | 13 +++++++------ 1 file changed, 7 insertions(+), 6 deletions(-) diff --git a/workers/grouper/src/data-filter.ts b/workers/grouper/src/data-filter.ts index a7c58c04..e55c492e 100644 --- a/workers/grouper/src/data-filter.ts +++ b/workers/grouper/src/data-filter.ts @@ -36,14 +36,14 @@ const MAX_CODE_LINE_LENGTH = 140; * @param callback - Function to call on each iteration */ function forAll(obj: Record, callback: (path: string[], key: string, obj: Record) => void): void { - const visit = (current, path: string[]): void => { + const visit = (current: Record, path: string[]): void => { for (const key in current) { if (!Object.prototype.hasOwnProperty.call(current, key)) { continue; } const value = current[key]; - if (!(typeof value === 'object' && !Array.isArray(value))) { + if (!(typeof value === 'object' && value !== null && !Array.isArray(value))) { callback(path, key, current); } else { /** @@ -51,7 +51,7 @@ function forAll(obj: Record, callback: (path: string[], key: st * This reduces GC pressure and memory usage for deeply nested objects */ const newPath = path.length < MAX_TRAVERSAL_DEPTH ? path.concat(key) : path; - visit(value, newPath); + visit(value as Record, newPath); } } }; @@ -264,11 +264,12 @@ export default class DataFilter { * * @param field - any object to iterate */ - private processField(field): void { - if (typeof field === 'string') { + private processField(field: unknown): void { + if (typeof field !== 'object' || field === null || Array.isArray(field)) { return; } - forAll(field, (_path, key, obj) => { + + forAll(field as Record, (_path, key, obj) => { obj[key] = this.filterPanNumbers(obj[key]); obj[key] = this.filterSensitiveData(key, obj[key]); }); From f5af8bc545191fddc425e6b649122bcd00651b8d Mon Sep 17 00:00:00 2001 From: Kuchizu <70284260+Kuchizu@users.noreply.github.com> Date: Wed, 26 Aug 2026 20:57:33 +0300 Subject: [PATCH 2/5] Use dailyEvents counter for all workspaces in limiter (#584) --- .env.sample | 3 - workers/limiter/src/index.ts | 128 ++++++++++++++-------------- workers/limiter/tests/index.test.ts | 33 +++---- 3 files changed, 77 insertions(+), 87 deletions(-) diff --git a/.env.sample b/.env.sample index a5ecb4dc..f2df9be2 100644 --- a/.env.sample +++ b/.env.sample @@ -55,8 +55,5 @@ HAWK_CATCHER_TOKEN= ## If true, Grouper worker will send messages about new events to Notifier worker IS_NOTIFIER_WORKER_ENABLED=false -## Comma-separated workspace ids that should use dailyEvents counters in Limiter quota checks. Use * for all workspaces. -LIMITER_DAILY_EVENTS_COUNTER_WORKSPACE_IDS= - ## Url for telegram notifications about workspace blocks and unblocks TELEGRAM_LIMITER_CHAT_URL= diff --git a/workers/limiter/src/index.ts b/workers/limiter/src/index.ts index cc38d768..9df8c970 100644 --- a/workers/limiter/src/index.ts +++ b/workers/limiter/src/index.ts @@ -266,7 +266,7 @@ export default class LimiterWorker extends Worker { const since = Math.floor(new Date(workspace.lastChargeDate).getTime() / MS_IN_SEC); - const workspaceEventsCount = await this.getWorkspaceEventsCount(workspace, projects, since); + const workspaceEventsCount = await this.dbHelper.getEventsCountByProjectsUsingDailyEvents(projects, since); this.logger.info(`workspace ${workspace._id} events count since last charge date: ${workspaceEventsCount}`); @@ -328,68 +328,70 @@ export default class LimiterWorker extends Worker { }; } - /** - * Returns workspace events count using the default raw counter or the - * dailyEvents-based counter when it is explicitly enabled for the workspace. - * - * For enabled workspaces both counters are computed and their results with - * timings are reported to Telegram to compare the algorithms during the - * testing period. The old counter is used as a fallback if the new one fails. - * - * @param workspace - workspace to count events for - * @param projects - workspace projects - * @param since - timestamp of the time from which we count the events - */ - private async getWorkspaceEventsCount( - workspace: WorkspaceWithTariffPlan, - projects: ProjectDBScheme[], - since: number - ): Promise { - if (!this.shouldUseDailyEventsCounter(workspace._id.toString())) { - return this.dbHelper.getEventsCountByProjects(projects, since); - } - - const oldAlgoStartedAt = Date.now(); - const oldAlgoCount = await this.dbHelper.getEventsCountByProjects(projects, since); - const oldAlgoTook = (Date.now() - oldAlgoStartedAt) / MS_IN_SEC; - - try { - const newAlgoStartedAt = Date.now(); - const newAlgoCount = await this.dbHelper.getEventsCountByProjectsUsingDailyEvents(projects, since); - const newAlgoTook = (Date.now() - newAlgoStartedAt) / MS_IN_SEC; - - telegram.sendMessage( - `Workspace ${workspace.name} event count:\n` + - `Old algo: ${oldAlgoCount}, took ${oldAlgoTook}sec\n` + - `New algo: ${newAlgoCount}, took ${newAlgoTook}sec`, - telegram.TelegramBotURLs.Limiter - ); - - return newAlgoCount; - } catch (error) { - HawkCatcher.send(error, { - workspaceId: workspace._id.toString(), - }); - - return oldAlgoCount; - } - } - - /** - * Checks whether dailyEvents-based quota counting is enabled for the workspace - * via LIMITER_DAILY_EVENTS_COUNTER_WORKSPACE_IDS environment variable — - * comma-separated workspace ids or `*` to enable it for every workspace. - * - * @param workspaceId - workspace id - */ - private shouldUseDailyEventsCounter(workspaceId: string): boolean { - const enabledWorkspaceIds = (process.env.LIMITER_DAILY_EVENTS_COUNTER_WORKSPACE_IDS || '') - .split(',') - .map(id => id.trim()) - .filter(Boolean); - - return enabledWorkspaceIds.includes('*') || enabledWorkspaceIds.includes(workspaceId); - } + // Old raw counter with the opt-in switch, kept in case we need to roll back + // + // /** + // * Returns workspace events count using the default raw counter or the + // * dailyEvents-based counter when it is explicitly enabled for the workspace. + // * + // * For enabled workspaces both counters are computed and their results with + // * timings are reported to Telegram to compare the algorithms during the + // * testing period. The old counter is used as a fallback if the new one fails. + // * + // * @param workspace - workspace to count events for + // * @param projects - workspace projects + // * @param since - timestamp of the time from which we count the events + // */ + // private async getWorkspaceEventsCount( + // workspace: WorkspaceWithTariffPlan, + // projects: ProjectDBScheme[], + // since: number + // ): Promise { + // if (!this.shouldUseDailyEventsCounter(workspace._id.toString())) { + // return this.dbHelper.getEventsCountByProjects(projects, since); + // } + // + // const oldAlgoStartedAt = Date.now(); + // const oldAlgoCount = await this.dbHelper.getEventsCountByProjects(projects, since); + // const oldAlgoTook = (Date.now() - oldAlgoStartedAt) / MS_IN_SEC; + // + // try { + // const newAlgoStartedAt = Date.now(); + // const newAlgoCount = await this.dbHelper.getEventsCountByProjectsUsingDailyEvents(projects, since); + // const newAlgoTook = (Date.now() - newAlgoStartedAt) / MS_IN_SEC; + // + // telegram.sendMessage( + // `Workspace ${workspace.name} event count:\n` + + // `Old algo: ${oldAlgoCount}, took ${oldAlgoTook}sec\n` + + // `New algo: ${newAlgoCount}, took ${newAlgoTook}sec`, + // telegram.TelegramBotURLs.Limiter + // ); + // + // return newAlgoCount; + // } catch (error) { + // HawkCatcher.send(error, { + // workspaceId: workspace._id.toString(), + // }); + // + // return oldAlgoCount; + // } + // } + // + // /** + // * Checks whether dailyEvents-based quota counting is enabled for the workspace + // * via LIMITER_DAILY_EVENTS_COUNTER_WORKSPACE_IDS environment variable — + // * comma-separated workspace ids or `*` to enable it for every workspace. + // * + // * @param workspaceId - workspace id + // */ + // private shouldUseDailyEventsCounter(workspaceId: string): boolean { + // const enabledWorkspaceIds = (process.env.LIMITER_DAILY_EVENTS_COUNTER_WORKSPACE_IDS || '') + // .split(',') + // .map(id => id.trim()) + // .filter(Boolean); + // + // return enabledWorkspaceIds.includes('*') || enabledWorkspaceIds.includes(workspaceId); + // } /** * Method that formats project list to html used in report messages diff --git a/workers/limiter/tests/index.test.ts b/workers/limiter/tests/index.test.ts index 1c6c7648..3fd54a5a 100644 --- a/workers/limiter/tests/index.test.ts +++ b/workers/limiter/tests/index.test.ts @@ -331,7 +331,7 @@ describe('Limiter worker', () => { expect(reportMessage).toContain(`${project1.name} (id: ${project1._id})`); }); - test('Should compute both counters and report the comparison to Telegram when dailyEvents counter is enabled', async () => { + test('Should count events via dailyEvents counters for every workspace', async () => { /** * Arrange */ @@ -349,7 +349,7 @@ describe('Limiter worker', () => { }); /** - * Bucket for the day after the boundary day — counted only by the new algorithm + * Bucket for the day after the boundary day — counted via dailyEvents */ await db.collection(`dailyEvents:${project._id.toString()}`).insertOne({ groupHash: 'ade987831d0d0d167aeea685b49db164eb4e113fd027858eef7f69d049357f62', @@ -357,23 +357,17 @@ describe('Limiter worker', () => { count: 7, }); - process.env.LIMITER_DAILY_EVENTS_COUNTER_WORKSPACE_IDS = workspace._id.toString(); - /** * Act */ - try { - const worker = new LimiterWorker(); - - await worker.start(); - await worker.handle(REGULAR_WORKSPACES_CHECK_EVENT); - await worker.finish(); - } finally { - delete process.env.LIMITER_DAILY_EVENTS_COUNTER_WORKSPACE_IDS; - } + const worker = new LimiterWorker(); + + await worker.start(); + await worker.handle(REGULAR_WORKSPACES_CHECK_EVENT); + await worker.finish(); /** - * Assert — the new counter result is saved, both results are reported with timings + * Assert */ const workspaceInDatabase = await workspaceCollection.findOne({ _id: workspace._id, @@ -381,13 +375,10 @@ describe('Limiter worker', () => { expect(workspaceInDatabase.billingPeriodEventsCount).toBe(12); // 5 boundary-day events + 7 from dailyEvents - const comparisonMessage = (telegram.sendMessage as jest.Mock).mock.calls - .map(call => call[0]) - .find(message => message.includes('Old algo')); - - expect(comparisonMessage).toContain(`Workspace ${workspace.name} event count:`); - expect(comparisonMessage).toMatch(/Old algo: 5, took [\d.]+sec/); - expect(comparisonMessage).toMatch(/New algo: 12, took [\d.]+sec/); + /** + * Counters comparison is not reported to Telegram anymore + */ + expect(telegram.sendMessage).not.toHaveBeenCalled(); }); test('Should not send a report when no projects are blocked or unblocked', async () => { From 39ae9b7036ffe58d4d386c50c15acd34e963dc06 Mon Sep 17 00:00:00 2001 From: Kuchizu <70284260+Kuchizu@users.noreply.github.com> Date: Thu, 17 Sep 2026 17:05:41 +0300 Subject: [PATCH 3/5] Stop limiter from opening every events collection each hour (#587) * Stop limiter from opening every events collection each hour * Comment limiter dailyEvents aggregation sections * Report sampled limiter counter validation to Telegram * Compare old and new limiter algo for listed workspaces instead of sampling * Use the previous dailyEvents query for old algo comparison --- .env.sample | 3 + workers/limiter/src/dbHelper.ts | 85 ++++++++++++++++++++++-- workers/limiter/src/index.ts | 38 ++++++++++- workers/limiter/tests/dbHelper.test.ts | 90 +++++++++++++++++++++++++- workers/limiter/tests/index.test.ts | 48 ++++++++++++++ 5 files changed, 258 insertions(+), 6 deletions(-) diff --git a/.env.sample b/.env.sample index f2df9be2..5c19ff64 100644 --- a/.env.sample +++ b/.env.sample @@ -57,3 +57,6 @@ IS_NOTIFIER_WORKER_ENABLED=false ## Url for telegram notifications about workspace blocks and unblocks TELEGRAM_LIMITER_CHAT_URL= + +## Workspace ids to compare old and new limiter counts in telegram +LIMITER_COMPARE_COUNTERS_WORKSPACE_IDS= diff --git a/workers/limiter/src/dbHelper.ts b/workers/limiter/src/dbHelper.ts index 984dd335..0d0e71e5 100644 --- a/workers/limiter/src/dbHelper.ts +++ b/workers/limiter/src/dbHelper.ts @@ -184,6 +184,7 @@ export class DbHelper { * increments `count` for originals and repetitions alike); only the * partial day containing `since` is counted from the raw collections, * since dailyEvents buckets have day granularity and lastChargeDate does not. + * Raw collections are skipped if that day's bucket is empty. * * @param project - project to check * @param since - timestamp of the time from which we count the events @@ -191,6 +192,83 @@ export class DbHelper { public async getEventsCountByProjectUsingDailyEvents( project: ProjectDBScheme, since: number + ): Promise { + try { + const projectId = project._id.toString(); + const dailyEventsCollection = this.eventsDbConnection.collection('dailyEvents:' + projectId); + const boundaryDayTimestamp = since - (since % SEC_IN_DAY); + const firstFullDayTimestamp = this.getFirstFullDailyEventsTimestamp(since); + + const [ counters ] = await dailyEventsCollection + .aggregate<{ boundaryDay: number; fullDays: number }>([ + /** buckets from the day containing `since` onwards */ + { $match: { groupingTimestamp: { $gte: boundaryDayTimestamp } } }, + { + $group: { + _id: null, + /** whole boundary day, only gates the raw count below */ + boundaryDay: { + $sum: { $cond: [ { $lt: ['$groupingTimestamp', firstFullDayTimestamp] }, '$count', 0] }, + }, + /** days after the boundary day */ + fullDays: { + $sum: { $cond: [ { $gte: ['$groupingTimestamp', firstFullDayTimestamp] }, '$count', 0] }, + }, + }, + }, + ], { + /** one table per project instead of racing all groupingTimestamp indexes */ + hint: { $natural: 1 }, + }) + .toArray(); + + /** no buckets in the billing period */ + if (!counters) { + return 0; + } + + /** the bucket spans the whole day, so the part after `since` is counted from raw events */ + const boundaryDayCount = counters.boundaryDay > 0 + ? await this.getRawEventsCountByProject(project, { + timestamp: { + $gt: since, + $lt: firstFullDayTimestamp, + }, + }) + : 0; + + return boundaryDayCount + counters.fullDays; + } catch (e) { + HawkCatcher.send(e); + throw new CriticalError(e); + } + } + + /** + * Calculates total events count for all provided projects since the specific date + * using dailyEvents counters for full days. + * + * @param projects - projects to calculate for + * @param since - timestamp of the time from which we count the events + */ + public async getEventsCountByProjectsUsingDailyEvents(projects: ProjectDBScheme[], since: number): Promise { + const sum = (array: number[]): number => array.reduce((acc, val) => acc + val, 0); + + return Promise.all(projects.map( + project => this.getEventsCountByProjectUsingDailyEvents(project, since) + )) + .then(sum); + } + + /** + * Previous query, kept for rollout comparison + * + * @param project - project to check + * @param since - timestamp of the time from which we count the events + */ + public async getEventsCountByProjectUsingDailyEventsOld( + project: ProjectDBScheme, + since: number ): Promise { try { const projectId = project._id.toString(); @@ -231,17 +309,16 @@ export class DbHelper { } /** - * Calculates total events count for all provided projects since the specific date - * using dailyEvents counters for full days. + * Previous query, kept for rollout comparison * * @param projects - projects to calculate for * @param since - timestamp of the time from which we count the events */ - public async getEventsCountByProjectsUsingDailyEvents(projects: ProjectDBScheme[], since: number): Promise { + public async getEventsCountByProjectsUsingDailyEventsOld(projects: ProjectDBScheme[], since: number): Promise { const sum = (array: number[]): number => array.reduce((acc, val) => acc + val, 0); return Promise.all(projects.map( - project => this.getEventsCountByProjectUsingDailyEvents(project, since) + project => this.getEventsCountByProjectUsingDailyEventsOld(project, since) )) .then(sum); } diff --git a/workers/limiter/src/index.ts b/workers/limiter/src/index.ts index 9df8c970..1b3ca6cd 100644 --- a/workers/limiter/src/index.ts +++ b/workers/limiter/src/index.ts @@ -266,7 +266,7 @@ export default class LimiterWorker extends Worker { const since = Math.floor(new Date(workspace.lastChargeDate).getTime() / MS_IN_SEC); - const workspaceEventsCount = await this.dbHelper.getEventsCountByProjectsUsingDailyEvents(projects, since); + const workspaceEventsCount = await this.getWorkspaceEventsCount(workspace, projects, since); this.logger.info(`workspace ${workspace._id} events count since last charge date: ${workspaceEventsCount}`); @@ -328,6 +328,42 @@ export default class LimiterWorker extends Worker { }; } + /** + * Counts workspace events, comparing with the previous query for LIMITER_COMPARE_COUNTERS_WORKSPACE_IDS + * + * @param workspace - workspace to count events for + * @param projects - workspace projects + * @param since - timestamp of the time from which we count the events + */ + private async getWorkspaceEventsCount( + workspace: WorkspaceWithTariffPlan, + projects: ProjectDBScheme[], + since: number + ): Promise { + const compareWorkspaceIds = (process.env.LIMITER_COMPARE_COUNTERS_WORKSPACE_IDS || '').split(',').map(id => id.trim()); + + if (!compareWorkspaceIds.includes(workspace._id.toString())) { + return this.dbHelper.getEventsCountByProjectsUsingDailyEvents(projects, since); + } + + const oldAlgoStartedAt = Date.now(); + const oldAlgoCount = await this.dbHelper.getEventsCountByProjectsUsingDailyEventsOld(projects, since); + const oldAlgoTook = (Date.now() - oldAlgoStartedAt) / MS_IN_SEC; + + const newAlgoStartedAt = Date.now(); + const newAlgoCount = await this.dbHelper.getEventsCountByProjectsUsingDailyEvents(projects, since); + const newAlgoTook = (Date.now() - newAlgoStartedAt) / MS_IN_SEC; + + telegram.sendMessage( + `Workspace ${workspace.name} event count:\n` + + `Old algo: ${oldAlgoCount}, took ${oldAlgoTook}sec\n` + + `New algo: ${newAlgoCount}, took ${newAlgoTook}sec`, + telegram.TelegramBotURLs.Limiter + ); + + return newAlgoCount; + } + // Old raw counter with the opt-in switch, kept in case we need to roll back // // /** diff --git a/workers/limiter/tests/dbHelper.test.ts b/workers/limiter/tests/dbHelper.test.ts index dd26baa7..210d486f 100644 --- a/workers/limiter/tests/dbHelper.test.ts +++ b/workers/limiter/tests/dbHelper.test.ts @@ -23,6 +23,9 @@ const BOUNDARY_DAY_TIMESTAMP = 1585756800; */ const NEXT_MIDNIGHT_AFTER_LAST_CHARGE = 1585785600; +/** 2020-04-01T00:00:00Z */ +const BOUNDARY_DAY_MIDNIGHT = NEXT_MIDNIGHT_AFTER_LAST_CHARGE - 86400; + describe('DbHelper', () => { let connection: MongoClient; let db: Db; @@ -132,6 +135,17 @@ describe('DbHelper', () => { await repetitionsCollection.insertMany(mockedEvents); } + /** as grouper does */ + const boundaryDayEventsCount = parameters.eventsToMock + (parameters.repetitionsToMock ?? 0); + + if (boundaryDayEventsCount > 0) { + await dailyEventsCollection.insertOne({ + groupHash: 'ade987831d0d0d167aeea685b49db164eb4e113fd027858eef7f69d049357f62', + groupingTimestamp: BOUNDARY_DAY_MIDNIGHT, + count: boundaryDayEventsCount, + }); + } + if (parameters.dailyEventsToMock?.length > 0) { await dailyEventsCollection.insertMany(parameters.dailyEventsToMock.map(bucket => ({ groupHash: 'ade987831d0d0d167aeea685b49db164eb4e113fd027858eef7f69d049357f62', @@ -711,7 +725,7 @@ describe('DbHelper', () => { dailyEventsToMock: [ /** bucket of the boundary day itself must not be counted */ { - groupingTimestamp: NEXT_MIDNIGHT_AFTER_LAST_CHARGE - 86400, + groupingTimestamp: BOUNDARY_DAY_MIDNIGHT, count: 100, }, ], @@ -779,6 +793,80 @@ describe('DbHelper', () => { */ expect(count).toBe(7); }); + + test('Should not query raw collections when the boundary day bucket is empty', async () => { + /** + * Arrange + */ + const workspace = createWorkspaceMock({ + plan: mockedPlans.eventsLimit10, + billingPeriodEventsCount: 0, + lastChargeDate: new Date(), + }); + const project = createProjectMock({ workspaceId: workspace._id }); + const since = Math.floor(LAST_CHARGE_DATE.getTime() / MS_IN_SEC); + + await fillDatabaseWithMockedData({ + workspace, + project, + eventsToMock: 0, + dailyEventsToMock: [ + { + groupingTimestamp: NEXT_MIDNIGHT_AFTER_LAST_CHARGE, + count: 4, + }, + ], + }); + + const collectionSpy = jest.spyOn(db, 'collection'); + + /** + * Act + */ + const count = await dbHelper.getEventsCountByProjectUsingDailyEvents(project, since); + + /** + * Assert + */ + expect(count).toBe(4); + expect(collectionSpy.mock.calls.map(([ name ]) => name)).toEqual([ `dailyEvents:${project._id.toString()}` ]); + + collectionSpy.mockRestore(); + }); + + test('Should return zero for a project without dailyEvents collection', async () => { + const project = createProjectMock({ workspaceId: new ObjectId() }); + const since = Math.floor(LAST_CHARGE_DATE.getTime() / MS_IN_SEC); + + const count = await dbHelper.getEventsCountByProjectUsingDailyEvents(project, since); + + expect(count).toBe(0); + }); + }); + + describe('getEventsCountByProjectUsingDailyEventsOld', () => { + test('Should count raw boundary-day events even without a dailyEvents bucket', async () => { + const project = createProjectMock({ workspaceId: new ObjectId() }); + const since = Math.floor(LAST_CHARGE_DATE.getTime() / MS_IN_SEC); + + await fillDatabaseWithMockedData({ + project, + eventsToMock: 0, + dailyEventsToMock: [ + { + groupingTimestamp: NEXT_MIDNIGHT_AFTER_LAST_CHARGE, + count: 3, + }, + ], + }); + await db.collection(`events:${project._id.toString()}`).insertMany([createEventMock(), createEventMock()]); + + const oldCount = await dbHelper.getEventsCountByProjectUsingDailyEventsOld(project, since); + const newCount = await dbHelper.getEventsCountByProjectUsingDailyEvents(project, since); + + expect(oldCount).toBe(5); + expect(newCount).toBe(3); + }); }); describe('getEventsCountByProjectsUsingDailyEvents', () => { diff --git a/workers/limiter/tests/index.test.ts b/workers/limiter/tests/index.test.ts index 3fd54a5a..d233a06f 100644 --- a/workers/limiter/tests/index.test.ts +++ b/workers/limiter/tests/index.test.ts @@ -138,6 +138,13 @@ describe('Limiter worker', () => { } await repetitionsCollection.insertMany(mockedEvents); } + + /** as grouper does */ + await db.collection(`dailyEvents:${parameters.project._id.toString()}`).insertOne({ + groupHash: 'ade987831d0d0d167aeea685b49db164eb4e113fd027858eef7f69d049357f62', + groupingTimestamp: NEXT_MIDNIGHT_AFTER_LAST_CHARGE - 86400, + count: parameters.eventsToMock + (parameters.repetitionsToMock ?? 0), + }); }; beforeAll(async () => { @@ -381,6 +388,47 @@ describe('Limiter worker', () => { expect(telegram.sendMessage).not.toHaveBeenCalled(); }); + test('Should report old and new algo counts for workspaces listed for comparison', async () => { + const workspace = createWorkspaceMock({ + plan: mockedPlans.eventsLimit10000, + billingPeriodEventsCount: 0, + lastChargeDate: LAST_CHARGE_DATE, + }); + const project = createProjectMock({ workspaceId: workspace._id }); + + await fillDatabaseWithMockedData({ + workspace, + project, + eventsToMock: 5, + }); + + await db.collection(`dailyEvents:${project._id.toString()}`).deleteMany({}); + + process.env.LIMITER_COMPARE_COUNTERS_WORKSPACE_IDS = workspace._id.toString(); + + const worker = new LimiterWorker(); + + try { + await worker.start(); + await worker.handle(REGULAR_WORKSPACES_CHECK_EVENT); + await worker.finish(); + } finally { + delete process.env.LIMITER_COMPARE_COUNTERS_WORKSPACE_IDS; + } + + const workspaceInDatabase = await workspaceCollection.findOne({ + _id: workspace._id, + }); + + expect(workspaceInDatabase.billingPeriodEventsCount).toBe(0); + expect(telegram.sendMessage).toHaveBeenCalledTimes(1); + + const reportMessage = (telegram.sendMessage as jest.Mock).mock.calls[0][0]; + + expect(reportMessage).toContain('Old algo: 5'); + expect(reportMessage).toContain('New algo: 0'); + }); + test('Should not send a report when no projects are blocked or unblocked', async () => { /** * Arrange From a6b17578b80a114a0578454bb8d1a04132d439eb Mon Sep 17 00:00:00 2001 From: Kuchizu <70284260+Kuchizu@users.noreply.github.com> Date: Thu, 17 Sep 2026 22:26:06 +0300 Subject: [PATCH 4/5] Remove old limiter algo comparison (#591) --- .env.sample | 3 -- workers/limiter/src/dbHelper.ts | 63 -------------------------- workers/limiter/src/index.ts | 38 +--------------- workers/limiter/tests/dbHelper.test.ts | 25 ---------- workers/limiter/tests/index.test.ts | 41 ----------------- 5 files changed, 1 insertion(+), 169 deletions(-) diff --git a/.env.sample b/.env.sample index 5c19ff64..f2df9be2 100644 --- a/.env.sample +++ b/.env.sample @@ -57,6 +57,3 @@ IS_NOTIFIER_WORKER_ENABLED=false ## Url for telegram notifications about workspace blocks and unblocks TELEGRAM_LIMITER_CHAT_URL= - -## Workspace ids to compare old and new limiter counts in telegram -LIMITER_COMPARE_COUNTERS_WORKSPACE_IDS= diff --git a/workers/limiter/src/dbHelper.ts b/workers/limiter/src/dbHelper.ts index 0d0e71e5..e6948e06 100644 --- a/workers/limiter/src/dbHelper.ts +++ b/workers/limiter/src/dbHelper.ts @@ -260,69 +260,6 @@ export class DbHelper { .then(sum); } - /** - * Previous query, kept for rollout comparison - * - * @param project - project to check - * @param since - timestamp of the time from which we count the events - */ - public async getEventsCountByProjectUsingDailyEventsOld( - project: ProjectDBScheme, - since: number - ): Promise { - try { - const projectId = project._id.toString(); - const dailyEventsCollection = this.eventsDbConnection.collection('dailyEvents:' + projectId); - const firstFullDayTimestamp = this.getFirstFullDailyEventsTimestamp(since); - - const boundaryDayQuery = { - timestamp: { - $gt: since, - $lt: firstFullDayTimestamp, - }, - }; - - const [boundaryDayCount, dailyCounters] = await Promise.all([ - since < firstFullDayTimestamp - ? this.getRawEventsCountByProject(project, boundaryDayQuery) - : 0, - dailyEventsCollection - .aggregate<{ count: number }>([ - { $match: { groupingTimestamp: { $gte: firstFullDayTimestamp } } }, - { - $group: { - _id: null, - count: { $sum: '$count' }, - }, - }, - ]) - .toArray(), - ]); - - const fullDaysCount = dailyCounters.length > 0 ? dailyCounters[0].count : 0; - - return boundaryDayCount + fullDaysCount; - } catch (e) { - HawkCatcher.send(e); - throw new CriticalError(e); - } - } - - /** - * Previous query, kept for rollout comparison - * - * @param projects - projects to calculate for - * @param since - timestamp of the time from which we count the events - */ - public async getEventsCountByProjectsUsingDailyEventsOld(projects: ProjectDBScheme[], since: number): Promise { - const sum = (array: number[]): number => array.reduce((acc, val) => acc + val, 0); - - return Promise.all(projects.map( - project => this.getEventsCountByProjectUsingDailyEventsOld(project, since) - )) - .then(sum); - } - /** * Returns all projects from Database or projects of the specified workspace * diff --git a/workers/limiter/src/index.ts b/workers/limiter/src/index.ts index 1b3ca6cd..9df8c970 100644 --- a/workers/limiter/src/index.ts +++ b/workers/limiter/src/index.ts @@ -266,7 +266,7 @@ export default class LimiterWorker extends Worker { const since = Math.floor(new Date(workspace.lastChargeDate).getTime() / MS_IN_SEC); - const workspaceEventsCount = await this.getWorkspaceEventsCount(workspace, projects, since); + const workspaceEventsCount = await this.dbHelper.getEventsCountByProjectsUsingDailyEvents(projects, since); this.logger.info(`workspace ${workspace._id} events count since last charge date: ${workspaceEventsCount}`); @@ -328,42 +328,6 @@ export default class LimiterWorker extends Worker { }; } - /** - * Counts workspace events, comparing with the previous query for LIMITER_COMPARE_COUNTERS_WORKSPACE_IDS - * - * @param workspace - workspace to count events for - * @param projects - workspace projects - * @param since - timestamp of the time from which we count the events - */ - private async getWorkspaceEventsCount( - workspace: WorkspaceWithTariffPlan, - projects: ProjectDBScheme[], - since: number - ): Promise { - const compareWorkspaceIds = (process.env.LIMITER_COMPARE_COUNTERS_WORKSPACE_IDS || '').split(',').map(id => id.trim()); - - if (!compareWorkspaceIds.includes(workspace._id.toString())) { - return this.dbHelper.getEventsCountByProjectsUsingDailyEvents(projects, since); - } - - const oldAlgoStartedAt = Date.now(); - const oldAlgoCount = await this.dbHelper.getEventsCountByProjectsUsingDailyEventsOld(projects, since); - const oldAlgoTook = (Date.now() - oldAlgoStartedAt) / MS_IN_SEC; - - const newAlgoStartedAt = Date.now(); - const newAlgoCount = await this.dbHelper.getEventsCountByProjectsUsingDailyEvents(projects, since); - const newAlgoTook = (Date.now() - newAlgoStartedAt) / MS_IN_SEC; - - telegram.sendMessage( - `Workspace ${workspace.name} event count:\n` + - `Old algo: ${oldAlgoCount}, took ${oldAlgoTook}sec\n` + - `New algo: ${newAlgoCount}, took ${newAlgoTook}sec`, - telegram.TelegramBotURLs.Limiter - ); - - return newAlgoCount; - } - // Old raw counter with the opt-in switch, kept in case we need to roll back // // /** diff --git a/workers/limiter/tests/dbHelper.test.ts b/workers/limiter/tests/dbHelper.test.ts index 210d486f..bf9aa46c 100644 --- a/workers/limiter/tests/dbHelper.test.ts +++ b/workers/limiter/tests/dbHelper.test.ts @@ -844,31 +844,6 @@ describe('DbHelper', () => { }); }); - describe('getEventsCountByProjectUsingDailyEventsOld', () => { - test('Should count raw boundary-day events even without a dailyEvents bucket', async () => { - const project = createProjectMock({ workspaceId: new ObjectId() }); - const since = Math.floor(LAST_CHARGE_DATE.getTime() / MS_IN_SEC); - - await fillDatabaseWithMockedData({ - project, - eventsToMock: 0, - dailyEventsToMock: [ - { - groupingTimestamp: NEXT_MIDNIGHT_AFTER_LAST_CHARGE, - count: 3, - }, - ], - }); - await db.collection(`events:${project._id.toString()}`).insertMany([createEventMock(), createEventMock()]); - - const oldCount = await dbHelper.getEventsCountByProjectUsingDailyEventsOld(project, since); - const newCount = await dbHelper.getEventsCountByProjectUsingDailyEvents(project, since); - - expect(oldCount).toBe(5); - expect(newCount).toBe(3); - }); - }); - describe('getEventsCountByProjectsUsingDailyEvents', () => { test('Should count events, repetitions and dailyEvents for multiple projects', async () => { /** diff --git a/workers/limiter/tests/index.test.ts b/workers/limiter/tests/index.test.ts index d233a06f..f49c987b 100644 --- a/workers/limiter/tests/index.test.ts +++ b/workers/limiter/tests/index.test.ts @@ -388,47 +388,6 @@ describe('Limiter worker', () => { expect(telegram.sendMessage).not.toHaveBeenCalled(); }); - test('Should report old and new algo counts for workspaces listed for comparison', async () => { - const workspace = createWorkspaceMock({ - plan: mockedPlans.eventsLimit10000, - billingPeriodEventsCount: 0, - lastChargeDate: LAST_CHARGE_DATE, - }); - const project = createProjectMock({ workspaceId: workspace._id }); - - await fillDatabaseWithMockedData({ - workspace, - project, - eventsToMock: 5, - }); - - await db.collection(`dailyEvents:${project._id.toString()}`).deleteMany({}); - - process.env.LIMITER_COMPARE_COUNTERS_WORKSPACE_IDS = workspace._id.toString(); - - const worker = new LimiterWorker(); - - try { - await worker.start(); - await worker.handle(REGULAR_WORKSPACES_CHECK_EVENT); - await worker.finish(); - } finally { - delete process.env.LIMITER_COMPARE_COUNTERS_WORKSPACE_IDS; - } - - const workspaceInDatabase = await workspaceCollection.findOne({ - _id: workspace._id, - }); - - expect(workspaceInDatabase.billingPeriodEventsCount).toBe(0); - expect(telegram.sendMessage).toHaveBeenCalledTimes(1); - - const reportMessage = (telegram.sendMessage as jest.Mock).mock.calls[0][0]; - - expect(reportMessage).toContain('Old algo: 5'); - expect(reportMessage).toContain('New algo: 0'); - }); - test('Should not send a report when no projects are blocked or unblocked', async () => { /** * Arrange From 9d9c98a632be04d37d59684152d2a0080e58cb2d Mon Sep 17 00:00:00 2001 From: Kuchizu <70284260+Kuchizu@users.noreply.github.com> Date: Sat, 26 Sep 2026 19:57:49 +0300 Subject: [PATCH 5/5] Exit runner when a worker fails to start (#594) --- runner.ts | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/runner.ts b/runner.ts index 1c19f97b..e5f20d56 100644 --- a/runner.ts +++ b/runner.ts @@ -154,7 +154,9 @@ class WorkerRunner { utils.sendReport(worker.constructor.name + ' failed to start'); - await this.stopWorker(worker); + await this.finishAll(); + + process.exit(1); } }) );