Fix rate limits bug

This commit is contained in:
Andrey Sobolev
2025-10-08 10:38:20 +07:00
parent a09ee082ed
commit c535ee63fd
7 changed files with 1226 additions and 44 deletions
@@ -0,0 +1,10 @@
{
"changes": [
{
"packageName": "@hcengineering/core",
"comment": "Fix issue with time rate limiter",
"type": "patch"
}
],
"packageName": "@hcengineering/core"
}
+41 -41
View File
@@ -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
}
]
@@ -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<string> => {
await new Promise((resolve) => setTimeout(resolve, 10))
results.push(id)
return id
}
const addOp = async (id: string): Promise<void> => {
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<number> => {
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<string | number> = []
const stringOp = async (val: string): Promise<string> => {
await new Promise((resolve) => setTimeout(resolve, 5))
results.push(val)
return val
}
const numberOp = async (val: number): Promise<number> => {
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<number> => {
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)
})
})
})
+344
View File
@@ -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)
})
})
})
@@ -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<string> => {
order.push('op1-start')
await new Promise((resolve) => setTimeout(resolve, 20))
order.push('op1-end')
return 'op1'
}
const op2 = async (): Promise<string> => {
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)
})
})
})
+1 -2
View File
@@ -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<void> {
while (this.active > 0) {
console.log('wait', this.active)
await new Promise<void>((resolve) => {
this.notify.push(resolve)
})
+1 -1
View File
@@ -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"