From 5d9d5d657e76518214a290b2194865c0a8ab5766 Mon Sep 17 00:00:00 2001 From: Nicholas Carter Date: Sun, 2 Aug 2026 11:13:19 +0000 Subject: [PATCH] test: cover issue 8 forwarder residuals --- src/middleware/forwarder.ts | 19 ++- tests/middleware/forwarder.test.ts | 191 +++++++++++++++++++++++++++++ 2 files changed, 208 insertions(+), 2 deletions(-) diff --git a/src/middleware/forwarder.ts b/src/middleware/forwarder.ts index e1dadcf..d693236 100644 --- a/src/middleware/forwarder.ts +++ b/src/middleware/forwarder.ts @@ -484,10 +484,25 @@ export class Forwarder { // Check for retryable status if (isRetryableStatus(response.status)) { const retryAfter = response.headers.get('retry-after'); - const delay = retryAfter - ? parseInt(retryAfter, 10) * 1000 + const retryAfterSeconds = retryAfter ? Number.parseInt(retryAfter, 10) : NaN; + const delay = Number.isFinite(retryAfterSeconds) + ? retryAfterSeconds * 1000 : calculateRetryDelay(attempt, this.config.retryDelayMs); + const error: UpstreamError = new Error( + `Upstream error: ${response.status} ${response.statusText}` + ); + error.statusCode = response.status; + error.upstreamStatus = response.status; + error.retryable = true; + error.details = await response.text().catch(() => 'Unknown error'); + lastError = error; + + if (attempt >= this.config.maxRetries) { + this.handleError(error, req, res); + return; + } + console.warn( `[Forwarder] Model ${model} returned ${response.status}. Retrying in ${delay}ms (attempt ${attempt + 1}/${this.config.maxRetries})...` ); diff --git a/tests/middleware/forwarder.test.ts b/tests/middleware/forwarder.test.ts index b94bae3..300cd69 100644 --- a/tests/middleware/forwarder.test.ts +++ b/tests/middleware/forwarder.test.ts @@ -1,4 +1,5 @@ import { describe, it, expect, beforeEach, afterEach, jest } from '@jest/globals'; +import { EventEmitter } from 'events'; import { Request, Response, NextFunction } from 'express'; import request from 'supertest'; import express, { Application } from 'express'; @@ -379,6 +380,196 @@ describe('Forwarder Middleware', () => { expect(response.body.error).toContain('aborted'); }); + it('should retry on 504 status', async () => { + mockFetch + .mockResolvedValueOnce({ + ok: false, + status: 504, + statusText: 'Gateway Timeout', + headers: new Headers({ 'retry-after': '0' }), + text: async () => 'Gateway timeout', + }) + .mockResolvedValueOnce({ + ok: true, + status: 200, + headers: new Headers({ 'content-type': 'application/json' }), + json: async () => ({ + id: 'chatcmpl-504', + object: 'chat.completion', + created: Date.now(), + model: 'gpt-4', + choices: [ + { index: 0, message: { role: 'assistant', content: 'ok' }, finish_reason: 'stop' }, + ], + }), + }); + + app.post('/chat/completions', createForwarder({ baseUrl: 'http://localhost:20132/v1' })); + + const response = await request(app) + .post('/chat/completions') + .send({ model: 'gpt-4', messages: [{ role: 'user', content: 'Hello' }] }) + .expect(200); + + expect(response.body.id).toBe('chatcmpl-504'); + expect(mockFetch).toHaveBeenCalledTimes(2); + }); + + it('should return last retryable error after exhausting retries', async () => { + mockFetch.mockResolvedValue({ + ok: false, + status: 504, + statusText: 'Gateway Timeout', + headers: new Headers({ 'retry-after': '0' }), + text: async () => 'Gateway timeout', + }); + + app.post( + '/chat/completions', + createForwarder({ baseUrl: 'http://localhost:20132/v1', maxRetries: 1 }) + ); + + const response = await request(app) + .post('/chat/completions') + .send({ model: 'gpt-4', messages: [{ role: 'user', content: 'Hello' }] }) + .expect(504); + + expect(response.body).toMatchObject({ + upstreamStatus: 504, + retryable: true, + details: 'Gateway timeout', + }); + expect(mockFetch).toHaveBeenCalledTimes(2); + }); + + it('should ignore invalid Retry-After and use configured backoff', async () => { + const setTimeoutSpy = jest.spyOn(global, 'setTimeout'); + mockFetch + .mockResolvedValueOnce({ + ok: false, + status: 503, + statusText: 'Service Unavailable', + headers: new Headers({ 'retry-after': 'soon' }), + text: async () => 'try later', + }) + .mockResolvedValueOnce({ + ok: true, + status: 200, + headers: new Headers({ 'content-type': 'application/json' }), + json: async () => ({ + id: 'chatcmpl-backoff', + object: 'chat.completion', + created: Date.now(), + model: 'gpt-4', + choices: [ + { index: 0, message: { role: 'assistant', content: 'ok' }, finish_reason: 'stop' }, + ], + }), + }); + + app.post( + '/chat/completions', + createForwarder({ baseUrl: 'http://localhost:20132/v1', retryDelayMs: 0 }) + ); + + await request(app) + .post('/chat/completions') + .send({ model: 'gpt-4', messages: [{ role: 'user', content: 'Hello' }] }) + .expect(200); + + expect(setTimeoutSpy).toHaveBeenCalledWith(expect.any(Function), 0); + setTimeoutSpy.mockRestore(); + }); + + it('should abort upstream fetch when the client disconnects', async () => { + mockFetch.mockImplementationOnce( + (_url, options) => + new Promise((_resolve, reject) => { + options.signal.addEventListener('abort', () => { + const error = new Error('Aborted') as Error & { name: string }; + error.name = 'AbortError'; + reject(error); + }); + }) + ); + + const middleware = createForwarder({ baseUrl: 'http://localhost:20132/v1', timeoutMs: 1000 }); + const req = new EventEmitter() as Request; + req.body = { model: 'gpt-4', messages: [{ role: 'user', content: 'Hello' }] }; + req.headers = {}; + const res = { + status: jest.fn().mockReturnThis(), + json: jest.fn().mockReturnThis(), + } as unknown as Response; + + const done = middleware(req, res, jest.fn() as NextFunction); + req.emit('close'); + await done; + + expect(res.status).toHaveBeenCalledWith(408); + expect(res.json).toHaveBeenCalledWith(expect.objectContaining({ retryable: false })); + }); + + it('should forward partial and mixed SSE chunks', async () => { + const chunks = [ + 'event: keepalive\n\n', + 'data: {"id":"chatcmpl-123","object":"chat.completion.chunk","created":123,"model":"gpt-4","choices":[{"index":0,"delta":{"content":"Hel"},"finish_reason":null}]}\n', + '\ndata: [DONE]\n\n', + ]; + let chunkIndex = 0; + + mockFetch.mockResolvedValueOnce({ + ok: true, + status: 200, + headers: new Headers({ 'content-type': 'text/event-stream' }), + body: { + getReader: () => ({ + read: async () => + chunkIndex < chunks.length + ? { done: false, value: new TextEncoder().encode(chunks[chunkIndex++]) } + : { done: true }, + }), + }, + }); + + app.post('/chat/completions', createForwarder({ baseUrl: 'http://localhost:20132/v1' })); + + const response = await request(app) + .post('/chat/completions') + .send({ model: 'gpt-4', messages: [{ role: 'user', content: 'Hello' }], stream: true }) + .expect(200); + + expect(response.text).toContain('event: keepalive'); + expect(response.text).toContain('data: [DONE]'); + }); + + it('should call onError for upstream failures', async () => { + const onError = jest.fn(); + mockFetch.mockResolvedValueOnce({ + ok: false, + status: 500, + statusText: 'Internal Server Error', + headers: new Headers({}), + text: async () => 'boom', + }); + + app.post( + '/chat/completions', + createForwarder({ baseUrl: 'http://localhost:20132/v1', onError }) + ); + + await request(app) + .post('/chat/completions') + .send({ model: 'gpt-4', messages: [{ role: 'user', content: 'Hello' }] }) + .expect(500); + + expect(onError).toHaveBeenCalledWith( + expect.objectContaining({ upstreamStatus: 500, details: 'boom' }), + expect.anything(), + expect.anything() + ); + }); + it('should track usage when enabled', async () => { const mockResponse = { id: 'chatcmpl-123',