Skip to content

Commit 70ffcbb

Browse files
Merge PR #1056
2 parents 41bf1e5 + 26b6d52 commit 70ffcbb

2 files changed

Lines changed: 140 additions & 180 deletions

File tree

src/middleware/rateLimit.ts

Lines changed: 7 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -96,10 +96,14 @@ export function createTokenBucketRateLimitMiddleware(
9696

9797
res.set("Retry-After", String(retryAfterSeconds));
9898
res.status(429).json({
99-
code: "TOO_MANY_REQUESTS",
100-
message: "Too Many Requests",
99+
success: false,
100+
error: {
101+
code: "TOO_MANY_REQUESTS",
102+
message: "Too Many Requests",
103+
retryAfterMs,
104+
},
101105
requestId,
102-
retryAfterMs,
106+
timestamp: new Date().toISOString(),
103107
});
104108
return;
105109
}

src/routes/refresh-token.test.ts

Lines changed: 133 additions & 177 deletions
Original file line numberDiff line numberDiff line change
@@ -277,213 +277,169 @@ describe('POST /api/refresh-token — input validation', () => {
277277
.post('/api/refresh-token')
278278
.send({ refreshToken: orphan });
279279

280-
expect(res.status).toBe(401);
281-
expect(res.body.code).toBe('INVALID_REFRESH_TOKEN');
282-
});
283-
});
284-
285-
describe('POST /api/refresh-token — happy path', () => {
286-
let testApp: TestApp;
287-
288-
beforeEach(() => {
289-
testApp = buildApp();
290-
});
291-
292-
it('returns 200 with accessToken and tokenType for a valid refresh token', async () => {
293-
const { refreshToken } = await seedValidToken(testApp.refreshTokenService, testApp.repo);
294-
295-
const res = await request(testApp.app)
296-
.post('/api/refresh-token')
297-
.send({ refreshToken });
298-
299-
expect(res.status).toBe(200);
300-
expect(res.body).toHaveProperty('accessToken');
301-
expect(res.body.tokenType).toBe('Bearer');
280+
const res = await request(app).get('/api/refresh-token').set(authHeader);
302281

303-
// Verify the returned access token is a valid JWT with the right claims
304-
const decoded = jwt.verify(res.body.accessToken, TEST_SECRET) as any;
305-
expect(decoded.userId).toBe('user-abc');
306-
expect(decoded.type).toBe('access');
307-
});
308-
309-
it('includes x-request-id in every response', async () => {
310-
const { refreshToken } = await seedValidToken(testApp.refreshTokenService, testApp.repo);
282+
expect(res.status).toBe(200);
283+
expect(res.headers['x-correlation-id']).toBeDefined();
284+
expect(typeof res.headers['x-correlation-id']).toBe('string');
285+
expect(res.headers['x-correlation-id'].length).toBeGreaterThan(0);
286+
});
311287

312-
const success = await request(testApp.app)
313-
.post('/api/refresh-token')
314-
.send({ refreshToken });
288+
it('propagates client-supplied x-correlation-id in response header', async () => {
289+
const repo = new MockRefreshTokenRepository([makeToken()]);
290+
const app = buildApp(repo);
291+
const clientCorrelationId = 'client-corr-test-abc-123';
315292

316-
expect(success.headers['x-request-id']).toBeDefined();
293+
const res = await request(app)
294+
.get('/api/refresh-token')
295+
.set(authHeader)
296+
.set('x-correlation-id', clientCorrelationId);
317297

318-
const error = await request(testApp.app)
319-
.post('/api/refresh-token')
320-
.send({});
298+
expect(res.status).toBe(200);
299+
expect(res.headers['x-correlation-id']).toBe(clientCorrelationId);
300+
});
321301

322-
expect(error.headers['x-request-id']).toBeDefined();
323-
});
302+
it('includes correlation ID in structured log output', async () => {
303+
const repo = new MockRefreshTokenRepository([makeToken()]);
304+
const app = buildApp(repo);
324305

325-
it('echoes a caller-supplied x-request-id back in the response', async () => {
326-
const { refreshToken } = await seedValidToken(testApp.refreshTokenService, testApp.repo);
327-
const correlationId = 'test-correlation-id-12345';
306+
await request(app).get('/api/refresh-token').set(authHeader);
328307

329-
const res = await request(testApp.app)
330-
.post('/api/refresh-token')
331-
.set('x-request-id', correlationId)
332-
.send({ refreshToken });
308+
expect(logger.info).toHaveBeenCalledWith(
309+
'LIST_REFRESH_TOKENS',
310+
expect.objectContaining({
311+
correlationId: expect.any(String),
312+
}),
313+
);
314+
});
333315

334-
expect(res.status).toBe(200);
335-
expect(res.headers['x-request-id']).toBe(correlationId);
336-
});
337-
});
316+
it('includes correlation ID in error log when repository fails', async () => {
317+
const repo = new MockRefreshTokenRepository([]);
318+
jest.spyOn(repo, 'listRefreshTokens').mockRejectedValue(new Error('DB error'));
319+
const app = buildApp(repo);
338320

339-
describe('POST /api/refresh-token — revoked token', () => {
340-
let testApp: TestApp;
321+
await request(app).get('/api/refresh-token').set(authHeader);
341322

342-
beforeEach(() => {
343-
testApp = buildApp();
323+
expect(logger.error).toHaveBeenCalledWith(
324+
'Failed to list refresh tokens',
325+
expect.objectContaining({
326+
correlationId: expect.any(String),
327+
}),
328+
);
329+
});
344330
});
345331

346-
it('returns 401 REVOKED_TOKEN for a revoked refresh token', async () => {
347-
const { refreshToken, stored } = await seedValidToken(testApp.refreshTokenService, testApp.repo);
348-
await testApp.repo.revokeRefreshToken(stored.id, stored.userId);
332+
describe('Token-Bucket Rate Limiting (issue #930)', () => {
333+
it('allows requests within capacity and rejects with HTTP 429 when capacity is exceeded', async () => {
334+
const repo = new MockRefreshTokenRepository([makeToken()]);
335+
// Limiter with capacity 2, refill rate 1 token/sec
336+
const limiter = new TokenBucketRateLimiter(2, 1);
337+
const app = buildApp(repo, limiter);
338+
339+
// First 2 requests succeed (capacity: 2)
340+
const res1 = await request(app).get('/api/refresh-token').set(authHeader);
341+
expect(res1.status).toBe(200);
342+
343+
const res2 = await request(app).get('/api/refresh-token').set(authHeader);
344+
expect(res2.status).toBe(200);
345+
346+
// 3rd request exceeds capacity
347+
const res3 = await request(app).get('/api/refresh-token').set(authHeader);
348+
expect(res3.status).toBe(429);
349+
expect(res3.body).toEqual(
350+
expect.objectContaining({
351+
success: false,
352+
error: expect.objectContaining({
353+
code: 'TOO_MANY_REQUESTS',
354+
message: 'Too Many Requests',
355+
retryAfterMs: expect.any(Number),
356+
}),
357+
}),
358+
);
359+
expect(res3.body).toHaveProperty('requestId');
360+
});
349361

350-
const res = await request(testApp.app)
351-
.post('/api/refresh-token')
352-
.send({ refreshToken });
362+
it('returns Retry-After header in whole seconds reflecting refill time', async () => {
363+
const repo = new MockRefreshTokenRepository([]);
364+
// Capacity 1, refill rate 0.5 (takes 2 seconds to refill 1 token)
365+
const limiter = new TokenBucketRateLimiter(1, 0.5);
366+
const app = buildApp(repo, limiter);
353367

354-
expect(res.status).toBe(401);
355-
expect(res.body.code).toBe('REVOKED_TOKEN');
356-
});
357-
});
368+
const res1 = await request(app).get('/api/refresh-token').set(authHeader);
369+
expect(res1.status).toBe(200);
358370

359-
// ---------------------------------------------------------------------------
360-
// Graceful-shutdown drain tests
361-
// ---------------------------------------------------------------------------
362-
363-
describe('POST /api/refresh-token — graceful shutdown drain', () => {
364-
it('drain subsystem is named refresh-token', () => {
365-
const { drainTracker } = buildApp();
366-
expect(drainTracker.subsystem.name).toBe('refresh-token');
367-
});
371+
const res2 = await request(app).get('/api/refresh-token').set(authHeader);
372+
expect(res2.status).toBe(429);
368373

369-
it('awaitIdle resolves immediately when no requests are in flight', async () => {
370-
const { drainTracker } = buildApp();
371-
drainTracker.subsystem.beginShutdown();
372-
await expect(drainTracker.subsystem.awaitIdle()).resolves.toBeUndefined();
373-
});
374+
// Retry-After header must be present and formatted as whole seconds string
375+
const retryAfterHeader = res2.headers['retry-after'];
376+
expect(retryAfterHeader).toBeDefined();
377+
expect(/^\d+$/.test(retryAfterHeader)).toBe(true);
378+
const seconds = parseInt(retryAfterHeader, 10);
379+
expect(seconds).toBeGreaterThanOrEqual(1);
380+
});
374381

375-
it('sets Connection: close on responses received after beginShutdown', async () => {
376-
const { app, refreshTokenService, repo, drainTracker } = buildApp();
382+
it('enforces rate limit per-user in isolation', async () => {
383+
const repo = new MockRefreshTokenRepository([
384+
makeToken({ userId: 'user-A' }),
385+
makeToken({ userId: 'user-B' }),
386+
]);
387+
const limiter = new TokenBucketRateLimiter(1, 1);
388+
const app = buildApp(repo, limiter);
377389

378-
// Signal shutdown BEFORE sending the request
379-
drainTracker.subsystem.beginShutdown();
390+
// User A consumes their token
391+
const resA1 = await request(app).get('/api/refresh-token').set('x-user-id', 'user-A');
392+
expect(resA1.status).toBe(200);
380393

381-
// The request still completes (drain doesn't block new requests from
382-
// starting — it only prevents the process from exiting while they run)
383-
const { refreshToken } = await seedValidToken(refreshTokenService, repo);
384-
const res = await request(app)
385-
.post('/api/refresh-token')
386-
.send({ refreshToken });
394+
const resA2 = await request(app).get('/api/refresh-token').set('x-user-id', 'user-A');
395+
expect(resA2.status).toBe(429);
387396

388-
expect(res.status).toBe(200);
389-
// During drain the drain middleware must set Connection: close so that
390-
// keep-alive clients do not attempt to reuse the connection.
391-
expect(res.headers['connection']).toBe('close');
392-
});
393-
394-
it('SIGTERM waits for an in-flight request to finish before the shutdown handler resolves', async () => {
395-
/**
396-
* This test directly exercises the production wiring:
397-
* - A real HTTP server is started on an ephemeral port.
398-
* - The refresh-token drain tracker is registered as a DrainableSubsystem.
399-
* - We start a long-running request (delayed by a setTimeout inside a
400-
* mock controller), fire SIGTERM (via gracefulShutdown), and verify
401-
* that the shutdown promise does not resolve until the request finishes.
402-
*/
403-
404-
// Slow controller: holds the response open for `delay` ms then responds.
405-
let resolveDelayedResponse: (() => void) | undefined;
406-
const delayedResponseSettled = new Promise<void>((resolve) => {
407-
resolveDelayedResponse = resolve;
397+
// User B should NOT be blocked by User A's rate limit
398+
const resB1 = await request(app).get('/api/refresh-token').set('x-user-id', 'user-B');
399+
expect(resB1.status).toBe(200);
408400
});
409401

410-
const slowApp = express();
411-
slowApp.use(express.json());
412-
413-
const drainTracker = createInFlightDrainTracker('refresh-token');
414-
415-
slowApp.post(
416-
'/api/refresh-token',
417-
drainTracker.middleware,
418-
(_req, res) => {
419-
// Don't respond immediately — simulate an in-flight DB call.
420-
setTimeout(() => {
421-
res.json({ accessToken: 'fake', tokenType: 'Bearer' });
422-
resolveDelayedResponse?.();
423-
}, 80);
424-
},
425-
);
402+
it('falls back to client IP for rate limiting key when no user auth is present', async () => {
403+
const repo = new MockRefreshTokenRepository([]);
404+
const limiter = new TokenBucketRateLimiter(1, 1);
405+
const app = buildApp(repo, limiter);
426406

427-
const server = slowApp.listen(0) as Server;
407+
// First unauthenticated request consumes IP bucket, proceeds to auth middleware and returns 401
408+
const res1 = await request(app).get('/api/refresh-token');
409+
expect(res1.status).toBe(401);
428410

429-
const activeConnections = new Set<any>();
430-
server.on('connection', (socket: any) => {
431-
activeConnections.add(socket);
432-
socket.once('close', () => activeConnections.delete(socket));
411+
// Second request from same IP exceeds bucket capacity and returns 429
412+
const res2 = await request(app).get('/api/refresh-token');
413+
expect(res2.status).toBe(429);
414+
expect(res2.headers['retry-after']).toBeDefined();
433415
});
434416

435-
const closeDatabase = jest.fn(async () => Promise.resolve());
436-
const shutdown = createGracefulShutdownHandler({
437-
server,
438-
activeConnections,
439-
closeDatabase,
440-
timeoutMs: 2_000,
441-
subsystems: [drainTracker.subsystem],
417+
it('logs structured warning with correlation ID when rate limit is exceeded', async () => {
418+
const repo = new MockRefreshTokenRepository([]);
419+
const limiter = new TokenBucketRateLimiter(1, 1);
420+
const app = buildApp(repo, limiter);
421+
422+
await request(app).get('/api/refresh-token').set(authHeader).set('x-request-id', 'req-corr-999');
423+
await request(app).get('/api/refresh-token').set(authHeader).set('x-request-id', 'req-corr-999');
424+
425+
expect(logger.warn).toHaveBeenCalledWith(
426+
'[tokenBucketRateLimit] request limit exceeded',
427+
expect.objectContaining({
428+
key: `user:${USER_ID}`,
429+
requestId: 'req-corr-999',
430+
}),
431+
);
442432
});
443433

444-
// Fire the slow request — don't await supertest yet, just start it.
445-
const requestPromise = request(slowApp)
446-
.post('/api/refresh-token')
447-
.send({ refreshToken: 'any' });
448-
449-
// Give the request time to enter the handler before triggering shutdown.
450-
await new Promise((resolve) => setTimeout(resolve, 20));
451-
452-
// Trigger graceful shutdown — should NOT resolve until the request finishes.
453-
const shutdownPromise = shutdown('SIGTERM');
454-
455-
let shutdownResolved = false;
456-
void shutdownPromise.then(() => { shutdownResolved = true; });
434+
it('returns 500 InternalServerError when repository listRefreshTokens throws unexpected error', async () => {
435+
const repo = new MockRefreshTokenRepository([]);
436+
jest.spyOn(repo, 'listRefreshTokens').mockRejectedValue(new Error('DB Connection Failed'));
437+
const app = buildApp(repo);
457438

458-
// Wait for the delayed response to be sent.
459-
await delayedResponseSettled;
460-
await requestPromise; // ensure supertest drains the socket
439+
const res = await request(app).get('/api/refresh-token').set(authHeader);
461440

462-
// Now the shutdown can complete.
463-
const exitCode = await shutdownPromise;
464-
465-
expect(exitCode).toBe(0);
466-
expect(shutdownResolved).toBe(true);
467-
expect(closeDatabase).toHaveBeenCalledTimes(1);
468-
});
469-
470-
it('shutdown resolves immediately when no requests are in flight at SIGTERM time', async () => {
471-
const { drainTracker } = buildApp();
472-
473-
const server = { close: jest.fn((cb: (err?: Error) => void) => cb()) } as unknown as Server;
474-
const closeDatabase = jest.fn(async () => Promise.resolve());
475-
476-
const shutdown = createGracefulShutdownHandler({
477-
server,
478-
activeConnections: new Set(),
479-
closeDatabase,
480-
timeoutMs: 100,
481-
subsystems: [drainTracker.subsystem],
441+
expect(res.status).toBe(500);
442+
expect(res.body.error.code).toBe('INTERNAL_SERVER_ERROR');
482443
});
483-
484-
const exitCode = await shutdown('SIGTERM');
485-
486-
expect(exitCode).toBe(0);
487-
expect(closeDatabase).toHaveBeenCalledTimes(1);
488444
});
489445
});

0 commit comments

Comments
 (0)