返回源码地图

packages/llm/llm-deepseek/src/adapter.ts

main snapshot · da00f7f5358f · 正文引用章节 09 / 09;完整原文可核对,不声称全文件人工逐行审计

完整原文供逐行核对;页面收录不代表每行都经过人工语义审核。MIT 许可见 许可证。

1/** Direct Messages transport with one cancellable lifecycle per model request. */
2
3import { attributionHeaders, LlmAdapter, LlmError } from '@deepseek-ai/dsh-llm'
4import type { GenerateOptions, ImageAttachmentAccessResolver, PreparedAdapterCall, StreamChunk } from '@deepseek-ai/dsh-llm'
5import type { DeepSeekLlmApiJson } from '@deepseek-ai/dsh-deepseek-llm-api-extensions'
6import { idleWatchdog, timeoutOf } from '@deepseek-ai/dsh-timeout'
7import { modelInfo } from './model-info.ts'
8import type { DeepSeekAdapterOptions, DeepSeekConnectionOptions as Connection } from './types.ts'
9import { DeepSeekFileStore } from './file-store.ts'
10import { MESSAGES_FILES_BETA, MESSAGES_TOOL_CHANGES_BETA, messagesApiRoot } from './messages-api.ts'
11import { FileResolutionFailure, RequestFiles } from './request-files.ts'
12import { prepareRequestExtensions } from './request-extensions.ts'
13import { imagePricing, inlineImages, prepareFileIds, prepareImages } from './images.ts'
14import { serialize } from './serialize.ts'
15import { parseSse } from './sse.ts'
16import { translate } from './translate.ts'
17import { providerError, providerErrorDetail } from './transport.ts'
18
19/** DeepSeek provider using Messages content and native thinking replay. */
20export class DeepSeekAdapter<C extends Connection = Connection> extends LlmAdapter {
21 private readonly files: DeepSeekFileStore
22 private readonly imageAccess: ImageAttachmentAccessResolver = (ref) => {
23 const attachments = this.dependencies.resolveAttachments?.()
24 return attachments === undefined ? undefined : this.dependencies.resolveImageAccess?.(attachments, ref)
25 }
26
27 constructor(private readonly dependencies: DeepSeekAdapterOptions<C>) {
28 super()
29 this.files = dependencies.resolveFiles?.() ?? new DeepSeekFileStore()
30 }
31
32 override providerInfo(provider: string) { return { id: provider, name: this.dependencies.providerName ?? 'DeepSeek' } }
33 override providerRetryPolicy(_provider: string) { return this.dependencies.options().retryPolicy }
34 override async listModels(provider: string) {
35 return this.dependencies.discoverModels?.(provider) ?? []
36 }
37 override resolveModel(provider: string, model: string, _signal?: AbortSignal) {
38 return Promise.resolve(modelInfo(this.dependencies.options(), provider, model))
39 }
40 override imageRequestPricing(_provider: string, model: string) {
41 return imagePricing(this.dependencies.options(), model, this.imageAccess)
42 }
43 override prepareCall(provider: string, model: string, _signal?: AbortSignal): Promise<PreparedAdapterCall> {
44 const connection = this.dependencies.options()
45 return Promise.resolve({ model: modelInfo(connection, provider, model), stream: options => this.generate(options, connection) })
46 }
47 stream(options: GenerateOptions): AsyncIterable<StreamChunk> {
48 return this.generate(options, this.dependencies.options())
49 }
50
51 private async * generate(options: GenerateOptions, connection: C): AsyncGenerator<StreamChunk> {
52 const consumer = new AbortController()
53 const signal = options.signal === undefined ? consumer.signal : AbortSignal.any([consumer.signal, options.signal])
54 using watchdog = idleWatchdog(signal, connection.streamIdleTimeoutMs, 'MESSAGES_IDLE')
55 const iterator = this.request(options, connection, watchdog.signal, () => { watchdog.pulse() })
56 try {
57 while (true) {
58 const next = await watchdog.next(iterator)
59 if (next.done) return
60 yield next.value
61 }
62 } catch (error) {
63 if (timeoutOf(watchdog.signal, 'MESSAGES_IDLE') !== undefined) throw new LlmError('DeepSeek Messages stream idle timeout', 'TIMEOUT', { cause: error })
64 if (options.signal?.aborted) throw new LlmError('DeepSeek Messages request aborted', 'ABORTED', { cause: error })
65 if (error instanceof LlmError) throw error
66 throw new LlmError('DeepSeek Messages transport failed', 'TRANSPORT', { cause: error })
67 } finally {
68 consumer.abort()
69 try { await iterator.return(undefined) } catch (_abortedRequestCleanup) {
70 // The request already settled; aborting its reader cannot replace that outcome.
71 }
72 }
73 }
74
75 private async * request(
76 options: GenerateOptions, connection: C, signal: AbortSignal, activity: () => void,
77 ): AsyncGenerator<StreamChunk> {
78 signal.throwIfAborted()
79 const { messages, versions } = await prepareImages(
80 options.messages, connection, options.model, this.dependencies.resolveAttachments?.(), this.imageAccess, signal,
81 )
82 const auth = await this.dependencies.resolveAuth(connection)
83 try {
84 const files = new RequestFiles(this.files, {
85 baseURL: connection.baseURL, headers: auth.headers,
86 },
87 connection.filePolicy, connection.filesApiTimeoutMs, signal, activity)
88 let inline = false
89 while (true) {
90 signal.throwIfAborted()
91 files.beginAttempt()
92 let fileIds: Awaited<ReturnType<typeof prepareFileIds>> | undefined
93 if (!inline) {
94 try {
95 fileIds = await prepareFileIds(messages, versions, files)
96 } catch (error) {
97 if (!(error instanceof FileResolutionFailure)) throw error
98 inline = true
99 continue
100 }
101 }
102 const history = inline ? inlineImages(messages, versions, connection) : messages
103 const body = serialize(options, connection, history, versions, this.imageAccess, (reason) => {
104 this.dependencies.onReplayDegrade?.({ provider: options.provider, model: options.model, reason })
105 }, fileIds)
106 const extensions = await prepareRequestExtensions(body as Readonly<Record<string, DeepSeekLlmApiJson>>, {
107 signal,
108 ...options.sessionId === undefined ? {} : { sessionId: String(options.sessionId) },
109 ...options.purpose === undefined ? {} : { purpose: options.purpose },
110 }, this.dependencies.prepareExtensions, (fields, error) => {
111 this.dependencies.onExtensionsOmitted?.({ provider: options.provider, model: options.model, fields, error })
112 })
113 signal.throwIfAborted()
114 const betas = [
115 ...fileIds !== undefined && fileIds.size > 0 ? [MESSAGES_FILES_BETA] : [],
116 ...body.messages.some(message => message.content.some(block => block.type === 'tool_addition' || block.type === 'tool_removal'))
117 ? [MESSAGES_TOOL_CHANGES_BETA]
118 : [],
119 ]
120 const response = await fetch(`${messagesApiRoot(connection.baseURL)}/messages`, {
121 method: 'POST', signal, body: extensions.payload, redirect: 'error',
122 headers: {
123 ...attributionHeaders(),
124 'content-type': 'application/json', 'accept': 'text/event-stream',
125 ...auth.headers,
126 'anthropic-version': '2023-06-01',
127 ...betas.length === 0 ? {} : { 'anthropic-beta': betas.join(',') },
128 'x-deepseek-harness-user-id': this.dependencies.resolveUserId(),
129 ...options.sessionId === undefined ? {} : { 'x-deepseek-harness-session-id': String(options.sessionId) },
130 ...options.purpose === 'compaction' ? { 'x-deepseek-harness-compact': '1' } : {},
131 },
132 })
133 if (!response.ok) {
134 const text = await response.text()
135 let raw: unknown
136 try { raw = JSON.parse(text) } catch (_nonJsonGatewayError) {
137 // HTTP status is authoritative when a gateway does not return JSON.
138 }
139 const detail = providerErrorDetail(raw)
140 if (await files.retry(detail)) continue
141 const failure = providerError(raw, response.status, response.headers)
142 const message = files.errorMessage(response.status, failure.message, detail)
143 throw new LlmError(message, failure.code, { ...failure.failure, cause: new Error(text) })
144 }
145 await extensions.accept()
146 if (response.body === null) throw new LlmError('DeepSeek Messages returned no response body', 'EMPTY_RESPONSE')
147 yield* translate(parseSse(response.body, activity), options.model)
148 return
149 }
150 } catch (error) {
151 if (auth.onRequestError !== undefined) {
152 let mapped: unknown
153 try { mapped = await auth.onRequestError(error) }
154 catch (_credentialUpdateFailed) { throw error }
155 throw mapped
156 }
157 throw error
158 }
159 }
160}