From 695000b710fbd33c48a11bdc5c5a22c8b72c0eee Mon Sep 17 00:00:00 2001 From: dash <1549469775@qq.com> Date: Sun, 26 Apr 2026 13:16:14 +0800 Subject: [PATCH] feat: add xllm package with streaming provider adapters 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: Cursor --- docs/superpowers/specs/2026-04-26-xllm-design.md | 236 +++++++++++++++++++++ packages/example/package.json | 1 + packages/example/src/index.ts | 47 ++++ packages/xllm/README.md | 25 +++ packages/xllm/package.json | 23 ++ packages/xllm/src/client/create-xllm.ts | 21 ++ packages/xllm/src/client/generate.ts | 17 ++ packages/xllm/src/client/stream.ts | 53 +++++ packages/xllm/src/core/errors.ts | 31 +++ packages/xllm/src/core/types.ts | 89 ++++++++ packages/xllm/src/index.test.ts | 90 ++++++++ packages/xllm/src/index.ts | 15 ++ packages/xllm/src/providers/deepseek.adapter.ts | 45 ++++ .../src/providers/openai-compatible.adapter.ts | 233 ++++++++++++++++++++ packages/xllm/src/providers/registry.ts | 11 + packages/xllm/src/providers/types.ts | 30 +++ packages/xllm/src/runtime/config.ts | 63 ++++++ packages/xllm/src/runtime/http.ts | 45 ++++ packages/xllm/src/runtime/sse.ts | 32 +++ packages/xllm/vitest.config.ts | 15 ++ 20 files changed, 1122 insertions(+) create mode 100644 docs/superpowers/specs/2026-04-26-xllm-design.md create mode 100644 packages/xllm/README.md create mode 100644 packages/xllm/package.json create mode 100644 packages/xllm/src/client/create-xllm.ts create mode 100644 packages/xllm/src/client/generate.ts create mode 100644 packages/xllm/src/client/stream.ts create mode 100644 packages/xllm/src/core/errors.ts create mode 100644 packages/xllm/src/core/types.ts create mode 100644 packages/xllm/src/index.test.ts create mode 100644 packages/xllm/src/index.ts create mode 100644 packages/xllm/src/providers/deepseek.adapter.ts create mode 100644 packages/xllm/src/providers/openai-compatible.adapter.ts create mode 100644 packages/xllm/src/providers/registry.ts create mode 100644 packages/xllm/src/providers/types.ts create mode 100644 packages/xllm/src/runtime/config.ts create mode 100644 packages/xllm/src/runtime/http.ts create mode 100644 packages/xllm/src/runtime/sse.ts create mode 100644 packages/xllm/vitest.config.ts diff --git a/docs/superpowers/specs/2026-04-26-xllm-design.md b/docs/superpowers/specs/2026-04-26-xllm-design.md new file mode 100644 index 0000000..9cfdc3f --- /dev/null +++ b/docs/superpowers/specs/2026-04-26-xllm-design.md @@ -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` +- 请求级覆盖字段: + - `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`,统一事件集合: + +- `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 切换 +- 非流式返回结构与流式聚合结果语义一致 +- 关键错误能归一化并携带可诊断信息 diff --git a/packages/example/package.json b/packages/example/package.json index 3c85079..47141dc 100644 --- a/packages/example/package.json +++ b/packages/example/package.json @@ -8,6 +8,7 @@ "@dm/crypto-wasm": "workspace:*", "@dm/core": "workspace:*", "@dm/dx": "workspace:*", + "@dm/xllm": "workspace:*", "crypto-js": "4.2.0", "sm-crypto": "0.4.0" }, diff --git a/packages/example/src/index.ts b/packages/example/src/index.ts index e41bdb8..0ba0e00 100644 --- a/packages/example/src/index.ts +++ b/packages/example/src/index.ts @@ -1,4 +1,5 @@ import { decrypt, encrypt } from "@dm/crypto-wasm"; +import { createXllm, type XProviderName } from "@dm/xllm"; import { encrypt as sm4Encrypt } from "./sm4"; /** 与 sm-crypto@0.4.0 + crypto-js@4.2.0(ECB/PKCS#7、密钥见 sm4.js)对同一段 UTF-8 明文应得到相同 Base64。 */ @@ -15,3 +16,49 @@ console.log("[hello-wasm] validateText === cipherJs:", validateText === cipherJs const recovered = decrypt(cipherWasm); console.log("[hello-wasm] round-trip OK:", recovered === plain); + +const env = import.meta.env as Record; +const xllmProvider = (env.VITE_XLLM_PROVIDER as XProviderName | undefined) ?? "deepseek"; +const xllmModel = env.VITE_XLLM_MODEL ?? "deepseek-chat"; +const xllmApiKey = env.VITE_XLLM_API_KEY; +const xllmBaseUrl = env.VITE_XLLM_BASE_URL; + +async function runXllmStreamExample(): Promise { + if (!xllmApiKey) { + console.log( + "[xllm-example] skip: set VITE_XLLM_API_KEY to run stream demo. Optional: VITE_XLLM_PROVIDER, VITE_XLLM_MODEL, VITE_XLLM_BASE_URL", + ); + return; + } + + const xllm = createXllm({ + provider: xllmProvider, + model: xllmModel, + apiKey: xllmApiKey, + baseURL: xllmBaseUrl, + }); + + console.log(`[xllm-example] provider=${xllmProvider}, model=${xllmModel}`); + let streamedText = ""; + + for await (const event of xllm.stream({ + messages: [ + { + role: "user", + content: [{ type: "text", text: "请用一句话介绍你自己。" }], + }, + ], + })) { + if (event.type === "text.delta") { + streamedText += event.text; + console.log("[xllm-example] delta:", event.text); + } + if (event.type === "response.done") { + console.log("[xllm-example] done:", streamedText); + } + } +} + +runXllmStreamExample().catch((error) => { + console.error("[xllm-example] failed:", error); +}); diff --git a/packages/xllm/README.md b/packages/xllm/README.md new file mode 100644 index 0000000..01f80d3 --- /dev/null +++ b/packages/xllm/README.md @@ -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); +} +``` diff --git a/packages/xllm/package.json b/packages/xllm/package.json new file mode 100644 index 0000000..c997eea --- /dev/null +++ b/packages/xllm/package.json @@ -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" + } +} diff --git a/packages/xllm/src/client/create-xllm.ts b/packages/xllm/src/client/create-xllm.ts new file mode 100644 index 0000000..907a3be --- /dev/null +++ b/packages/xllm/src/client/create-xllm.ts @@ -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; + generate(request: XRequest): Promise; + with(overrides: Partial): XllmClient; +} + +export const createXllm = (options: XClientOptions = {}): XllmClient => ({ + stream(request: XRequest) { + return stream(options, request); + }, + generate(request: XRequest) { + return generate(options, request); + }, + with(overrides: Partial) { + return createXllm({ ...options, ...overrides }); + }, +}); diff --git a/packages/xllm/src/client/generate.ts b/packages/xllm/src/client/generate.ts new file mode 100644 index 0000000..be83f7a --- /dev/null +++ b/packages/xllm/src/client/generate.ts @@ -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 => { + 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 }; +}; diff --git a/packages/xllm/src/client/stream.ts b/packages/xllm/src/client/stream.ts new file mode 100644 index 0000000..f1403f6 --- /dev/null +++ b/packages/xllm/src/client/stream.ts @@ -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 { + 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; + } + } +}; diff --git a/packages/xllm/src/core/errors.ts b/packages/xllm/src/core/errors.ts new file mode 100644 index 0000000..3de828f --- /dev/null +++ b/packages/xllm/src/core/errors.ts @@ -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; + } +} diff --git a/packages/xllm/src/core/types.ts b/packages/xllm/src/core/types.ts new file mode 100644 index 0000000..d9e8756 --- /dev/null +++ b/packages/xllm/src/core/types.ts @@ -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; +} + +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; + 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; +} + +export interface ResolvedConfig { + provider: XProviderName; + model: string; + apiKey: string; + baseURL: string; + fetchImpl: typeof fetch; + headers: Record; +} diff --git a/packages/xllm/src/index.test.ts b/packages/xllm/src/index.test.ts new file mode 100644 index 0000000..e764414 --- /dev/null +++ b/packages/xllm/src/index.test.ts @@ -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 => { + const encoder = new TextEncoder(); + return new ReadableStream({ + 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 => event.type === "text.delta") + .map((event) => event.text) + .join(""); + expect(text).toBe("hello"); + expect(events.some((event) => event.type === "response.done")).toBe(true); + }); +}); diff --git a/packages/xllm/src/index.ts b/packages/xllm/src/index.ts new file mode 100644 index 0000000..53145dd --- /dev/null +++ b/packages/xllm/src/index.ts @@ -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"; diff --git a/packages/xllm/src/providers/deepseek.adapter.ts b/packages/xllm/src/providers/deepseek.adapter.ts new file mode 100644 index 0000000..012a198 --- /dev/null +++ b/packages/xllm/src/providers/deepseek.adapter.ts @@ -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, + }); + }, +}; diff --git a/packages/xllm/src/providers/openai-compatible.adapter.ts b/packages/xllm/src/providers/openai-compatible.adapter.ts new file mode 100644 index 0000000..d606ccc --- /dev/null +++ b/packages/xllm/src/providers/openai-compatible.adapter.ts @@ -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> => { + 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 => { + 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 = { + 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 }); diff --git a/packages/xllm/src/providers/registry.ts b/packages/xllm/src/providers/registry.ts new file mode 100644 index 0000000..42d6900 --- /dev/null +++ b/packages/xllm/src/providers/registry.ts @@ -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 = { + "openai-compatible": openAICompatibleAdapter, + deepseek: deepseekAdapter, +}; + +export const getProviderAdapter = (provider: XProviderName): ProviderAdapter => adapterByName[provider]; diff --git a/packages/xllm/src/providers/types.ts b/packages/xllm/src/providers/types.ts new file mode 100644 index 0000000..708722d --- /dev/null +++ b/packages/xllm/src/providers/types.ts @@ -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; + body: Record; +} + +export interface StreamState { + started: boolean; + doneEmitted: boolean; + toolCallsByIndex: Map; +} + +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; +} diff --git a/packages/xllm/src/runtime/config.ts b/packages/xllm/src/runtime/config.ts new file mode 100644 index 0000000..4bd0f32 --- /dev/null +++ b/packages/xllm/src/runtime/config.ts @@ -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, + }; +}; diff --git a/packages/xllm/src/runtime/http.ts b/packages/xllm/src/runtime/http.ts new file mode 100644 index 0000000..9eebee9 --- /dev/null +++ b/packages/xllm/src/runtime/http.ts @@ -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 => { + 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 => { + 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, + }); +}; diff --git a/packages/xllm/src/runtime/sse.ts b/packages/xllm/src/runtime/sse.ts new file mode 100644 index 0000000..c56b039 --- /dev/null +++ b/packages/xllm/src/runtime/sse.ts @@ -0,0 +1,32 @@ +export async function* parseSSE(stream: ReadableStream): AsyncGenerator { + 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(); + } +} diff --git a/packages/xllm/vitest.config.ts b/packages/xllm/vitest.config.ts new file mode 100644 index 0000000..30bcf01 --- /dev/null +++ b/packages/xllm/vitest.config.ts @@ -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", + }, +});