diff --git a/common/changes/@hcengineering/core/main_2025-10-08-03-37.json b/common/changes/@hcengineering/core/main_2025-10-08-03-37.json new file mode 100644 index 0000000000..1398c67aa6 --- /dev/null +++ b/common/changes/@hcengineering/core/main_2025-10-08-03-37.json @@ -0,0 +1,10 @@ +{ + "changes": [ + { + "packageName": "@hcengineering/core", + "comment": "Fix issue with time rate limiter", + "type": "patch" + } + ], + "packageName": "@hcengineering/core" +} \ No newline at end of file diff --git a/common/config/rush/version-policies.json b/common/config/rush/version-policies.json index d03b73a3b6..00501e9e33 100644 --- a/common/config/rush/version-policies.json +++ b/common/config/rush/version-policies.json @@ -20,21 +20,21 @@ // * SemVer range is usually restricted to a single version. // */ // "definitionName": "lockStepVersion", - // + // // /** // * (Required) The name that will be used for the "versionPolicyName" field in rush.json. // * This name is also used command-line parameters such as "--version-policy" // * and "--to-version-policy". // */ // "policyName": "MyBigFramework", - // + // // /** // * (Required) The current version. All packages belonging to the set should have this version // * in the current branch. When bumping versions, Rush uses this to determine the next version. // * (The "version" field in package.json is NOT considered.) // */ // "version": "1.0.0", - // + // // /** // * (Required) The type of bump that will be performed when publishing the next release. // * When creating a release branch in Git, this field should be updated according to the @@ -43,7 +43,7 @@ // * Valid values are: "prerelease", "preminor", "minor", "patch", "major" // */ // "nextBump": "prerelease", - // + // // /** // * (Optional) If specified, all packages in the set share a common CHANGELOG.md file. // * This file is stored with the specified "main" project, which must be a member of the set. @@ -52,7 +52,7 @@ // * package in the set. // */ // "mainProject": "my-app", - // + // // /** // * (Optional) If enabled, the "rush change" command will prompt the user for their email address // * and include it in the JSON change files. If an organization maintains multiple repos, tracking @@ -63,40 +63,40 @@ // */ // // "includeEmailInChangeFile": true // }, - // - // { - // /** - // * (Required) Indicates the kind of version policy being defined ("lockStepVersion" or "individualVersion"). - // * - // * The "individualVersion" mode specifies that the projects will use "individual versioning". - // * This is the typical NPM model where each package has an independent version number - // * and CHANGELOG.md file. Although a single CI definition is responsible for publishing the - // * packages, they otherwise don't have any special relationship. The version bumping will - // * depend on how developers answer the "rush change" questions for each package that - // * is changed. - // */ - // "definitionName": "individualVersion", - // - // "policyName": "MyRandomLibraries", - // - // /** - // * (Optional) This can be used to enforce that all packages in the set must share a common - // * major version number, e.g. because they are from the same major release branch. - // * It can also be used to discourage people from accidentally making "MAJOR" SemVer changes - // * inappropriately. The minor/patch version parts will be bumped independently according - // * to the types of changes made to each project, according to the "rush change" command. - // */ - // "lockedMajor": 3, - // - // /** - // * (Optional) When publishing is managed by Rush, by default the "rush change" command will - // * request changes for any projects that are modified by a pull request. These change entries - // * will produce a CHANGELOG.md file. If you author your CHANGELOG.md manually or announce updates - // * in some other way, set "exemptFromRushChange" to true to tell "rush change" to ignore the projects - // * belonging to this version policy. - // */ - // "exemptFromRushChange": false, - // - // // "includeEmailInChangeFile": true - // } + // + { + /** + * (Required) Indicates the kind of version policy being defined ("lockStepVersion" or "individualVersion"). + * + * The "individualVersion" mode specifies that the projects will use "individual versioning". + * This is the typical NPM model where each package has an independent version number + * and CHANGELOG.md file. Although a single CI definition is responsible for publishing the + * packages, they otherwise don't have any special relationship. The version bumping will + * depend on how developers answer the "rush change" questions for each package that + * is changed. + */ + "definitionName": "individualVersion", + + "policyName": "Huly.core", + + /** + * (Optional) This can be used to enforce that all packages in the set must share a common + * major version number, e.g. because they are from the same major release branch. + * It can also be used to discourage people from accidentally making "MAJOR" SemVer changes + * inappropriately. The minor/patch version parts will be bumped independently according + * to the types of changes made to each project, according to the "rush change" command. + */ + "lockedMajor": 0.7, + + /** + * (Optional) When publishing is managed by Rush, by default the "rush change" command will + * request changes for any projects that are modified by a pull request. These change entries + * will produce a CHANGELOG.md file. If you author your CHANGELOG.md manually or announce updates + * in some other way, set "exemptFromRushChange" to true to tell "rush change" to ignore the projects + * belonging to this version policy. + */ + "exemptFromRushChange": false + + // "includeEmailInChangeFile": true + } ] diff --git a/packages/core/src/__tests__/limiter-edge-cases.test.ts b/packages/core/src/__tests__/limiter-edge-cases.test.ts new file mode 100644 index 0000000000..395e3fbfe8 --- /dev/null +++ b/packages/core/src/__tests__/limiter-edge-cases.test.ts @@ -0,0 +1,380 @@ +// +// Copyright © 2023 Hardcore Engineering Inc. +// +// Licensed under the Eclipse Public License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. You may +// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// +// See the License for the specific language governing permissions and +// limitations under the License. +// + +import { RateLimiter, TimeRateLimiter } from '../utils' + +describe('RateLimiter and TimeRateLimiter - Advanced Edge Cases', () => { + describe('RateLimiter - Memory and Resource Management', () => { + it('should not leak memory in notify array', async () => { + const limiter = new RateLimiter(1) + const mockFn = jest.fn().mockImplementation(async () => { + await new Promise((resolve) => setTimeout(resolve, 10)) + return 'result' + }) + + // Execute multiple operations that will queue + const operations = Array(10) + .fill(0) + .map(() => limiter.exec(mockFn)) + + await Promise.all(operations) + + // notify array should be empty after all operations complete + expect(limiter.notify.length).toBe(0) + }) + + it('should clean up processingQueue correctly after many operations', async () => { + const limiter = new RateLimiter(5) + const mockFn = jest.fn().mockImplementation(async (args?: any) => { + await new Promise((resolve) => setTimeout(resolve, Math.random() * 10)) + return args?.id + }) + + for (let batch = 0; batch < 5; batch++) { + const operations = Array(20) + .fill(0) + .map((_, i) => limiter.exec(mockFn, { id: batch * 20 + i })) + + await Promise.all(operations) + expect(limiter.processingQueue.size).toBe(0) + } + + expect(mockFn).toHaveBeenCalledTimes(100) + }) + + it('should handle interleaved exec and add operations', async () => { + const limiter = new RateLimiter(2) + const results: string[] = [] + + const execOp = async (id: string): Promise => { + await new Promise((resolve) => setTimeout(resolve, 10)) + results.push(id) + return id + } + + const addOp = async (id: string): Promise => { + await new Promise((resolve) => setTimeout(resolve, 10)) + results.push(id) + } + + await Promise.all([ + limiter.exec(async () => await execOp('exec1')), + limiter.add(async () => { + await addOp('add1') + }), + limiter.exec(async () => await execOp('exec2')), + limiter.add(async () => { + await addOp('add2') + }) + ]) + + await limiter.waitProcessing() + + expect(results).toHaveLength(4) + expect(results).toContain('exec1') + expect(results).toContain('exec2') + expect(results).toContain('add1') + expect(results).toContain('add2') + }) + + it('should handle rapid creation and destruction of operations', async () => { + const limiter = new RateLimiter(3) + let successCount = 0 + let errorCount = 0 + + const operations = Array(30) + .fill(0) + .map(async (_, i) => { + const mockFn = jest.fn().mockImplementation(async () => { + await new Promise((resolve) => setTimeout(resolve, 5)) + if (i % 5 === 0) { + throw new Error(`Error ${i}`) + } + return `Result ${i}` + }) + + try { + await limiter.exec(mockFn) + successCount++ + } catch (err) { + errorCount++ + } + }) + + await Promise.all(operations) + + expect(successCount).toBe(24) + expect(errorCount).toBe(6) + expect(limiter.processingQueue.size).toBe(0) + }) + }) + + describe('TimeRateLimiter - Memory and Resource Management', () => { + it('should not leak memory in notify array', async () => { + jest.useFakeTimers() + const limiter = new TimeRateLimiter(2, 1000) + const mockFn = jest.fn().mockImplementation(async () => { + await new Promise((resolve) => setTimeout(resolve, 10)) + return 'result' + }) + + // Execute operations that will need to wait + const operations = [limiter.exec(mockFn), limiter.exec(mockFn), limiter.exec(mockFn), limiter.exec(mockFn)] + + jest.advanceTimersByTime(20) + await Promise.resolve() + await Promise.resolve() + + jest.advanceTimersByTime(1001) + await Promise.resolve() + await Promise.resolve() + + jest.advanceTimersByTime(20) + await Promise.resolve() + await Promise.resolve() + + await Promise.all(operations) + + // notify array should be empty after all operations complete + expect(limiter.notify.length).toBe(0) + jest.useRealTimers() + }, 10000) + + it('should not accumulate executions indefinitely', async () => { + jest.useFakeTimers() + const limiter = new TimeRateLimiter(5, 1000) + const mockFn = jest.fn().mockResolvedValue('result') + + // Execute many batches + for (let batch = 0; batch < 10; batch++) { + const operations = Array(5) + .fill(0) + .map(() => limiter.exec(mockFn)) + + await Promise.all(operations) + + jest.advanceTimersByTime(1100) + } + + // Executions should be cleaned up + // Only the most recent batch should remain (or less) + expect(limiter.executions.length).toBeLessThanOrEqual(5) + jest.useRealTimers() + }) + + it('should handle operations that never resolve gracefully', async () => { + const limiter = new TimeRateLimiter(2, 1000) + + const normalOp = jest.fn().mockImplementation(async () => { + await new Promise((resolve) => setTimeout(resolve, 10)) + return 'result' + }) + + // Note: We can't actually test hanging operations without causing issues + // Instead, test that slow operations don't block faster ones + const slowOp = jest.fn().mockImplementation(async () => { + await new Promise((resolve) => setTimeout(resolve, 100)) + return 'slow' + }) + + void limiter.exec(slowOp) + expect(limiter.active).toBe(1) + + // Other operations should still work + const result = await limiter.exec(normalOp) + expect(result).toBe('result') + expect(normalOp).toHaveBeenCalledTimes(1) + }, 10000) + }) + + describe('RateLimiter - Boundary Conditions', () => { + it('should handle operations that complete immediately', async () => { + const limiter = new RateLimiter(3) + const mockFn = jest.fn().mockResolvedValue('instant') + + const results = await Promise.all([limiter.exec(mockFn), limiter.exec(mockFn), limiter.exec(mockFn)]) + + expect(results).toEqual(['instant', 'instant', 'instant']) + expect(limiter.processingQueue.size).toBe(0) + }) + + it('should maintain correct state when operations fail at different stages', async () => { + const limiter = new RateLimiter(2) + const errors: string[] = [] + + const operations = [ + limiter + .exec(async () => { + await new Promise((resolve) => setTimeout(resolve, 10)) + throw new Error('error1') + }) + .catch((e) => { + errors.push(e.message) + }), + + limiter.exec(async () => { + await new Promise((resolve) => setTimeout(resolve, 5)) + return 'success1' + }), + + limiter + .exec(async () => { + throw new Error('error2') + }) + .catch((e) => { + errors.push(e.message) + }), + + limiter.exec(async () => { + await new Promise((resolve) => setTimeout(resolve, 15)) + return 'success2' + }) + ] + + const results = await Promise.all(operations) + + expect(errors).toHaveLength(2) + expect(errors).toContain('error1') + expect(errors).toContain('error2') + expect(results.filter(Boolean)).toContain('success1') + expect(results.filter(Boolean)).toContain('success2') + expect(limiter.processingQueue.size).toBe(0) + }, 10000) + }) + + describe('TimeRateLimiter - Boundary Conditions', () => { + it('should handle operations completing in reverse order', async () => { + jest.useFakeTimers() + const limiter = new TimeRateLimiter(3, 1000) + const completionOrder: number[] = [] + + const createOp = (id: number, delay: number) => async (): Promise => { + await new Promise((resolve) => setTimeout(resolve, delay)) + completionOrder.push(id) + return id + } + + const operations = [limiter.exec(createOp(1, 30)), limiter.exec(createOp(2, 20)), limiter.exec(createOp(3, 10))] + + jest.advanceTimersByTime(11) + await Promise.resolve() + + jest.advanceTimersByTime(10) + await Promise.resolve() + + jest.advanceTimersByTime(10) + await Promise.resolve() + + await Promise.all(operations) + + // Operations should complete in reverse order of their delays + expect(completionOrder).toEqual([3, 2, 1]) + jest.useRealTimers() + }) + + it('should handle rate limit at exact boundaries', async () => { + jest.useFakeTimers() + const limiter = new TimeRateLimiter(2, 1000) + const mockFn = jest.fn().mockResolvedValue('result') + + // Execute exactly at the rate + void limiter.exec(mockFn) + void limiter.exec(mockFn) + + expect(mockFn).toHaveBeenCalledTimes(2) + expect(limiter.active).toBe(2) + + // One more should wait + const thirdOp = limiter.exec(mockFn) + expect(mockFn).toHaveBeenCalledTimes(2) + + // Advance to exactly the period + jest.advanceTimersByTime(1000) + await Promise.resolve() + await Promise.resolve() + + await thirdOp + expect(mockFn).toHaveBeenCalledTimes(3) + jest.useRealTimers() + }) + }) + + describe('Integration Scenarios', () => { + it('should handle mixed RateLimiter operations with varying complexities', async () => { + const limiter = new RateLimiter(3) + const results: Array = [] + + const stringOp = async (val: string): Promise => { + await new Promise((resolve) => setTimeout(resolve, 5)) + results.push(val) + return val + } + + const numberOp = async (val: number): Promise => { + await new Promise((resolve) => setTimeout(resolve, 10)) + results.push(val) + return val + } + + const objectOp = async (obj: { id: number }): Promise<{ id: number }> => { + await new Promise((resolve) => setTimeout(resolve, 7)) + results.push(obj.id) + return obj + } + + await Promise.all([ + limiter.exec(async () => await stringOp('a')), + limiter.exec(async () => await numberOp(1)), + limiter.exec(async () => await objectOp({ id: 100 })), + limiter.exec(async () => await stringOp('b')), + limiter.exec(async () => await numberOp(2)) + ]) + + expect(results).toHaveLength(5) + expect(results).toContain('a') + expect(results).toContain('b') + expect(results).toContain(1) + expect(results).toContain(2) + expect(results).toContain(100) + }) + + it('should handle nested rate limiters', async () => { + const outerLimiter = new RateLimiter(2) + const innerLimiter = new RateLimiter(1) + + let executionCount = 0 + + const nestedOp = async (id: number): Promise => { + return await innerLimiter.exec(async () => { + executionCount++ + await new Promise((resolve) => setTimeout(resolve, 5)) + return id + }) + } + + const results = await Promise.all([ + outerLimiter.exec(async () => await nestedOp(1)), + outerLimiter.exec(async () => await nestedOp(2)), + outerLimiter.exec(async () => await nestedOp(3)), + outerLimiter.exec(async () => await nestedOp(4)) + ]) + + expect(results).toEqual([1, 2, 3, 4]) + expect(executionCount).toBe(4) + expect(outerLimiter.processingQueue.size).toBe(0) + expect(innerLimiter.processingQueue.size).toBe(0) + }) + }) +}) diff --git a/packages/core/src/__tests__/limits.test.ts b/packages/core/src/__tests__/limits.test.ts index 2fc323d9f2..2bfc263c15 100644 --- a/packages/core/src/__tests__/limits.test.ts +++ b/packages/core/src/__tests__/limits.test.ts @@ -16,6 +16,21 @@ import { TimeRateLimiter } from '../utils' describe('TimeRateLimiter', () => { + describe('constructor', () => { + it('should initialize with correct rate and period', () => { + const limiter = new TimeRateLimiter(5, 2000) + expect(limiter.rate).toBe(5) + expect(limiter.period).toBe(2000) + expect(limiter.active).toBe(0) + expect(limiter.executions).toEqual([]) + }) + + it('should use default period of 1000ms', () => { + const limiter = new TimeRateLimiter(3) + expect(limiter.period).toBe(1000) + }) + }) + it('should limit rate of executions', async () => { jest.useFakeTimers() const limiter = new TimeRateLimiter(2, 1000) // 2 executions per second @@ -38,6 +53,7 @@ describe('TimeRateLimiter', () => { expect(mockFn).toHaveBeenCalledTimes(4) await Promise.all(operations) + jest.useRealTimers() }) it('should cleanup old executions', async () => { @@ -57,6 +73,7 @@ describe('TimeRateLimiter', () => { await limiter.exec(mockFn) expect(mockFn).toHaveBeenCalledTimes(2) expect(limiter.executions.length).toBe(1) // Old execution should be cleaned up + jest.useRealTimers() }) it('should handle concurrent operations', async () => { @@ -89,6 +106,7 @@ describe('TimeRateLimiter', () => { expect(mockFn).toHaveBeenCalledTimes(3) await operations + jest.useRealTimers() }) it('should wait for processing to complete', async () => { @@ -114,5 +132,331 @@ describe('TimeRateLimiter', () => { await waitPromise await operation expect(limiter.active).toBe(0) + jest.useRealTimers() + }) + + describe('execution tracking', () => { + it('should track running executions correctly', async () => { + jest.useFakeTimers() + const limiter = new TimeRateLimiter(3, 1000) + const mockFn = jest.fn().mockImplementation(async () => { + await new Promise((resolve) => setTimeout(resolve, 100)) + return 'result' + }) + + const op1 = limiter.exec(mockFn) + expect(limiter.executions.length).toBe(1) + expect(limiter.executions[0].running).toBe(true) + + const op2 = limiter.exec(mockFn) + expect(limiter.executions.length).toBe(2) + + jest.advanceTimersByTime(101) + await Promise.resolve() + await Promise.resolve() + + await op1 + expect(limiter.executions[0].running).toBe(false) + + await op2 + jest.useRealTimers() + }) + + it('should mark executions as not running after completion', async () => { + const limiter = new TimeRateLimiter(2, 1000) + const mockFn = jest.fn().mockImplementation(async () => { + await new Promise((resolve) => setTimeout(resolve, 10)) + return 'result' + }) + + await limiter.exec(mockFn) + + const execution = limiter.executions[0] + expect(execution.running).toBe(false) + }) + + it('should mark executions as not running even on error', async () => { + const limiter = new TimeRateLimiter(2, 1000) + const mockFn = jest.fn().mockRejectedValue(new Error('test error')) + + await expect(limiter.exec(mockFn)).rejects.toThrow('test error') + + const execution = limiter.executions[0] + expect(execution.running).toBe(false) + }) + }) + + describe('cleanup behavior', () => { + it('should cleanup executions older than period', async () => { + jest.useFakeTimers() + const limiter = new TimeRateLimiter(3, 1000) + const mockFn = jest.fn().mockResolvedValue('result') + + await limiter.exec(mockFn) + expect(limiter.executions.length).toBe(1) + + jest.advanceTimersByTime(1100) + + await limiter.exec(mockFn) + // After cleanup, only the new execution should remain + expect(limiter.executions.length).toBe(1) + jest.useRealTimers() + }) + + it('should keep running executions regardless of time', async () => { + jest.useFakeTimers() + const limiter = new TimeRateLimiter(2, 1000) + let resolveOp: any + const mockFn = jest.fn().mockImplementation( + async () => + await new Promise((resolve) => { + resolveOp = resolve + }) + ) + + const op = limiter.exec(mockFn) + expect(limiter.executions.length).toBe(1) + + jest.advanceTimersByTime(2000) + + // Start another operation to trigger cleanup + const mockFn2 = jest.fn().mockResolvedValue('result') + await limiter.exec(mockFn2) + + // The first operation should still be tracked because it's running + const runningExecution = limiter.executions.find((e) => e.running) + expect(runningExecution).toBeDefined() + + resolveOp('done') + await op + jest.useRealTimers() + }) + }) + + describe('rate limiting behavior', () => { + it('should allow executions up to rate within period', async () => { + jest.useFakeTimers() + const limiter = new TimeRateLimiter(3, 1000) + const mockFn = jest.fn().mockResolvedValue('result') + + const ops = [limiter.exec(mockFn), limiter.exec(mockFn), limiter.exec(mockFn)] + + expect(mockFn).toHaveBeenCalledTimes(3) + await Promise.all(ops) + jest.useRealTimers() + }) + + it('should delay 4th execution when rate is 3', async () => { + jest.useFakeTimers() + const limiter = new TimeRateLimiter(3, 1000) + const mockFn = jest.fn().mockResolvedValue('result') + + void limiter.exec(mockFn) + void limiter.exec(mockFn) + void limiter.exec(mockFn) + const fourthOp = limiter.exec(mockFn) + + expect(mockFn).toHaveBeenCalledTimes(3) + + jest.advanceTimersByTime(1001) + await Promise.resolve() + await Promise.resolve() + + await fourthOp + expect(mockFn).toHaveBeenCalledTimes(4) + jest.useRealTimers() + }) + + it('should respect period for rate limiting', async () => { + jest.useFakeTimers() + const limiter = new TimeRateLimiter(2, 2000) // 2 per 2 seconds + const mockFn = jest.fn().mockResolvedValue('result') + + void limiter.exec(mockFn) + void limiter.exec(mockFn) + const thirdOp = limiter.exec(mockFn) + + expect(mockFn).toHaveBeenCalledTimes(2) + + // 1 second should not be enough + jest.advanceTimersByTime(1001) + await Promise.resolve() + expect(mockFn).toHaveBeenCalledTimes(2) + + // 2 seconds should allow the third + jest.advanceTimersByTime(1001) + await Promise.resolve() + await Promise.resolve() + + await thirdOp + expect(mockFn).toHaveBeenCalledTimes(3) + jest.useRealTimers() + }) + }) + + describe('error handling', () => { + it('should handle operation errors and continue', async () => { + const limiter = new TimeRateLimiter(2, 1000) + const errorFn = jest.fn().mockRejectedValue(new Error('operation error')) + const successFn = jest.fn().mockResolvedValue('success') + + await expect(limiter.exec(errorFn)).rejects.toThrow('operation error') + expect(limiter.active).toBe(0) + + const result = await limiter.exec(successFn) + expect(result).toBe('success') + }) + + it('should decrement active counter on error', async () => { + const limiter = new TimeRateLimiter(2, 1000) + const mockFn = jest.fn().mockRejectedValue(new Error('test error')) + + expect(limiter.active).toBe(0) + await expect(limiter.exec(mockFn)).rejects.toThrow('test error') + expect(limiter.active).toBe(0) + }) + + it('should handle synchronous throws', async () => { + const limiter = new TimeRateLimiter(2, 1000) + const mockFn = jest.fn().mockImplementation(() => { + throw new Error('sync error') + }) + + await expect(limiter.exec(mockFn)).rejects.toThrow('sync error') + expect(limiter.active).toBe(0) + }) + }) + + describe('waitProcessing', () => { + it('should resolve immediately when no operations are active', async () => { + const limiter = new TimeRateLimiter(2, 1000) + + await limiter.waitProcessing() + expect(limiter.active).toBe(0) + }) + + it('should wait for all active operations to complete', async () => { + jest.useFakeTimers() + const limiter = new TimeRateLimiter(2, 1000) + let completed = false + const mockFn = jest.fn().mockImplementation(async () => { + await new Promise((resolve) => setTimeout(resolve, 100)) + completed = true + return 'result' + }) + + void limiter.exec(mockFn) + expect(limiter.active).toBe(1) + + const waitPromise = limiter.waitProcessing() + + jest.advanceTimersByTime(101) + await Promise.resolve() + await Promise.resolve() + await Promise.resolve() + + await waitPromise + expect(completed).toBe(true) + expect(limiter.active).toBe(0) + jest.useRealTimers() + }) + }) + + describe('arguments passing', () => { + it('should pass arguments to operation', async () => { + const limiter = new TimeRateLimiter(2, 1000) + const mockFn = jest.fn().mockImplementation(async (args?: any) => { + return args?.value + }) + + const result = await limiter.exec(mockFn, { value: 42 }) + + expect(result).toBe(42) + expect(mockFn).toHaveBeenCalledWith({ value: 42 }) + }) + + it('should handle operations without arguments', async () => { + const limiter = new TimeRateLimiter(2, 1000) + const mockFn = jest.fn().mockResolvedValue('no-args') + + const result = await limiter.exec(mockFn) + + expect(result).toBe('no-args') + expect(mockFn).toHaveBeenCalledWith(undefined) + }) + }) + + describe('edge cases', () => { + it('should handle rate of 1 correctly', async () => { + jest.useFakeTimers() + const limiter = new TimeRateLimiter(1, 1000) + const mockFn = jest.fn().mockResolvedValue('result') + + void limiter.exec(mockFn) + expect(mockFn).toHaveBeenCalledTimes(1) + + const secondOp = limiter.exec(mockFn) + expect(mockFn).toHaveBeenCalledTimes(1) // Should wait + + jest.advanceTimersByTime(1001) + await Promise.resolve() + await Promise.resolve() + + await secondOp + expect(mockFn).toHaveBeenCalledTimes(2) + jest.useRealTimers() + }) + + it('should handle very short period', async () => { + jest.useFakeTimers() + const limiter = new TimeRateLimiter(2, 10) // 2 per 10ms + const mockFn = jest.fn().mockResolvedValue('result') + + void limiter.exec(mockFn) + void limiter.exec(mockFn) + const thirdOp = limiter.exec(mockFn) + + jest.advanceTimersByTime(11) + await Promise.resolve() + + await thirdOp + expect(mockFn).toHaveBeenCalledTimes(3) + jest.useRealTimers() + }) + + it('should handle large rate', async () => { + const limiter = new TimeRateLimiter(100, 1000) + const mockFn = jest.fn().mockResolvedValue('result') + + const operations = Array(50) + .fill(0) + .map(() => limiter.exec(mockFn)) + + await Promise.all(operations) + expect(mockFn).toHaveBeenCalledTimes(50) + }) + }) + + describe('stress test', () => { + it('should handle many operations correctly', async () => { + const limiter = new TimeRateLimiter(10, 100) + const results: number[] = [] + + const mockFn = jest.fn().mockImplementation(async (args?: any) => { + await new Promise((resolve) => setTimeout(resolve, 5)) + results.push(args?.id) + return args?.id + }) + + const operations = Array(30) + .fill(0) + .map((_, i) => limiter.exec(mockFn, { id: i })) + + await Promise.all(operations) + + expect(mockFn).toHaveBeenCalledTimes(30) + expect(results).toHaveLength(30) + expect(limiter.active).toBe(0) + }) }) }) diff --git a/packages/core/src/__tests__/rate-limiter.test.ts b/packages/core/src/__tests__/rate-limiter.test.ts new file mode 100644 index 0000000000..5ae22898a3 --- /dev/null +++ b/packages/core/src/__tests__/rate-limiter.test.ts @@ -0,0 +1,449 @@ +// +// Copyright © 2023 Hardcore Engineering Inc. +// +// Licensed under the Eclipse Public License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. You may +// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// +// See the License for the specific language governing permissions and +// limitations under the License. +// + +import { RateLimiter } from '../utils' + +describe('RateLimiter', () => { + describe('constructor', () => { + it('should create limiter with specified rate', () => { + const limiter = new RateLimiter(5) + expect(limiter.rate).toBe(5) + expect(limiter.idCounter).toBe(0) + expect(limiter.processingQueue.size).toBe(0) + }) + + it('should handle rate of 1', () => { + const limiter = new RateLimiter(1) + expect(limiter.rate).toBe(1) + }) + + it('should handle large rates', () => { + const limiter = new RateLimiter(1000) + expect(limiter.rate).toBe(1000) + }) + }) + + describe('exec', () => { + it('should execute single operation immediately', async () => { + const limiter = new RateLimiter(1) + const mockFn = jest.fn().mockResolvedValue('result') + + const result = await limiter.exec(mockFn) + + expect(mockFn).toHaveBeenCalledTimes(1) + expect(result).toBe('result') + expect(limiter.processingQueue.size).toBe(0) + }) + + it('should execute operations up to rate limit concurrently', async () => { + const limiter = new RateLimiter(3) + let runningCount = 0 + let maxRunning = 0 + + const mockFn = jest.fn().mockImplementation(async () => { + runningCount++ + maxRunning = Math.max(maxRunning, runningCount) + await new Promise((resolve) => setTimeout(resolve, 10)) + runningCount-- + return 'result' + }) + + const operations = [ + limiter.exec(mockFn), + limiter.exec(mockFn), + limiter.exec(mockFn), + limiter.exec(mockFn), + limiter.exec(mockFn) + ] + + await Promise.all(operations) + + expect(mockFn).toHaveBeenCalledTimes(5) + expect(maxRunning).toBeLessThanOrEqual(3) + }) + + it('should queue operations beyond rate limit', async () => { + const limiter = new RateLimiter(2) + const order: number[] = [] + + const createOperation = (id: number) => async () => { + order.push(id) + await new Promise((resolve) => setTimeout(resolve, 10)) + return id + } + + const operations = [ + limiter.exec(createOperation(1)), + limiter.exec(createOperation(2)), + limiter.exec(createOperation(3)), + limiter.exec(createOperation(4)) + ] + + await Promise.all(operations) + + expect(order).toEqual([1, 2, 3, 4]) + }) + + it('should handle operation errors correctly', async () => { + const limiter = new RateLimiter(2) + const mockFn = jest.fn().mockRejectedValue(new Error('Operation failed')) + + await expect(limiter.exec(mockFn)).rejects.toThrow('Operation failed') + expect(limiter.processingQueue.size).toBe(0) + }) + + it('should continue processing after error', async () => { + const limiter = new RateLimiter(1) + const successFn = jest.fn().mockResolvedValue('success') + const errorFn = jest.fn().mockRejectedValue(new Error('error')) + + await expect(limiter.exec(errorFn)).rejects.toThrow('error') + const result = await limiter.exec(successFn) + + expect(result).toBe('success') + expect(successFn).toHaveBeenCalledTimes(1) + }) + + it('should pass arguments to operations', async () => { + const limiter = new RateLimiter(1) + const mockFn = jest.fn().mockImplementation(async (args?: any) => { + return args?.value + }) + + const result = await limiter.exec(mockFn, { value: 42 }) + + expect(result).toBe(42) + expect(mockFn).toHaveBeenCalledWith({ value: 42 }) + }) + + it('should handle operations without arguments', async () => { + const limiter = new RateLimiter(1) + const mockFn = jest.fn().mockResolvedValue('no-args') + + const result = await limiter.exec(mockFn) + + expect(result).toBe('no-args') + expect(mockFn).toHaveBeenCalledWith(undefined) + }) + + it('should increment idCounter for each operation', async () => { + const limiter = new RateLimiter(10) + const mockFn = jest.fn().mockResolvedValue('result') + + await limiter.exec(mockFn) + expect(limiter.idCounter).toBe(1) + + await limiter.exec(mockFn) + expect(limiter.idCounter).toBe(2) + + await limiter.exec(mockFn) + expect(limiter.idCounter).toBe(3) + }) + + it('should notify waiting operations when slots become available', async () => { + const limiter = new RateLimiter(1) + const order: string[] = [] + + const op1 = async (): Promise => { + order.push('op1-start') + await new Promise((resolve) => setTimeout(resolve, 20)) + order.push('op1-end') + return 'op1' + } + + const op2 = async (): Promise => { + order.push('op2-start') + await new Promise((resolve) => setTimeout(resolve, 10)) + order.push('op2-end') + return 'op2' + } + + const results = await Promise.all([limiter.exec(op1), limiter.exec(op2)]) + + expect(results).toEqual(['op1', 'op2']) + expect(order).toEqual(['op1-start', 'op1-end', 'op2-start', 'op2-end']) + }) + + it('should handle rapid sequential operations', async () => { + const limiter = new RateLimiter(2) + const mockFn = jest.fn().mockResolvedValue('result') + const results = [] + + for (let i = 0; i < 10; i++) { + results.push(limiter.exec(mockFn)) + } + + await Promise.all(results) + expect(mockFn).toHaveBeenCalledTimes(10) + }) + + it('should clean up processingQueue after completion', async () => { + const limiter = new RateLimiter(5) + const mockFn = jest.fn().mockImplementation(async () => { + await new Promise((resolve) => setTimeout(resolve, 5)) + return 'result' + }) + + const operations = Array(5) + .fill(0) + .map(() => limiter.exec(mockFn)) + + await Promise.all(operations) + expect(limiter.processingQueue.size).toBe(0) + }) + }) + + describe('add', () => { + it('should add operation without waiting for result', async () => { + const limiter = new RateLimiter(1) + let executed = false + const mockFn = jest.fn().mockImplementation(async () => { + await new Promise((resolve) => setTimeout(resolve, 10)) + executed = true + return 'result' + }) + + await limiter.add(mockFn) + + // add returns immediately, but operation may not be complete + expect(mockFn).toHaveBeenCalledTimes(1) + + // Wait for operation to complete + await new Promise((resolve) => setTimeout(resolve, 20)) + expect(executed).toBe(true) + }) + + it('should handle errors with error handler', async () => { + const limiter = new RateLimiter(1) + const errorHandler = jest.fn() + const mockFn = jest.fn().mockRejectedValue(new Error('test error')) + + await limiter.add(mockFn, undefined, errorHandler) + + // Wait for operation to execute + await new Promise((resolve) => setTimeout(resolve, 10)) + + expect(errorHandler).toHaveBeenCalledWith(new Error('test error')) + }) + + it('should log errors when no error handler provided', async () => { + const limiter = new RateLimiter(1) + const consoleSpy = jest.spyOn(console, 'error').mockImplementation() + const mockFn = jest.fn().mockRejectedValue(new Error('test error')) + + await limiter.add(mockFn) + + // Wait for operation to execute + await new Promise((resolve) => setTimeout(resolve, 10)) + + expect(consoleSpy).toHaveBeenCalled() + consoleSpy.mockRestore() + }) + + it('should pass arguments to operation', async () => { + const limiter = new RateLimiter(1) + const mockFn = jest.fn().mockResolvedValue('result') + + await limiter.add(mockFn, { test: 'value' }) + + // Wait for operation to start + await new Promise((resolve) => setTimeout(resolve, 10)) + + expect(mockFn).toHaveBeenCalledWith({ test: 'value' }) + }) + + it('should queue multiple add operations', async () => { + const limiter = new RateLimiter(1) + const order: number[] = [] + + const createOp = (id: number) => async () => { + order.push(id) + await new Promise((resolve) => setTimeout(resolve, 10)) + } + + await limiter.add(createOp(1)) + await limiter.add(createOp(2)) + await limiter.add(createOp(3)) + + // Wait for all operations to complete + await limiter.waitProcessing() + + expect(order).toEqual([1, 2, 3]) + }) + }) + + describe('waitProcessing', () => { + it('should wait for all operations to complete', async () => { + const limiter = new RateLimiter(2) + const mockFn = jest.fn().mockImplementation(async () => { + await new Promise((resolve) => setTimeout(resolve, 20)) + return 'result' + }) + + const operations = [limiter.exec(mockFn), limiter.exec(mockFn), limiter.exec(mockFn)] + + expect(limiter.processingQueue.size).toBeGreaterThan(0) + + await limiter.waitProcessing() + + expect(limiter.processingQueue.size).toBe(0) + await Promise.all(operations) + }) + + it('should resolve immediately when no operations are processing', async () => { + const limiter = new RateLimiter(1) + + await limiter.waitProcessing() + + expect(limiter.processingQueue.size).toBe(0) + }) + + it('should wait for operations added via add method', async () => { + const limiter = new RateLimiter(1) + let completed = false + const mockFn = jest.fn().mockImplementation(async () => { + await new Promise((resolve) => setTimeout(resolve, 20)) + completed = true + }) + + await limiter.add(mockFn) + expect(completed).toBe(false) + + await limiter.waitProcessing() + expect(completed).toBe(true) + }) + }) + + describe('edge cases', () => { + it('should handle zero rate (should not happen but test defensively)', async () => { + const limiter = new RateLimiter(0) + const mockFn = jest.fn().mockResolvedValue('result') + + // This will hang forever, so we use Promise.race with timeout + const timeoutPromise = new Promise((resolve) => { + setTimeout(() => { + resolve('timeout') + }, 100) + }) + const result = await Promise.race([limiter.exec(mockFn), timeoutPromise]) + + expect(result).toBe('timeout') + }) + + it('should handle very long running operations', async () => { + const limiter = new RateLimiter(1) + const mockFn = jest.fn().mockImplementation(async () => { + await new Promise((resolve) => setTimeout(resolve, 50)) + return 'result' + }) + + const start = Date.now() + await limiter.exec(mockFn) + const duration = Date.now() - start + + expect(duration).toBeGreaterThanOrEqual(50) + }) + + it('should handle mixed sync and async operations', async () => { + const limiter = new RateLimiter(2) + const syncOp = jest.fn().mockResolvedValue('sync') + const asyncOp = jest.fn().mockImplementation(async () => { + await new Promise((resolve) => setTimeout(resolve, 10)) + return 'async' + }) + + const results = await Promise.all([limiter.exec(syncOp), limiter.exec(asyncOp)]) + + expect(results).toContain('sync') + expect(results).toContain('async') + }) + + it('should handle operations that throw synchronously', async () => { + const limiter = new RateLimiter(1) + const mockFn = jest.fn().mockImplementation(() => { + throw new Error('Sync error') + }) + + await expect(limiter.exec(mockFn)).rejects.toThrow('Sync error') + expect(limiter.processingQueue.size).toBe(0) + }) + + it('should handle notify array correctly when multiple operations wait', async () => { + const limiter = new RateLimiter(1) + const order: number[] = [] + + const createOp = (id: number) => async () => { + order.push(id) + await new Promise((resolve) => setTimeout(resolve, 10)) + return id + } + + // Start multiple operations that will need to wait + const operations = [ + limiter.exec(createOp(1)), + limiter.exec(createOp(2)), + limiter.exec(createOp(3)), + limiter.exec(createOp(4)) + ] + + await Promise.all(operations) + + expect(order).toEqual([1, 2, 3, 4]) + expect(limiter.notify.length).toBe(0) + }) + }) + + describe('concurrent stress test', () => { + it('should handle many concurrent operations correctly', async () => { + const limiter = new RateLimiter(5) + const mockFn = jest.fn().mockImplementation(async (args?: any) => { + await new Promise((resolve) => setTimeout(resolve, Math.random() * 10)) + return args?.id + }) + + const operations = Array(50) + .fill(0) + .map((_, i) => limiter.exec(mockFn, { id: i })) + + const results = await Promise.all(operations) + + expect(mockFn).toHaveBeenCalledTimes(50) + expect(results).toHaveLength(50) + expect(limiter.processingQueue.size).toBe(0) + }) + + it('should maintain rate limit under heavy load', async () => { + const limiter = new RateLimiter(3) + let currentCount = 0 + let maxConcurrent = 0 + + const mockFn = jest.fn().mockImplementation(async () => { + currentCount++ + maxConcurrent = Math.max(maxConcurrent, currentCount) + await new Promise((resolve) => setTimeout(resolve, 10)) + currentCount-- + }) + + const operations = Array(20) + .fill(0) + .map(() => limiter.exec(mockFn)) + + await Promise.all(operations) + + expect(maxConcurrent).toBeLessThanOrEqual(3) + expect(currentCount).toBe(0) + }) + }) +}) diff --git a/packages/core/src/utils.ts b/packages/core/src/utils.ts index 245cbd998b..905a460cc0 100644 --- a/packages/core/src/utils.ts +++ b/packages/core/src/utils.ts @@ -904,10 +904,10 @@ export class TimeRateLimiter { } const v = { time: Date.now(), running: true } + this.active++ try { this.executions.push(v) const p = op(args) - this.active++ return await p } finally { v.running = false @@ -922,7 +922,6 @@ export class TimeRateLimiter { async waitProcessing (): Promise { while (this.active > 0) { - console.log('wait', this.active) await new Promise((resolve) => { this.notify.push(resolve) }) diff --git a/rush.json b/rush.json index 0d9c6549bd..2115992aad 100644 --- a/rush.json +++ b/rush.json @@ -213,7 +213,7 @@ * your PR branch, and in this situation "rush change" will also automatically invoke "git fetch" * to retrieve the latest activity for the remote main branch. */ - "url": "https://github.com/hcengineering/huly.net", + "url": "https://github.com/hcengineering/huly.core", /** * The default branch name. This tells "rush change" which remote branch to compare against. * The default value is "main"