1
/** Direct Messages transport with one cancellable lifecycle per model request. */3
import { attributionHeaders, LlmAdapter, LlmError } from '@deepseek-ai/dsh-llm'4
import type { GenerateOptions, ImageAttachmentAccessResolver, PreparedAdapterCall, StreamChunk } from '@deepseek-ai/dsh-llm'5
import type { DeepSeekLlmApiJson } from '@deepseek-ai/dsh-deepseek-llm-api-extensions'6
import { idleWatchdog, timeoutOf } from '@deepseek-ai/dsh-timeout'7
import { modelInfo } from './model-info.ts'8
import type { DeepSeekAdapterOptions, DeepSeekConnectionOptions as Connection } from './types.ts'9
import { DeepSeekFileStore } from './file-store.ts'10
import { MESSAGES_FILES_BETA, MESSAGES_TOOL_CHANGES_BETA, messagesApiRoot } from './messages-api.ts'11
import { FileResolutionFailure, RequestFiles } from './request-files.ts'12
import { prepareRequestExtensions } from './request-extensions.ts'13
import { imagePricing, inlineImages, prepareFileIds, prepareImages } from './images.ts'14
import { serialize } from './serialize.ts'15
import { parseSse } from './sse.ts'16
import { translate } from './translate.ts'17
import { providerError, providerErrorDetail } from './transport.ts'19
/** DeepSeek provider using Messages content and native thinking replay. */20
export class DeepSeekAdapter<C extends Connection = Connection> extends LlmAdapter {21
private readonly files: DeepSeekFileStore22
private readonly imageAccess: ImageAttachmentAccessResolver = (ref) => {23
const attachments = this.dependencies.resolveAttachments?.()24
return attachments === undefined ? undefined : this.dependencies.resolveImageAccess?.(attachments, ref)25
}27
constructor(private readonly dependencies: DeepSeekAdapterOptions<C>) {28
super()29
this.files = dependencies.resolveFiles?.() ?? new DeepSeekFileStore()30
}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
}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) return60
yield next.value61
}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 error66
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
}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 = false89
while (true) {90
signal.throwIfAborted()91
files.beginAttempt()92
let fileIds: Awaited<ReturnType<typeof prepareFileIds>> | undefined93
if (!inline) {94
try {95
fileIds = await prepareFileIds(messages, versions, files)96
} catch (error) {97
if (!(error instanceof FileResolutionFailure)) throw error98
inline = true99
continue100
}101
}102
const history = inline ? inlineImages(messages, versions, connection) : messages103
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: unknown136
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)) continue141
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
return149
}150
} catch (error) {151
if (auth.onRequestError !== undefined) {152
let mapped: unknown153
try { mapped = await auth.onRequestError(error) }154
catch (_credentialUpdateFailed) { throw error }155
throw mapped156
}157
throw error158
}159
}160
}