Browse Source
Introduce @dm/xllm with a unified request/response model, streaming AsyncIterable events, and adapter-based support for OpenAI-compatible and DeepSeek backends. Add an example integration that demonstrates provider/model switching via Vite env variables and direct stream output consumption. Made-with: Cursormain
20 changed files with 1122 additions and 0 deletions
@ -0,0 +1,236 @@ |
|||
# xllm 统一大模型请求库设计文档 |
|||
|
|||
## 1. 背景与目标 |
|||
|
|||
在 `packages/xllm` 新建一个可跨 Node 与浏览器使用的大模型请求库,满足以下核心诉求: |
|||
|
|||
- 支持流式输出,并可直接消费增量内容 |
|||
- 提供供应商适配器,将请求和响应统一为一致结构 |
|||
- 开发体验简洁,支持便捷切换供应商与模型 |
|||
|
|||
首版范围明确为: |
|||
|
|||
- 运行环境:Node + 浏览器同一套 API |
|||
- 首批供应商:OpenAI-Compatible、DeepSeek |
|||
- 流式消费方式:`for await...of` |
|||
- 能力范围:文本对话 + 工具调用 + 多模态输入 |
|||
- 配置策略:参数优先,环境变量兜底 |
|||
|
|||
## 2. 非目标(首版不做) |
|||
|
|||
- 不引入复杂重试与熔断策略(仅保留最小错误归一化能力) |
|||
- 不实现供应商专有高级能力(如某家独有推理模式)的一致化抽象 |
|||
- 不做 UI 层封装,仅提供 SDK 级能力 |
|||
|
|||
## 3. 总体架构 |
|||
|
|||
采用“核心领域模型 + 供应商插件适配器”架构(推荐方案 A)。 |
|||
|
|||
### 3.1 模块划分 |
|||
|
|||
建议目录: |
|||
|
|||
- `packages/xllm/src/core/` |
|||
- 统一类型定义:`XRequest`、`XMessage`、`XStreamEvent`、`XResponse`、`XToolCall`、`XUsage` |
|||
- 统一错误模型:`XllmError` |
|||
- `packages/xllm/src/providers/` |
|||
- `openai-compatible.adapter.ts` |
|||
- `deepseek.adapter.ts` |
|||
- `registry.ts`(按 provider 名称注册/获取适配器) |
|||
- `packages/xllm/src/client/` |
|||
- `createXllm.ts`(实例创建与默认配置) |
|||
- `stream.ts`(流式执行) |
|||
- `generate.ts`(非流式执行) |
|||
- `packages/xllm/src/runtime/` |
|||
- `fetch.ts`(跨平台请求封装) |
|||
- `sse.ts`(SSE 增量解析) |
|||
- `env.ts`(环境变量读取) |
|||
- `packages/xllm/src/index.ts` |
|||
- 对外导出 API 与类型 |
|||
|
|||
### 3.2 核心设计原则 |
|||
|
|||
- 业务入口稳定:调用方只面对 `xllm.stream()` / `xllm.generate()` |
|||
- 差异内聚:供应商差异仅存在于 adapter 内 |
|||
- 可扩展:新增供应商只需实现同一 adapter 契约并注册 |
|||
|
|||
## 4. 统一数据模型 |
|||
|
|||
### 4.1 请求模型 |
|||
|
|||
`XRequest`(逻辑结构): |
|||
|
|||
- `messages: XMessage[]` |
|||
- `tools?: XToolDefinition[]` |
|||
- `toolChoice?: "auto" | "none" | { name: string }` |
|||
- `stream?: boolean` |
|||
- `temperature?: number` |
|||
- `topP?: number` |
|||
- `maxTokens?: number` |
|||
- `metadata?: Record<string, string>` |
|||
- 请求级覆盖字段: |
|||
- `provider?: "openai-compatible" | "deepseek"` |
|||
- `model?: string` |
|||
- `apiKey?: string` |
|||
- `baseURL?: string` |
|||
|
|||
`XMessage`: |
|||
|
|||
- `role: "system" | "user" | "assistant" | "tool"` |
|||
- `content: XContentPart[]` |
|||
- `toolCallId?: string`(role 为 `tool` 时用于关联) |
|||
|
|||
`XContentPart`: |
|||
|
|||
- `{ type: "text"; text: string }` |
|||
- `{ type: "image"; image: { url: string } }` |
|||
|
|||
> 首版仅强制支持 `url`(包含 http(s) 与 data URL),后续可扩展 file/blob 引用。 |
|||
|
|||
### 4.2 响应模型 |
|||
|
|||
`XResponse`: |
|||
|
|||
- `text: string` |
|||
- `toolCalls: XToolCall[]` |
|||
- `usage?: XUsage` |
|||
- `provider: string` |
|||
- `model: string` |
|||
- `finishReason?: string` |
|||
- `raw?: unknown`(可选调试字段) |
|||
|
|||
## 5. 流式事件协议 |
|||
|
|||
`stream()` 返回 `AsyncIterable<XStreamEvent>`,统一事件集合: |
|||
|
|||
- `response.start` |
|||
- 响应开始,包含 provider/model/requestId 等元数据 |
|||
- `text.delta` |
|||
- 文本增量,用于直接输出 |
|||
- `tool_call.delta` |
|||
- 工具调用参数增量 |
|||
- `tool_call.done` |
|||
- 单个工具调用完成,输出完整 name/args |
|||
- `response.usage` |
|||
- token 使用统计 |
|||
- `response.done` |
|||
- 完成事件,含 finish reason |
|||
|
|||
说明: |
|||
|
|||
- 默认将供应商 SSE 分片归一化为上述事件 |
|||
- 消费侧推荐只关心 `text.delta` 与 `response.done` 即可完成常见直出场景 |
|||
|
|||
## 6. 供应商适配器契约 |
|||
|
|||
定义统一接口 `ProviderAdapter`: |
|||
|
|||
- `name: string` |
|||
- `toProviderRequest(input: XRequest, config: ResolvedConfig): ProviderHttpRequest` |
|||
- `fromProviderResponse(raw: unknown): XResponse` |
|||
- `fromProviderStreamChunk(chunk: string, state: StreamState): XStreamEvent[]` |
|||
- `normalizeError(err: unknown): XllmError` |
|||
|
|||
适配职责: |
|||
|
|||
- 请求映射:统一字段 -> 供应商字段 |
|||
- 流式映射:供应商 chunk -> 统一事件 |
|||
- 错误映射:供应商错误 -> 统一错误码 |
|||
|
|||
## 7. 配置与切换策略 |
|||
|
|||
### 7.1 配置优先级 |
|||
|
|||
最终配置合并顺序(高到低): |
|||
|
|||
1. 请求级参数(`stream/generate` 调用时传入) |
|||
2. 实例默认参数(`createXllm`) |
|||
3. 环境变量兜底 |
|||
|
|||
### 7.2 建议环境变量 |
|||
|
|||
- `XLLM_API_KEY` |
|||
- `OPENAI_API_KEY` |
|||
- `DEEPSEEK_API_KEY` |
|||
- `XLLM_BASE_URL`(可选) |
|||
|
|||
### 7.3 切换体验 |
|||
|
|||
- `createXllm({ provider, model, ... })` 设全局默认 |
|||
- 单次调用可覆写 `{ provider, model, ... }` |
|||
- 同一业务逻辑不变,仅替换 provider/model 参数即可切换模型后端 |
|||
|
|||
## 8. 错误模型与稳定性 |
|||
|
|||
统一错误码建议: |
|||
|
|||
- `AUTH_ERROR` |
|||
- `RATE_LIMIT` |
|||
- `NETWORK_ERROR` |
|||
- `INVALID_REQUEST` |
|||
- `PROVIDER_ERROR` |
|||
|
|||
`XllmError` 结构包含: |
|||
|
|||
- `code` |
|||
- `message` |
|||
- `provider` |
|||
- `statusCode?` |
|||
- `requestId?` |
|||
- `raw?` |
|||
|
|||
流式错误处理约定: |
|||
|
|||
- 出现不可恢复错误时抛出 `XllmError` |
|||
- 已发出的增量内容由调用方决定是否展示或回滚 |
|||
|
|||
## 9. 测试策略 |
|||
|
|||
### 9.1 单元测试 |
|||
|
|||
- adapter 请求映射测试 |
|||
- adapter 非流式响应映射测试 |
|||
- adapter 流式 chunk 解析测试(含半包、粘包、结束包) |
|||
|
|||
### 9.2 集成测试 |
|||
|
|||
- mock fetch + mock SSE,覆盖 `stream` 与 `generate` 主流程 |
|||
- 断言统一事件顺序与最终响应一致性 |
|||
|
|||
### 9.3 契约测试 |
|||
|
|||
- 同一 `XRequest` 在不同 provider 下返回结构都满足统一断言 |
|||
- 确保新增 provider 不破坏现有统一接口语义 |
|||
|
|||
## 10. 包结构与发布约定(与仓库一致) |
|||
|
|||
新包遵循现有 monorepo 规范: |
|||
|
|||
- 包名建议:`@dm/xllm` |
|||
- 入口:`src/index.ts` |
|||
- `package.json` 中导出: |
|||
- `"development": "./src/index.ts"` |
|||
- `"import": "./dist/index.js"` |
|||
- `"types": "./dist/index.d.ts"` |
|||
- 脚本:`"dev": "dm dev"`、`"build": "dm build"` |
|||
|
|||
## 11. 里程碑与实施顺序 |
|||
|
|||
建议分 3 个小阶段: |
|||
|
|||
1. **M1(骨架)** |
|||
- 建立 core 类型、client 主入口、provider registry、基础测试框架 |
|||
2. **M2(能力)** |
|||
- 完成 openai-compatible 与 deepseek 的 stream/generate 适配 |
|||
- 打通文本、多模态、工具调用映射 |
|||
3. **M3(稳态)** |
|||
- 补齐错误归一化与契约测试 |
|||
- 完善 README 与示例用法 |
|||
|
|||
## 12. 验收标准 |
|||
|
|||
- 可在 Node 与浏览器中通过同一 API 调用 |
|||
- 能通过 `for await...of` 获取稳定的 `text.delta` 直出 |
|||
- OpenAI-Compatible 与 DeepSeek 可通过仅修改 provider/model 切换 |
|||
- 非流式返回结构与流式聚合结果语义一致 |
|||
- 关键错误能归一化并携带可诊断信息 |
|||
@ -0,0 +1,25 @@ |
|||
# @dm/xllm |
|||
|
|||
统一的大模型请求库,支持: |
|||
|
|||
- 流式输出(`for await...of`) |
|||
- OpenAI-Compatible 与 DeepSeek 供应商适配 |
|||
- 文本、多模态输入(图片 URL)、工具调用结构统一 |
|||
|
|||
## 快速开始 |
|||
|
|||
```ts |
|||
import { createXllm } from "@dm/xllm"; |
|||
|
|||
const xllm = createXllm({ |
|||
provider: "deepseek", |
|||
model: "deepseek-chat", |
|||
apiKey: process.env.DEEPSEEK_API_KEY, |
|||
}); |
|||
|
|||
for await (const event of xllm.stream({ |
|||
messages: [{ role: "user", content: [{ type: "text", text: "你好" }] }], |
|||
})) { |
|||
if (event.type === "text.delta") process.stdout.write(event.text); |
|||
} |
|||
``` |
|||
@ -0,0 +1,23 @@ |
|||
{ |
|||
"name": "@dm/xllm", |
|||
"type": "module", |
|||
"version": "0.0.1-alpha.1", |
|||
"scripts": { |
|||
"dev": "dm dev", |
|||
"build": "dm build" |
|||
}, |
|||
"main": "./dist/index.js", |
|||
"module": "./dist/index.js", |
|||
"types": "./dist/index.d.ts", |
|||
"files": [ |
|||
"dist" |
|||
], |
|||
"exports": { |
|||
".": { |
|||
"development": "./src/index.ts", |
|||
"import": "./dist/index.js", |
|||
"types": "./dist/index.d.ts" |
|||
}, |
|||
"./package.json": "./package.json" |
|||
} |
|||
} |
|||
@ -0,0 +1,21 @@ |
|||
import type { XClientOptions, XRequest, XResponse, XStreamEvent } from "../core/types"; |
|||
import { generate } from "./generate"; |
|||
import { stream } from "./stream"; |
|||
|
|||
export interface XllmClient { |
|||
stream(request: XRequest): AsyncIterable<XStreamEvent>; |
|||
generate(request: XRequest): Promise<XResponse>; |
|||
with(overrides: Partial<XClientOptions>): XllmClient; |
|||
} |
|||
|
|||
export const createXllm = (options: XClientOptions = {}): XllmClient => ({ |
|||
stream(request: XRequest) { |
|||
return stream(options, request); |
|||
}, |
|||
generate(request: XRequest) { |
|||
return generate(options, request); |
|||
}, |
|||
with(overrides: Partial<XClientOptions>) { |
|||
return createXllm({ ...options, ...overrides }); |
|||
}, |
|||
}); |
|||
@ -0,0 +1,17 @@ |
|||
import type { XClientOptions, XRequest, XResponse } from "../core/types"; |
|||
import { getProviderAdapter } from "../providers/registry"; |
|||
import { resolveConfig } from "../runtime/config"; |
|||
import { postJSON, throwForBadStatus } from "../runtime/http"; |
|||
|
|||
export const generate = async (options: XClientOptions, request: XRequest): Promise<XResponse> => { |
|||
const resolved = resolveConfig(request, options); |
|||
const adapter = getProviderAdapter(resolved.provider); |
|||
const providerRequest = adapter.toProviderRequest(request, resolved, false); |
|||
const response = await postJSON(resolved.fetchImpl, providerRequest, adapter); |
|||
|
|||
await throwForBadStatus(response, adapter); |
|||
|
|||
const payload = await response.json(); |
|||
const normalized = adapter.fromProviderResponse(payload, resolved.provider); |
|||
return { ...normalized, provider: resolved.provider, model: normalized.model || resolved.model }; |
|||
}; |
|||
@ -0,0 +1,53 @@ |
|||
import { XllmError } from "../core/errors"; |
|||
import type { XClientOptions, XRequest, XStreamEvent } from "../core/types"; |
|||
import { getProviderAdapter } from "../providers/registry"; |
|||
import type { StreamState } from "../providers/types"; |
|||
import { resolveConfig } from "../runtime/config"; |
|||
import { postJSON, throwForBadStatus } from "../runtime/http"; |
|||
import { parseSSE } from "../runtime/sse"; |
|||
|
|||
export const stream = async function* ( |
|||
options: XClientOptions, |
|||
request: XRequest, |
|||
): AsyncGenerator<XStreamEvent> { |
|||
const resolved = resolveConfig(request, options); |
|||
const adapter = getProviderAdapter(resolved.provider); |
|||
const providerRequest = adapter.toProviderRequest(request, resolved, true); |
|||
const response = await postJSON(resolved.fetchImpl, providerRequest, adapter); |
|||
|
|||
await throwForBadStatus(response, adapter); |
|||
|
|||
if (!response.body) { |
|||
throw new XllmError({ |
|||
code: "NETWORK_ERROR", |
|||
message: "Response body is empty in streaming mode", |
|||
provider: resolved.provider, |
|||
}); |
|||
} |
|||
|
|||
const streamState: StreamState = { started: false, doneEmitted: false, toolCallsByIndex: new Map() }; |
|||
for await (const data of parseSSE(response.body)) { |
|||
if (data === "[DONE]") { |
|||
if (!streamState.started) { |
|||
yield { type: "response.start", provider: resolved.provider, model: resolved.model }; |
|||
} |
|||
if (!streamState.doneEmitted) { |
|||
yield { type: "response.done" }; |
|||
streamState.doneEmitted = true; |
|||
} |
|||
break; |
|||
} |
|||
|
|||
let payload: unknown; |
|||
try { |
|||
payload = JSON.parse(data); |
|||
} catch { |
|||
continue; |
|||
} |
|||
|
|||
const events = adapter.fromProviderStreamChunk(payload, streamState); |
|||
for (const event of events) { |
|||
yield event; |
|||
} |
|||
} |
|||
}; |
|||
@ -0,0 +1,31 @@ |
|||
export type XllmErrorCode = |
|||
| "AUTH_ERROR" |
|||
| "RATE_LIMIT" |
|||
| "NETWORK_ERROR" |
|||
| "INVALID_REQUEST" |
|||
| "PROVIDER_ERROR"; |
|||
|
|||
export class XllmError extends Error { |
|||
public readonly code: XllmErrorCode; |
|||
public readonly provider: string; |
|||
public readonly statusCode?: number; |
|||
public readonly requestId?: string; |
|||
public readonly raw?: unknown; |
|||
|
|||
public constructor(params: { |
|||
code: XllmErrorCode; |
|||
message: string; |
|||
provider: string; |
|||
statusCode?: number; |
|||
requestId?: string; |
|||
raw?: unknown; |
|||
}) { |
|||
super(params.message); |
|||
this.name = "XllmError"; |
|||
this.code = params.code; |
|||
this.provider = params.provider; |
|||
this.statusCode = params.statusCode; |
|||
this.requestId = params.requestId; |
|||
this.raw = params.raw; |
|||
} |
|||
} |
|||
@ -0,0 +1,89 @@ |
|||
export type XProviderName = "openai-compatible" | "deepseek"; |
|||
|
|||
export type XMessageRole = "system" | "user" | "assistant" | "tool"; |
|||
|
|||
export type XTextPart = { type: "text"; text: string }; |
|||
export type XImagePart = { type: "image"; image: { url: string } }; |
|||
export type XContentPart = XTextPart | XImagePart; |
|||
|
|||
export type XMessage = |
|||
| { |
|||
role: "system" | "user" | "assistant"; |
|||
content: XContentPart[]; |
|||
} |
|||
| { |
|||
role: "tool"; |
|||
content: XContentPart[]; |
|||
toolCallId: string; |
|||
}; |
|||
|
|||
export interface XToolDefinition { |
|||
name: string; |
|||
description?: string; |
|||
parameters?: Record<string, unknown>; |
|||
} |
|||
|
|||
export interface XToolCall { |
|||
id: string; |
|||
name: string; |
|||
arguments: string; |
|||
} |
|||
|
|||
export type XToolChoice = "auto" | "none" | { name: string }; |
|||
|
|||
export interface XRequest { |
|||
messages: XMessage[]; |
|||
tools?: XToolDefinition[]; |
|||
toolChoice?: XToolChoice; |
|||
stream?: boolean; |
|||
temperature?: number; |
|||
topP?: number; |
|||
maxTokens?: number; |
|||
metadata?: Record<string, string>; |
|||
provider?: XProviderName; |
|||
model?: string; |
|||
apiKey?: string; |
|||
baseURL?: string; |
|||
} |
|||
|
|||
export interface XUsage { |
|||
promptTokens?: number; |
|||
completionTokens?: number; |
|||
totalTokens?: number; |
|||
} |
|||
|
|||
export interface XResponse { |
|||
text: string; |
|||
toolCalls: XToolCall[]; |
|||
usage?: XUsage; |
|||
provider: XProviderName; |
|||
model: string; |
|||
finishReason?: string; |
|||
raw?: unknown; |
|||
} |
|||
|
|||
export type XStreamEvent = |
|||
| { type: "response.start"; provider: XProviderName; model: string; requestId?: string } |
|||
| { type: "text.delta"; text: string } |
|||
| { type: "tool_call.delta"; id: string; name?: string; argumentsDelta?: string } |
|||
| { type: "tool_call.done"; toolCall: XToolCall } |
|||
| { type: "response.usage"; usage: XUsage } |
|||
| { type: "response.done"; finishReason?: string }; |
|||
|
|||
export interface XClientOptions { |
|||
provider?: XProviderName; |
|||
model?: string; |
|||
apiKey?: string; |
|||
baseURL?: string; |
|||
fetch?: typeof fetch; |
|||
headers?: Record<string, string>; |
|||
} |
|||
|
|||
export interface ResolvedConfig { |
|||
provider: XProviderName; |
|||
model: string; |
|||
apiKey: string; |
|||
baseURL: string; |
|||
fetchImpl: typeof fetch; |
|||
headers: Record<string, string>; |
|||
} |
|||
@ -0,0 +1,90 @@ |
|||
import { describe, expect, it } from "vitest"; |
|||
import { createXllm, type XStreamEvent } from "./index"; |
|||
|
|||
const toJsonResponse = (payload: unknown): Response => |
|||
new Response(JSON.stringify(payload), { |
|||
status: 200, |
|||
headers: { "content-type": "application/json" }, |
|||
}); |
|||
|
|||
const toSSEBody = (frames: string[]): ReadableStream<Uint8Array> => { |
|||
const encoder = new TextEncoder(); |
|||
return new ReadableStream<Uint8Array>({ |
|||
start(controller) { |
|||
for (const frame of frames) { |
|||
controller.enqueue(encoder.encode(frame)); |
|||
} |
|||
controller.close(); |
|||
}, |
|||
}); |
|||
}; |
|||
|
|||
describe("xllm", () => { |
|||
it("normalizes non-stream response", async () => { |
|||
const mockFetch: typeof fetch = async () => |
|||
toJsonResponse({ |
|||
model: "gpt-4o-mini", |
|||
choices: [ |
|||
{ |
|||
message: { |
|||
content: "hello", |
|||
}, |
|||
finish_reason: "stop", |
|||
}, |
|||
], |
|||
usage: { |
|||
prompt_tokens: 1, |
|||
completion_tokens: 1, |
|||
total_tokens: 2, |
|||
}, |
|||
}); |
|||
|
|||
const client = createXllm({ |
|||
provider: "openai-compatible", |
|||
model: "gpt-4o-mini", |
|||
apiKey: "test", |
|||
fetch: mockFetch, |
|||
}); |
|||
const result = await client.generate({ |
|||
messages: [{ role: "user", content: [{ type: "text", text: "hi" }] }], |
|||
}); |
|||
|
|||
expect(result.text).toBe("hello"); |
|||
expect(result.usage?.totalTokens).toBe(2); |
|||
}); |
|||
|
|||
it("streams text delta events", async () => { |
|||
const sseFrames = [ |
|||
'data: {"id":"r1","model":"deepseek-chat","choices":[{"delta":{"content":"he"},"finish_reason":null}]}\n\n', |
|||
'data: {"id":"r1","model":"deepseek-chat","choices":[{"delta":{"content":"llo"},"finish_reason":"stop"}]}\n\n', |
|||
"data: [DONE]\n\n", |
|||
]; |
|||
|
|||
const mockFetch: typeof fetch = async () => |
|||
new Response(toSSEBody(sseFrames), { |
|||
status: 200, |
|||
headers: { "content-type": "text/event-stream" }, |
|||
}); |
|||
|
|||
const client = createXllm({ |
|||
provider: "deepseek", |
|||
model: "deepseek-chat", |
|||
apiKey: "test", |
|||
fetch: mockFetch, |
|||
}); |
|||
|
|||
const events: XStreamEvent[] = []; |
|||
for await (const event of client.stream({ |
|||
messages: [{ role: "user", content: [{ type: "text", text: "hi" }] }], |
|||
})) { |
|||
events.push(event); |
|||
} |
|||
|
|||
const text = events |
|||
.filter((event): event is Extract<XStreamEvent, { type: "text.delta" }> => event.type === "text.delta") |
|||
.map((event) => event.text) |
|||
.join(""); |
|||
expect(text).toBe("hello"); |
|||
expect(events.some((event) => event.type === "response.done")).toBe(true); |
|||
}); |
|||
}); |
|||
@ -0,0 +1,15 @@ |
|||
export { createXllm } from "./client/create-xllm"; |
|||
export { XllmError } from "./core/errors"; |
|||
export type { |
|||
XClientOptions, |
|||
XContentPart, |
|||
XMessage, |
|||
XProviderName, |
|||
XRequest, |
|||
XResponse, |
|||
XStreamEvent, |
|||
XToolCall, |
|||
XToolChoice, |
|||
XToolDefinition, |
|||
XUsage, |
|||
} from "./core/types"; |
|||
@ -0,0 +1,45 @@ |
|||
import { XllmError } from "../core/errors"; |
|||
import type { ProviderAdapter } from "./types"; |
|||
import { openAICompatibleAdapter } from "./openai-compatible.adapter"; |
|||
|
|||
const DEFAULT_DEEPSEEK_BASE_URL = "https://api.deepseek.com"; |
|||
|
|||
export const deepseekAdapter: ProviderAdapter = { |
|||
name: "deepseek", |
|||
|
|||
toProviderRequest(input, config, stream) { |
|||
return openAICompatibleAdapter.toProviderRequest( |
|||
input, |
|||
{ |
|||
...config, |
|||
baseURL: config.baseURL || DEFAULT_DEEPSEEK_BASE_URL, |
|||
}, |
|||
stream, |
|||
); |
|||
}, |
|||
|
|||
fromProviderResponse(raw) { |
|||
const response = openAICompatibleAdapter.fromProviderResponse(raw, "deepseek"); |
|||
return { ...response, provider: "deepseek" }; |
|||
}, |
|||
|
|||
fromProviderStreamChunk(chunkData, state) { |
|||
const events = openAICompatibleAdapter.fromProviderStreamChunk(chunkData, state); |
|||
return events.map((event) => { |
|||
if (event.type !== "response.start") return event; |
|||
return { ...event, provider: "deepseek" }; |
|||
}); |
|||
}, |
|||
|
|||
normalizeError(error) { |
|||
const baseError = openAICompatibleAdapter.normalizeError(error); |
|||
return new XllmError({ |
|||
code: baseError.code, |
|||
message: baseError.message, |
|||
provider: "deepseek", |
|||
statusCode: baseError.statusCode, |
|||
requestId: baseError.requestId, |
|||
raw: baseError.raw, |
|||
}); |
|||
}, |
|||
}; |
|||
@ -0,0 +1,233 @@ |
|||
import { XllmError } from "../core/errors"; |
|||
import type { |
|||
XContentPart, |
|||
XMessage, |
|||
XProviderName, |
|||
XRequest, |
|||
XResponse, |
|||
XStreamEvent, |
|||
XToolCall, |
|||
XUsage, |
|||
} from "../core/types"; |
|||
import type { ProviderAdapter, StreamState } from "./types"; |
|||
|
|||
const DEFAULT_OPENAI_COMPATIBLE_BASE_URL = "https://api.openai.com/v1"; |
|||
|
|||
const toProviderContent = (parts: XContentPart[]): string | Array<Record<string, unknown>> => { |
|||
if (parts.length === 1 && parts[0]?.type === "text") { |
|||
return parts[0].text; |
|||
} |
|||
|
|||
return parts.map((part) => { |
|||
if (part.type === "text") { |
|||
return { type: "text", text: part.text }; |
|||
} |
|||
|
|||
return { type: "image_url", image_url: { url: part.image.url } }; |
|||
}); |
|||
}; |
|||
|
|||
const toProviderMessage = (message: XMessage): Record<string, unknown> => { |
|||
if (message.role === "tool") { |
|||
return { |
|||
role: "tool", |
|||
content: toProviderContent(message.content), |
|||
tool_call_id: message.toolCallId, |
|||
}; |
|||
} |
|||
|
|||
return { |
|||
role: message.role, |
|||
content: toProviderContent(message.content), |
|||
}; |
|||
}; |
|||
|
|||
const toUsage = (usage: any): XUsage | undefined => { |
|||
if (!usage) return undefined; |
|||
return { |
|||
promptTokens: usage.prompt_tokens, |
|||
completionTokens: usage.completion_tokens, |
|||
totalTokens: usage.total_tokens, |
|||
}; |
|||
}; |
|||
|
|||
const toToolCall = (toolCall: any): XToolCall => ({ |
|||
id: toolCall.id ?? `tool_call_${Math.random().toString(36).slice(2)}`, |
|||
name: toolCall.function?.name ?? "", |
|||
arguments: toolCall.function?.arguments ?? "", |
|||
}); |
|||
|
|||
const mapToolChoice = (toolChoice: XRequest["toolChoice"]): unknown => { |
|||
if (!toolChoice) return undefined; |
|||
if (toolChoice === "auto" || toolChoice === "none") return toolChoice; |
|||
return { type: "function", function: { name: toolChoice.name } }; |
|||
}; |
|||
|
|||
export const openAICompatibleAdapter: ProviderAdapter = { |
|||
name: "openai-compatible", |
|||
|
|||
toProviderRequest(input, config, stream) { |
|||
const body: Record<string, unknown> = { |
|||
model: config.model, |
|||
messages: input.messages.map(toProviderMessage), |
|||
stream, |
|||
temperature: input.temperature, |
|||
top_p: input.topP, |
|||
max_tokens: input.maxTokens, |
|||
metadata: input.metadata, |
|||
}; |
|||
|
|||
if (input.tools && input.tools.length > 0) { |
|||
body.tools = input.tools.map((tool) => ({ |
|||
type: "function", |
|||
function: { |
|||
name: tool.name, |
|||
description: tool.description, |
|||
parameters: tool.parameters ?? { type: "object", properties: {} }, |
|||
}, |
|||
})); |
|||
body.tool_choice = mapToolChoice(input.toolChoice); |
|||
} |
|||
|
|||
return { |
|||
method: "POST", |
|||
url: `${config.baseURL || DEFAULT_OPENAI_COMPATIBLE_BASE_URL}/chat/completions`, |
|||
headers: { |
|||
"content-type": "application/json", |
|||
authorization: `Bearer ${config.apiKey}`, |
|||
...config.headers, |
|||
}, |
|||
body, |
|||
}; |
|||
}, |
|||
|
|||
fromProviderResponse(raw, provider) { |
|||
const response = raw as any; |
|||
const choice = response?.choices?.[0] ?? {}; |
|||
const message = choice?.message ?? {}; |
|||
|
|||
return { |
|||
text: message?.content ?? "", |
|||
toolCalls: Array.isArray(message?.tool_calls) ? message.tool_calls.map(toToolCall) : [], |
|||
usage: toUsage(response?.usage), |
|||
provider, |
|||
model: response?.model ?? "", |
|||
finishReason: choice?.finish_reason, |
|||
raw: response, |
|||
}; |
|||
}, |
|||
|
|||
fromProviderStreamChunk(chunkData, state) { |
|||
const payload = chunkData as any; |
|||
const events: XStreamEvent[] = []; |
|||
|
|||
if (!state.started) { |
|||
state.started = true; |
|||
events.push({ |
|||
type: "response.start", |
|||
provider: "openai-compatible", |
|||
model: payload?.model ?? "", |
|||
requestId: payload?.id, |
|||
}); |
|||
} |
|||
|
|||
const choice = payload?.choices?.[0]; |
|||
if (!choice) return events; |
|||
|
|||
const delta = choice.delta ?? {}; |
|||
if (typeof delta.content === "string" && delta.content.length > 0) { |
|||
events.push({ type: "text.delta", text: delta.content }); |
|||
} |
|||
|
|||
if (Array.isArray(delta.tool_calls)) { |
|||
for (const toolDelta of delta.tool_calls) { |
|||
const index = toolDelta.index ?? 0; |
|||
const existing = state.toolCallsByIndex.get(index) ?? { |
|||
id: toolDelta.id ?? `tool_call_${index}`, |
|||
name: "", |
|||
arguments: "", |
|||
}; |
|||
|
|||
if (typeof toolDelta.function?.name === "string") { |
|||
existing.name = toolDelta.function.name; |
|||
} |
|||
if (typeof toolDelta.function?.arguments === "string") { |
|||
existing.arguments += toolDelta.function.arguments; |
|||
} |
|||
if (typeof toolDelta.id === "string") { |
|||
existing.id = toolDelta.id; |
|||
} |
|||
|
|||
state.toolCallsByIndex.set(index, existing); |
|||
events.push({ |
|||
type: "tool_call.delta", |
|||
id: existing.id, |
|||
name: existing.name || undefined, |
|||
argumentsDelta: |
|||
typeof toolDelta.function?.arguments === "string" |
|||
? toolDelta.function.arguments |
|||
: undefined, |
|||
}); |
|||
} |
|||
} |
|||
|
|||
if (payload?.usage) { |
|||
events.push({ type: "response.usage", usage: toUsage(payload.usage) ?? {} }); |
|||
} |
|||
|
|||
if (choice.finish_reason) { |
|||
for (const [, toolCall] of state.toolCallsByIndex) { |
|||
events.push({ type: "tool_call.done", toolCall }); |
|||
} |
|||
events.push({ type: "response.done", finishReason: choice.finish_reason }); |
|||
state.doneEmitted = true; |
|||
} |
|||
|
|||
return events; |
|||
}, |
|||
|
|||
normalizeError(error) { |
|||
const err = error as any; |
|||
const statusCode = err?.statusCode ?? err?.status; |
|||
const message = err?.message ?? "Provider request failed"; |
|||
if (statusCode === 401 || statusCode === 403) { |
|||
return new XllmError({ |
|||
code: "AUTH_ERROR", |
|||
message, |
|||
provider: "openai-compatible", |
|||
statusCode, |
|||
raw: err, |
|||
}); |
|||
} |
|||
if (statusCode === 429) { |
|||
return new XllmError({ |
|||
code: "RATE_LIMIT", |
|||
message, |
|||
provider: "openai-compatible", |
|||
statusCode, |
|||
raw: err, |
|||
}); |
|||
} |
|||
if (statusCode && statusCode >= 400 && statusCode < 500) { |
|||
return new XllmError({ |
|||
code: "INVALID_REQUEST", |
|||
message, |
|||
provider: "openai-compatible", |
|||
statusCode, |
|||
raw: err, |
|||
}); |
|||
} |
|||
return new XllmError({ |
|||
code: "PROVIDER_ERROR", |
|||
message, |
|||
provider: "openai-compatible", |
|||
statusCode, |
|||
raw: err, |
|||
}); |
|||
}, |
|||
}; |
|||
|
|||
export const withProviderName = ( |
|||
response: XResponse, |
|||
provider: XProviderName, |
|||
): XResponse => ({ ...response, provider }); |
|||
@ -0,0 +1,11 @@ |
|||
import type { XProviderName } from "../core/types"; |
|||
import { deepseekAdapter } from "./deepseek.adapter"; |
|||
import { openAICompatibleAdapter } from "./openai-compatible.adapter"; |
|||
import type { ProviderAdapter } from "./types"; |
|||
|
|||
const adapterByName: Record<XProviderName, ProviderAdapter> = { |
|||
"openai-compatible": openAICompatibleAdapter, |
|||
deepseek: deepseekAdapter, |
|||
}; |
|||
|
|||
export const getProviderAdapter = (provider: XProviderName): ProviderAdapter => adapterByName[provider]; |
|||
@ -0,0 +1,30 @@ |
|||
import type { XllmError } from "../core/errors"; |
|||
import type { |
|||
ResolvedConfig, |
|||
XProviderName, |
|||
XRequest, |
|||
XResponse, |
|||
XStreamEvent, |
|||
XToolCall, |
|||
} from "../core/types"; |
|||
|
|||
export interface ProviderHttpRequest { |
|||
url: string; |
|||
method: "POST"; |
|||
headers: Record<string, string>; |
|||
body: Record<string, unknown>; |
|||
} |
|||
|
|||
export interface StreamState { |
|||
started: boolean; |
|||
doneEmitted: boolean; |
|||
toolCallsByIndex: Map<number, XToolCall>; |
|||
} |
|||
|
|||
export interface ProviderAdapter { |
|||
name: XProviderName; |
|||
toProviderRequest(input: XRequest, config: ResolvedConfig, stream: boolean): ProviderHttpRequest; |
|||
fromProviderResponse(raw: unknown, provider: XProviderName): XResponse; |
|||
fromProviderStreamChunk(chunkData: unknown, state: StreamState): XStreamEvent[]; |
|||
normalizeError(error: unknown): XllmError; |
|||
} |
|||
@ -0,0 +1,63 @@ |
|||
import { XllmError } from "../core/errors"; |
|||
import type { ResolvedConfig, XClientOptions, XRequest } from "../core/types"; |
|||
|
|||
const env = (key: string): string | undefined => { |
|||
if (typeof globalThis.process !== "undefined" && globalThis.process?.env) { |
|||
return globalThis.process.env[key]; |
|||
} |
|||
return undefined; |
|||
}; |
|||
|
|||
const defaultModelByProvider = { |
|||
"openai-compatible": "gpt-4o-mini", |
|||
deepseek: "deepseek-chat", |
|||
} as const; |
|||
|
|||
const defaultBaseUrlByProvider = { |
|||
"openai-compatible": "https://api.openai.com/v1", |
|||
deepseek: "https://api.deepseek.com", |
|||
} as const; |
|||
|
|||
const resolveApiKey = (provider: "openai-compatible" | "deepseek", request?: XRequest, options?: XClientOptions) => { |
|||
const value = |
|||
request?.apiKey ?? |
|||
options?.apiKey ?? |
|||
env("XLLM_API_KEY") ?? |
|||
(provider === "deepseek" ? env("DEEPSEEK_API_KEY") : env("OPENAI_API_KEY")); |
|||
return value; |
|||
}; |
|||
|
|||
export const resolveConfig = (request: XRequest, options: XClientOptions): ResolvedConfig => { |
|||
const provider = request.provider ?? options.provider ?? "openai-compatible"; |
|||
const model = request.model ?? options.model ?? defaultModelByProvider[provider]; |
|||
const apiKey = resolveApiKey(provider, request, options); |
|||
const baseURL = |
|||
request.baseURL ?? options.baseURL ?? env("XLLM_BASE_URL") ?? defaultBaseUrlByProvider[provider]; |
|||
const fetchImpl = options.fetch ?? globalThis.fetch; |
|||
const headers = options.headers ?? {}; |
|||
|
|||
if (!fetchImpl) { |
|||
throw new XllmError({ |
|||
code: "NETWORK_ERROR", |
|||
message: "No fetch implementation available", |
|||
provider, |
|||
}); |
|||
} |
|||
|
|||
if (!apiKey) { |
|||
throw new XllmError({ |
|||
code: "AUTH_ERROR", |
|||
message: `Missing apiKey for provider: ${provider}`, |
|||
provider, |
|||
}); |
|||
} |
|||
|
|||
return { |
|||
provider, |
|||
model, |
|||
apiKey, |
|||
baseURL, |
|||
fetchImpl, |
|||
headers, |
|||
}; |
|||
}; |
|||
@ -0,0 +1,45 @@ |
|||
import { XllmError } from "../core/errors"; |
|||
import type { ProviderAdapter, ProviderHttpRequest } from "../providers/types"; |
|||
|
|||
const toNetworkError = (error: unknown, provider: string): XllmError => { |
|||
if (error instanceof XllmError) return error; |
|||
const err = error as Error; |
|||
return new XllmError({ |
|||
code: "NETWORK_ERROR", |
|||
message: err?.message || "Network request failed", |
|||
provider, |
|||
raw: error, |
|||
}); |
|||
}; |
|||
|
|||
export const postJSON = async ( |
|||
fetchImpl: typeof fetch, |
|||
request: ProviderHttpRequest, |
|||
adapter: ProviderAdapter, |
|||
): Promise<Response> => { |
|||
try { |
|||
return await fetchImpl(request.url, { |
|||
method: request.method, |
|||
headers: request.headers, |
|||
body: JSON.stringify(request.body), |
|||
}); |
|||
} catch (error) { |
|||
throw toNetworkError(error, adapter.name); |
|||
} |
|||
}; |
|||
|
|||
export const throwForBadStatus = async (response: Response, adapter: ProviderAdapter): Promise<void> => { |
|||
if (response.ok) return; |
|||
let raw: unknown = undefined; |
|||
const textBody = await response.text(); |
|||
try { |
|||
raw = JSON.parse(textBody); |
|||
} catch { |
|||
raw = textBody; |
|||
} |
|||
throw adapter.normalizeError({ |
|||
statusCode: response.status, |
|||
message: (raw as any)?.error?.message ?? `HTTP ${response.status}`, |
|||
raw, |
|||
}); |
|||
}; |
|||
@ -0,0 +1,32 @@ |
|||
export async function* parseSSE(stream: ReadableStream<Uint8Array>): AsyncGenerator<string> { |
|||
const reader = stream.getReader(); |
|||
const decoder = new TextDecoder(); |
|||
let buffer = ""; |
|||
|
|||
try { |
|||
while (true) { |
|||
const result = await reader.read(); |
|||
if (result.done) break; |
|||
buffer += decoder.decode(result.value, { stream: true }).replaceAll("\r\n", "\n"); |
|||
|
|||
let separatorIndex = buffer.indexOf("\n\n"); |
|||
while (separatorIndex >= 0) { |
|||
const frame = buffer.slice(0, separatorIndex); |
|||
buffer = buffer.slice(separatorIndex + 2); |
|||
separatorIndex = buffer.indexOf("\n\n"); |
|||
|
|||
const lines = frame |
|||
.split("\n") |
|||
.map((line) => line.trim()) |
|||
.filter((line) => line.startsWith("data:")); |
|||
for (const line of lines) { |
|||
const payload = line.slice(5).trim(); |
|||
if (!payload) continue; |
|||
yield payload; |
|||
} |
|||
} |
|||
} |
|||
} finally { |
|||
reader.releaseLock(); |
|||
} |
|||
} |
|||
@ -0,0 +1,15 @@ |
|||
import { defineProject } from "vitest/config"; |
|||
import { dirname } from "node:path"; |
|||
import { fileURLToPath } from "node:url"; |
|||
|
|||
const root = dirname(fileURLToPath(import.meta.url)); |
|||
|
|||
export default defineProject({ |
|||
root, |
|||
test: { |
|||
name: "xllm", |
|||
exclude: [], |
|||
include: ["src/**/*.test.ts"], |
|||
environment: "node", |
|||
}, |
|||
}); |
|||
Loading…
Reference in new issue