deepseek-responses-bridge.ts 6.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190
  1. import { createServer } from 'node:http'
  2. import type {
  3. IncomingMessage,
  4. Server,
  5. ServerResponse,
  6. } from 'node:http'
  7. import { completeResponsesEvents } from './responses-fixture.ts'
  8. const OFFICIAL_DEEPSEEK_BASE_URL = 'https://api.deepseek.com'
  9. const MAX_REQUEST_BYTES = 1_048_576
  10. /** One running test-only Responses-to-DeepSeek bridge. */
  11. export interface DeepSeekResponsesBridge {
  12. readonly baseUrl: string
  13. readonly completedRequests: number
  14. close(): Promise<void>
  15. }
  16. function readRequest(request: IncomingMessage): Promise<string> {
  17. return new Promise((resolve, reject) => {
  18. let body = ''
  19. request.setEncoding('utf8')
  20. request.on('data', (chunk: string) => {
  21. body += chunk
  22. if (Buffer.byteLength(body) > MAX_REQUEST_BYTES) {
  23. request.destroy(new Error('DeepSeek bridge request exceeded its byte limit'))
  24. }
  25. })
  26. request.on('end', () => { resolve(body) })
  27. request.on('error', reject)
  28. })
  29. }
  30. function responseInputTexts(body: Record<string, unknown>): string[] {
  31. if (!Array.isArray(body.input)) return []
  32. return body.input.flatMap((item): string[] => {
  33. if (item === null || typeof item !== 'object') return []
  34. const content = (item as Record<string, unknown>).content
  35. if (!Array.isArray(content)) return []
  36. return content.flatMap((part): string[] => (
  37. part !== null
  38. && typeof part === 'object'
  39. && typeof (part as Record<string, unknown>).text === 'string'
  40. ? [(part as Record<string, unknown>).text as string]
  41. : []
  42. ))
  43. })
  44. }
  45. function taskText(body: Record<string, unknown>): string {
  46. const input = responseInputTexts(body).join('\n')
  47. if (input.trim().length > 0) return input
  48. return typeof body.instructions === 'string' ? body.instructions : ''
  49. }
  50. function deepSeekBaseUrl(): string {
  51. const configured = (process.env.DEEPSEEK_BASE_URL ?? OFFICIAL_DEEPSEEK_BASE_URL)
  52. .replace(/\/+$/, '')
  53. if (configured !== OFFICIAL_DEEPSEEK_BASE_URL) {
  54. throw new Error('Codex DeepSeek e2e requires the official DeepSeek base URL')
  55. }
  56. return configured
  57. }
  58. async function completeWithDeepSeek(
  59. authorization: string,
  60. task: string,
  61. ): Promise<string> {
  62. const response = await fetch(`${deepSeekBaseUrl()}/chat/completions`, {
  63. method: 'POST',
  64. headers: {
  65. authorization,
  66. 'content-type': 'application/json',
  67. },
  68. body: JSON.stringify({
  69. model: 'deepseek-v4-flash',
  70. messages: [
  71. {
  72. role: 'system',
  73. content: 'Follow the user instruction and return only the requested nonce.',
  74. },
  75. { role: 'user', content: task },
  76. ],
  77. temperature: 0,
  78. max_tokens: 64,
  79. stream: false,
  80. }),
  81. })
  82. if (!response.ok) {
  83. void response.body?.cancel()
  84. throw new Error(`DeepSeek bridge upstream returned HTTP ${response.status}`)
  85. }
  86. const payload = await response.json() as {
  87. choices?: Array<{ message?: { content?: unknown } }>
  88. }
  89. const content = payload.choices?.[0]?.message?.content
  90. if (typeof content !== 'string' || content.trim().length === 0) {
  91. throw new Error('DeepSeek bridge upstream returned no text')
  92. }
  93. return content
  94. }
  95. function closeServer(server: Server): Promise<void> {
  96. return new Promise((resolve, reject) => {
  97. server.close((error) => {
  98. if (error !== undefined) reject(error)
  99. else resolve()
  100. })
  101. server.closeAllConnections()
  102. })
  103. }
  104. /**
  105. * Start the single-purpose loopback bridge used by the Codex credentialed e2e.
  106. * @param nonce - unique answer the incoming Responses task must request.
  107. * @returns loopback endpoint, completion count, and close operation.
  108. */
  109. export async function startDeepSeekResponsesBridge(
  110. nonce: string,
  111. ): Promise<DeepSeekResponsesBridge> {
  112. let seenRequests = 0
  113. let completedRequests = 0
  114. const openResponses = new Set<ServerResponse>()
  115. const server = createServer((request, response) => {
  116. openResponses.add(response)
  117. response.on('close', () => { openResponses.delete(response) })
  118. void (async () => {
  119. if (request.method !== 'POST' || request.url !== '/v1/responses') {
  120. response.writeHead(404)
  121. response.end()
  122. return
  123. }
  124. if (seenRequests !== 0) {
  125. response.writeHead(409)
  126. response.end()
  127. return
  128. }
  129. seenRequests += 1
  130. const authorization = request.headers.authorization
  131. if (
  132. typeof authorization !== 'string'
  133. || !authorization.startsWith('Bearer ')
  134. || authorization.length === 'Bearer '.length
  135. ) {
  136. throw new Error('Codex DeepSeek bridge received no bearer credential')
  137. }
  138. const body = JSON.parse(await readRequest(request)) as Record<string, unknown>
  139. const task = taskText(body)
  140. if (!task.includes(nonce)) {
  141. throw new Error('Codex DeepSeek bridge request omitted the expected nonce')
  142. }
  143. const text = await completeWithDeepSeek(authorization, task)
  144. completedRequests += 1
  145. response.writeHead(200, {
  146. 'content-type': 'text/event-stream',
  147. 'cache-control': 'no-cache',
  148. connection: 'keep-alive',
  149. 'x-request-id': 'req_deepseek_e2e',
  150. })
  151. for (const event of completeResponsesEvents(text)) {
  152. response.write(`data: ${JSON.stringify(event)}\n\n`)
  153. }
  154. response.end('data: [DONE]\n\n')
  155. })().catch(() => {
  156. if (!response.headersSent) {
  157. response.writeHead(502, { 'content-type': 'application/json' })
  158. }
  159. response.end(JSON.stringify({ error: { message: 'DeepSeek bridge request failed' } }))
  160. })
  161. })
  162. await new Promise<void>((resolve, reject) => {
  163. server.once('error', reject)
  164. server.listen(0, '127.0.0.1', () => {
  165. server.off('error', reject)
  166. resolve()
  167. })
  168. })
  169. const address = server.address()
  170. if (address === null || typeof address === 'string') {
  171. throw new Error('DeepSeek bridge did not acquire a TCP port')
  172. }
  173. return {
  174. baseUrl: `http://127.0.0.1:${address.port}/v1`,
  175. get completedRequests(): number { return completedRequests },
  176. async close(): Promise<void> {
  177. for (const response of openResponses) response.destroy()
  178. await closeServer(server)
  179. },
  180. }
  181. }