From 3ca994d2645f63e31af44d5f6bb37ef4151b52c9 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E6=B0=B4=E5=8C=96?= <168127639+suisanka@users.noreply.github.com> Date: Sun, 2 Aug 2026 05:09:15 +0800 Subject: [PATCH 1/6] refactor: replace event streams with parser --- README.md | 105 ++------ examples/echo.ts | 41 --- package.json | 19 +- pnpm-lock.yaml | 26 -- scripts/check-build.ts | 59 +++++ scripts/fetch-types.sh | 1 + scripts/split-zod.ts | 179 +++++++++++++ src/client/endpoint.ts | 11 - src/client/fetch.ts | 2 +- src/event.ts | 60 +++++ src/events/index.ts | 256 ------------------- src/events/internal.ts | 331 ------------------------ src/events/source.ts | 176 ------------- src/index.ts | 2 +- src/types.ts | 2 +- tests/client.test.ts | 209 +--------------- tests/event.test.ts | 119 +++++++++ tests/events.coverage.test.ts | 287 --------------------- tests/events.logic.test.ts | 304 ---------------------- tests/events.test.ts | 458 ---------------------------------- tests/fetch.test.ts | 4 +- tests/helpers/async.ts | 49 ---- tests/helpers/transports.ts | 89 ------- tests/schema-split.test.ts | 21 ++ tsdown.config.ts | 8 +- 25 files changed, 475 insertions(+), 2343 deletions(-) delete mode 100644 examples/echo.ts create mode 100644 scripts/check-build.ts create mode 100644 scripts/split-zod.ts create mode 100644 src/event.ts delete mode 100644 src/events/index.ts delete mode 100644 src/events/internal.ts delete mode 100644 src/events/source.ts create mode 100644 tests/event.test.ts delete mode 100644 tests/events.coverage.test.ts delete mode 100644 tests/events.logic.test.ts delete mode 100644 tests/events.test.ts delete mode 100644 tests/helpers/transports.ts create mode 100644 tests/schema-split.test.ts diff --git a/README.md b/README.md index 4b02c8e..e632474 100644 --- a/README.md +++ b/README.md @@ -3,18 +3,12 @@ [![CI](https://github.com/SaltifyDev/milky-tea/actions/workflows/ci.yml/badge.svg?branch=main)](https://github.com/SaltifyDev/milky-tea/actions/workflows/ci.yml) [![Coverage](https://img.shields.io/endpoint?url=https%3A%2F%2Fraw.githubusercontent.com%2FSaltifyDev%2Fmilky-tea%2Fbadges%2Fcoverage-badge.json)](https://github.com/SaltifyDev/milky-tea/actions/workflows/ci.yml) -Milky 的 TypeScript SDK,提供类型安全的 API 调用和事件流支持。 +Milky 的 TypeScript SDK,提供类型安全的 API 调用和事件解析。 ## 安装 ```bash -npm i @saltify/milky-tea @saltify/milky-types -``` - -如果运行环境不支持 EventSource(例如 Node.js 环境)且需要 SSE 支持,则需要安装 `eventsource`: - -```bash -npm i eventsource +pnpm add @saltify/milky-tea zod ``` ## 使用方法 @@ -46,99 +40,28 @@ await client.group.quitGroup({ group_id: 10001 }, { timeout: false }) 在这里,第二个参数是可选的,可以覆盖默认的 `baseURL`、`token`、`timeout` 等设置。 -### 监听事件 +### 解析事件 -通过 `client.event()` 创建一个事件连接,支持 WebSocket 和 SSE 两种连接方式。连接模式有如下几种: - -- `websocket`:仅使用 WebSocket -- `sse`:仅使用 Server-Sent Events -- `auto`:兼容旧版本的保留值,不再支持,传入后会报错;请显式使用 `websocket` 或 `sse` +SDK 不负责创建或管理事件连接。通过 SSE、WebSocket、WebHook 或其他方式收到事件后,将反序列化后的对象传给 `resolveMilkyEvent`: ```ts -const source = client.event('websocket', { - reconnect: { - interval: 1000, - attempts: 'always', - }, -}) +import { resolveMilkyEvent } from '@saltify/milky-tea/event' -// 监听连接打开 -source.on('open', () => { - console.log('connected') -}) - -// 监听所有事件 -source.on('push', (event) => { - console.log(event.event_type, event) -}) - -// 监听特定类型的事件 -source.on('foobar', (event) => { - console.log(event.message.content) -}) - -// 监听错误 -source.on('error', (event) => { - console.error(event.message) -}) +const event = await resolveMilkyEvent(JSON.parse(payload)) -// 使用 async iteration -for await (const event of source) { - console.log(event.event_type) - if (shouldStop) +switch (event.event_type) { + case 'message_receive': + console.log(event.data) + break + case 'bot_offline': + console.log(event.data.reason) break } - -source.close() ``` -**注意**: 事件对象是深度只读的(immutable),所有嵌套属性都被冻结,无法修改。 - -### `createMilkyEventSource` - -如果需要更底层的事件源控制,可以使用 `createMilkyEventSource` 直接创建事件源。 - -```ts -import { createMilkyEventSource } from '@saltify/milky-tea' - -// 使用连接类型和选项 -const source = createMilkyEventSource('websocket', { - baseURL: 'https://milky.example.com', - token: process.env.MILKY_TOKEN, - timeout: 15000, - reconnect: { - interval: 1000, - attempts: 5, - }, -}) - -// 或使用自定义传输工厂 -const source = createMilkyEventSource( - async (options, signal) => { - // 返回 WebSocket 或 EventSource 实例 - return new WebSocket('wss://milky.example.com/event') - }, - { - timeout: 10000, - }, -) - -source.on('open', () => console.log('Connected')) -source.on('push', event => console.log(event)) -source.close() -``` +也可以从包根入口导入。推荐使用 `@saltify/milky-tea/event`,以便打包器完全隔离客户端代码和 API schema。 -**参数**: - -- `kind`: 连接类型 (`'websocket'` | `'sse'`;`'auto'` 为兼容保留值,传入会报错) -- `factory`: 自定义传输工厂函数 -- `options`: - - `baseURL`: 服务器地址(使用 kind 时必需) - - `token`: 访问令牌 - - `timeout`: 连接超时时间(默认 15000ms) - - `reconnect`: 重连配置 - - `interval`: 重连间隔(毫秒) - - `attempts`: 重连次数(`'always'` 或数字) +`resolveMilkyEvent` 使用生成的 Zod schema 校验输入。校验结果会移除未知字段并返回深拷贝,但不会冻结返回对象;校验失败时会抛出带有 Zod 错误原因的异常。 ### `createMilkyFetch` diff --git a/examples/echo.ts b/examples/echo.ts deleted file mode 100644 index f4c5e9a..0000000 --- a/examples/echo.ts +++ /dev/null @@ -1,41 +0,0 @@ -/* eslint-disable no-console */ -import { createMilkyClient } from '@/index' - -const milky = createMilkyClient({ - baseURL: 'http://127.0.0.1:3100', -}) - -const source = milky.event() - -source.on('open', () => console.log('Milky connected')) - -source.on('message_receive', async ({ self_id, data }) => { - if (data.segments.length < 2 || data.message_scene !== 'group' || data.sender_id === self_id) { - return - } - - const [first, second] = data.segments - if (first.type !== 'mention' || first.data.user_id !== self_id || second.type !== 'text') { - return - } - - const content = second.data.text.trim() - - console.log(`Mentioned by ${data.group_member.nickname} (${data.group_member.user_id}), content: ${content}`) - - await milky.message.sendGroupMessage({ - group_id: data.peer_id, - message: [ - { - type: 'reply', - data: { message_seq: data.message_seq }, - }, - { - type: 'text', - data: { - text: content, - }, - }, - ], - }) -}) diff --git a/package.json b/package.json index e735a53..0d590b8 100644 --- a/package.json +++ b/package.json @@ -10,8 +10,16 @@ "type": "git", "url": "git+https://github.com/SaltifyDev/milky-tea.git" }, + "sideEffects": false, "exports": { - ".": "./dist/index.mjs", + ".": { + "types": "./dist/index.d.mts", + "import": "./dist/index.mjs" + }, + "./event": { + "types": "./dist/event.d.mts", + "import": "./dist/event.mjs" + }, "./package.json": "./package.json" }, "types": "./dist/index.d.mts", @@ -24,6 +32,7 @@ "generate-api": "sh ./scripts/fetch-types.sh", "test": "vitest run", "test:coverage": "vitest run --coverage", + "check:bundle": "pnpm run build && pnpm exec jiti ./scripts/check-build.ts", "typecheck": "tsc --noEmit", "prepublishOnly": "pnpm run build", "fmt": "eslint . --fix", @@ -31,20 +40,13 @@ "bump": "bumpp" }, "peerDependencies": { - "eventsource": "^4.1.0", "zod": "^4.4.3" }, "peerDependenciesMeta": { - "eventsource": { - "optional": true - }, "zod": { "optional": true } }, - "dependencies": { - "mitt": "^3.0.1" - }, "devDependencies": { "@antfu/eslint-config": "^7.7.3", "@types/node": "^26.1.2", @@ -53,7 +55,6 @@ "bumpp": "^10.4.1", "change-case": "^5.4.4", "eslint": "^10.8.0", - "eventsource": "^4.1.0", "jiti": "^2.7.0", "tsdown": "^0.21.10", "typescript": "^5.9.3", diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index b079d33..bde9ae5 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -7,10 +7,6 @@ settings: importers: .: - dependencies: - mitt: - specifier: ^3.0.1 - version: 3.0.1 devDependencies: '@antfu/eslint-config': specifier: ^7.7.3 @@ -33,9 +29,6 @@ importers: eslint: specifier: ^10.8.0 version: 10.8.0(jiti@2.7.0) - eventsource: - specifier: ^4.1.0 - version: 4.1.0 jiti: specifier: ^2.7.0 version: 2.7.0 @@ -1227,14 +1220,6 @@ packages: resolution: {integrity: sha512-kVscqXk4OCp68SZ0dkgEKVi6/8ij300KBWTJq32P/dYeWTSwK41WyTxalN1eRmA5Z9UU/LX9D7FWSmV9SAYx6g==} engines: {node: '>=0.10.0'} - eventsource-parser@3.1.0: - resolution: {integrity: sha512-kJezFj9YFAMLeORyi7aCLxLbD5/qWMQnoMVlVPyHIll7lgRJCc3JVln9Vgl9nwQi0YkMnhdGTMNn7CkRRAptMg==} - engines: {node: '>=18.0.0'} - - eventsource@4.1.0: - resolution: {integrity: sha512-2GuF51iuHX6A9xdTccMTsNb7VO0lHZihApxhvQzJB5A03DvHDd2FQepodbMaztPBmBcE/ox7o2gqaxGhYB9LhQ==} - engines: {node: '>=20.0.0'} - expect-type@1.4.0: resolution: {integrity: sha512-KfYbmpRm0VbLjEvVa9yGwCi9GI34xvi7A/HXYWQO65CSD2u3MczUJSuwXKFIxlGsgBQizV9q5J9NHj4VG0n+pA==} engines: {node: '>=12.0.0'} @@ -1693,9 +1678,6 @@ packages: resolution: {integrity: sha512-tEBHqDnIoM/1rXME1zgka9g6Q2lcoCkxHLuc7ODJ5BxbP5d4c2Z5cGgtXAku59200Cx7diuHTOYfSBD8n6mm8A==} engines: {node: '>=16 || 14 >=14.17'} - mitt@3.0.1: - resolution: {integrity: sha512-vKivATfr97l2/QBCYAkXYDbrIWPM2IIKEl7YPhjCvKlG3kE2gm+uBo6nEXK3M5/Ffh/FLpKExzOQ3JJoJGFKBw==} - mlly@1.8.2: resolution: {integrity: sha512-d+ObxMQFmbt10sretNDytwt85VrbkhhUA/JBGm1MPaWJ65Cl4wOgLaB1NYvJSZ0Ef03MMEU/0xpPMXUIQ29UfA==} @@ -3381,12 +3363,6 @@ snapshots: esutils@2.0.3: {} - eventsource-parser@3.1.0: {} - - eventsource@4.1.0: - dependencies: - eventsource-parser: 3.1.0 - expect-type@1.4.0: {} exsolve@1.1.1: {} @@ -3960,8 +3936,6 @@ snapshots: minipass@7.1.3: {} - mitt@3.0.1: {} - mlly@1.8.2: dependencies: acorn: 8.18.0 diff --git a/scripts/check-build.ts b/scripts/check-build.ts new file mode 100644 index 0000000..47a4fb0 --- /dev/null +++ b/scripts/check-build.ts @@ -0,0 +1,59 @@ +import { readdirSync, readFileSync } from 'node:fs' +import { basename, join } from 'node:path' +import process from 'node:process' + +const distPath = './dist' +const modules = new Map( + readdirSync(distPath) + .filter(file => file.endsWith('.mjs')) + .map(file => [file, readFileSync(join(distPath, file), 'utf8')]), +) + +function requireModule(name: string): string { + const contents = modules.get(name) + if (contents == null) { + throw new Error(`Missing build module ${name}`) + } + return contents +} + +function findModule(fragment: string): [string, string] { + const matches = [...modules].filter(([, contents]) => contents.includes(fragment)) + if (matches.length !== 1) { + throw new Error(`Expected exactly one build module containing ${JSON.stringify(fragment)}, found ${matches.length}`) + } + return matches[0]! +} + +function assert(condition: unknown, message: string): asserts condition { + if (!condition) { + throw new Error(message) + } +} + +const indexEntry = requireModule('index.mjs') +const eventEntry = requireModule('event.mjs') +const [apiSchemaName, apiSchema] = findModule('const zodApiCategories =') +const [eventSchemaName, eventSchema] = findModule('const Event = z.discriminatedUnion("event_type"') +const [commonSchemaName, commonSchema] = findModule('const IncomingMessage = z.discriminatedUnion("message_scene"') + +assert(eventEntry.includes('import("./zod-event-'), 'event entry must load its schema dynamically') +assert(!eventEntry.includes('endpointNamesByCategory'), 'event entry must not include client metadata') +assert(!eventEntry.includes('zodApiCategories'), 'event entry must not include API schemas') +assert(indexEntry.includes('import("./zod-api-'), 'client entry must load API schemas dynamically') +assert(!indexEntry.includes('from "zod"'), 'root entry must not statically load zod') +assert(!eventSchema.includes('zodApiCategories'), 'event schema chunk must not include API schemas') +assert(!apiSchema.includes('discriminatedUnion("event_type"'), 'API schema chunk must not include the event union') +assert(commonSchema.includes('from "zod"'), 'common schema chunk must own its zod dependency') + +for (const [name, contents] of modules) { + assert(!contents.includes('from "mitt"'), `${name} must not depend on mitt`) + assert(!contents.includes('from "eventsource"'), `${name} must not depend on eventsource`) +} + +process.stdout.write(`${[ + `event entry: ${basename('event.mjs')}`, + `event schema: ${eventSchemaName}`, + `API schema: ${apiSchemaName}`, + `shared schema: ${commonSchemaName}`, +].join('\n')}\n`) diff --git a/scripts/fetch-types.sh b/scripts/fetch-types.sh index e031ebc..33a08be 100644 --- a/scripts/fetch-types.sh +++ b/scripts/fetch-types.sh @@ -8,3 +8,4 @@ echo; echo 'Fetching zod types...' curl -o './src/gen/zod.ts' 'https://milky-next.ntqqrev.org/raw/typescript/zod.txt' pnpm exec jiti ./scripts/generate-meta.ts +pnpm exec jiti ./scripts/split-zod.ts diff --git a/scripts/split-zod.ts b/scripts/split-zod.ts new file mode 100644 index 0000000..87174d3 --- /dev/null +++ b/scripts/split-zod.ts @@ -0,0 +1,179 @@ +import { readFileSync, writeFileSync } from 'node:fs' +import ts from 'typescript' + +const sourcePath = './src/gen/zod.ts' +const commonOutputPath = './src/gen/zod-common.ts' +const eventOutputPath = './src/gen/zod-event.ts' +const apiOutputPath = './src/gen/zod-api.ts' + +const sourceText = readFileSync(sourcePath, 'utf8') +const sourceFile = ts.createSourceFile( + sourcePath, + sourceText, + ts.ScriptTarget.Latest, + true, + ts.ScriptKind.TS, +) + +interface RuntimeDeclaration { + readonly names: readonly string[] + readonly statement: ts.Statement + readonly dependencies: Set +} + +function getRuntimeDeclarationNames(statement: ts.Statement): string[] { + if (ts.isFunctionDeclaration(statement) && statement.name) { + return [statement.name.text] + } + + if (!ts.isVariableStatement(statement)) { + return [] + } + + return statement.declarationList.declarations.flatMap((declaration) => { + if (!ts.isIdentifier(declaration.name)) { + throw new TypeError(`Unsupported generated declaration in ${sourcePath}`) + } + + return declaration.name.text + }) +} + +const runtimeDeclarations: RuntimeDeclaration[] = [] +const declarationByName = new Map() + +for (const statement of sourceFile.statements) { + const names = getRuntimeDeclarationNames(statement) + if (names.length === 0) { + continue + } + + const declaration: RuntimeDeclaration = { + names, + statement, + dependencies: new Set(), + } + + runtimeDeclarations.push(declaration) + for (const name of names) { + if (declarationByName.has(name)) { + throw new TypeError(`Duplicate generated declaration ${name} in ${sourcePath}`) + } + declarationByName.set(name, declaration) + } +} + +function collectDependencies(declaration: RuntimeDeclaration): void { + function visit(node: ts.Node): void { + if (ts.isIdentifier(node)) { + const dependency = declarationByName.get(node.text) + if (dependency && dependency !== declaration) { + declaration.dependencies.add(node.text) + } + } + + ts.forEachChild(node, visit) + } + + visit(declaration.statement) +} + +for (const declaration of runtimeDeclarations) { + collectDependencies(declaration) +} + +function resolveClosure(rootName: string): Set { + const root = declarationByName.get(rootName) + if (!root) { + throw new TypeError(`Missing generated declaration ${rootName} in ${sourcePath}`) + } + + const closure = new Set() + const pending = [root] + + while (pending.length > 0) { + const declaration = pending.pop()! + if (closure.has(declaration)) { + continue + } + + closure.add(declaration) + for (const dependencyName of declaration.dependencies) { + pending.push(declarationByName.get(dependencyName)!) + } + } + + return closure +} + +const eventClosure = resolveClosure('Event') +const apiClosure = resolveClosure('zodApiCategories') +const commonDeclarations = new Set( + [...eventClosure].filter(declaration => apiClosure.has(declaration)), +) +const eventDeclarations = new Set( + [...eventClosure].filter(declaration => !commonDeclarations.has(declaration)), +) +const apiDeclarations = new Set( + [...apiClosure].filter(declaration => !commonDeclarations.has(declaration)), +) + +function collectExternalDependencies( + declarations: ReadonlySet, + externalDeclarations: ReadonlySet, +): string[] { + const dependencies = new Set() + + for (const declaration of declarations) { + for (const dependencyName of declaration.dependencies) { + const dependency = declarationByName.get(dependencyName)! + if (externalDeclarations.has(dependency)) { + dependencies.add(dependencyName) + } + } + } + + return [...dependencies].sort() +} + +function renderModule( + declarations: ReadonlySet, + commonImports: readonly string[] = [], +): string { + const lines = [ + '// Generated by scripts/split-zod.ts. Do not edit.', + 'import { z } from \'zod\'', + ] + + if (commonImports.length > 0) { + lines.push(`import { ${commonImports.join(', ')} } from './zod-common'`) + } + + lines.push('') + + for (const declaration of runtimeDeclarations) { + if (!declarations.has(declaration)) { + continue + } + + lines.push(sourceText.slice(declaration.statement.getFullStart(), declaration.statement.end).trim(), '') + } + + return lines.join('\n') +} + +writeFileSync(commonOutputPath, renderModule(commonDeclarations)) +writeFileSync( + eventOutputPath, + renderModule( + eventDeclarations, + collectExternalDependencies(eventDeclarations, commonDeclarations), + ), +) +writeFileSync( + apiOutputPath, + renderModule( + apiDeclarations, + collectExternalDependencies(apiDeclarations, commonDeclarations), + ), +) diff --git a/src/client/endpoint.ts b/src/client/endpoint.ts index 4a025c0..41784fd 100644 --- a/src/client/endpoint.ts +++ b/src/client/endpoint.ts @@ -1,24 +1,14 @@ import type { MilkyFetch, MilkyFetchCreateOptions, MilkyFetchOptions } from '@/client/fetch' -import type { MilkyEventSource, MilkyEventSourceOptions } from '@/events' -import type { MilkyEventSourceConnectionKind } from '@/events/source' import type { MilkyApiCategories, MilkyCamelCase, MilkyClientEndpointNames, MilkyRawEndpoints } from '@/types' import { createMilkyFetch } from '@/client/fetch' -import { createMilkyEventSource } from '@/events' import { clientEndpointNames } from '@/types' function createProxy(options: MilkyFetchCreateOptions): any { const milkyFetch = createMilkyFetch(options) - const event = (kind?: MilkyEventSourceConnectionKind, eventOptions?: MilkyEventSourceOptions) => - createMilkyEventSource(kind ?? 'websocket', { - ...eventOptions, - baseURL: options.baseURL, - token: eventOptions?.token ?? options.token, - }) const cachedEndpoints = new Map() return new Proxy({ fetch: milkyFetch, - event, }, { get(target, prop) { if (!Object.hasOwn(clientEndpointNames, prop)) { @@ -66,7 +56,6 @@ function createProxy(options: MilkyFetchCreateOptions): any { export type MilkyClient = { readonly fetch: MilkyFetch - readonly event: (kind?: MilkyEventSourceConnectionKind, options?: MilkyEventSourceOptions) => MilkyEventSource } & { readonly [K in keyof MilkyApiCategories]: { readonly [E in keyof MilkyApiCategories[K]['apis'] as MilkyCamelCase]: diff --git a/src/client/fetch.ts b/src/client/fetch.ts index 3594b91..f0a83e2 100644 --- a/src/client/fetch.ts +++ b/src/client/fetch.ts @@ -61,7 +61,7 @@ function isMissingZodError(error: unknown, seen = new Set()): boolean { } async function resolveMilkyProto(): Promise { - milkyProtoPromise ??= import('@/gen/zod') + milkyProtoPromise ??= import('@/gen/zod-api') .then(module => createMilkyProto(module.zodApiCategories)) .catch((error) => { milkyProtoPromise = undefined diff --git a/src/event.ts b/src/event.ts new file mode 100644 index 0000000..35084ff --- /dev/null +++ b/src/event.ts @@ -0,0 +1,60 @@ +import type { Event as MilkyEvent } from '@/gen/types' + +interface MilkyEventSchema { + safeParseAsync: (value: unknown) => Promise< + | { success: true, data: MilkyEvent } + | { success: false, error: Error } + > +} + +let milkyEventSchemaPromise: Promise | undefined + +function isMissingZodError(error: unknown, seen = new Set()): boolean { + if (seen.has(error)) { + return false + } + seen.add(error) + + const message = error instanceof Error ? error.message : String(error) + if (message.includes('zod') && ( + message.includes('Cannot find') + || message.includes('ERR_MODULE_NOT_FOUND') + || message.includes('module') + )) { + return true + } + + const cause = error != null && typeof error === 'object' && 'cause' in error + ? (error as { cause?: unknown }).cause + : undefined + + return cause != null && isMissingZodError(cause, seen) +} + +async function resolveMilkyEventSchema(): Promise { + milkyEventSchemaPromise ??= import('@/gen/zod-event') + .then(module => module.Event) + .catch((error) => { + milkyEventSchemaPromise = undefined + if (isMissingZodError(error)) { + throw new Error('milky: zod is required to resolve events', { cause: error }) + } + + throw error + }) + + return milkyEventSchemaPromise +} + +export async function resolveMilkyEvent(obj: unknown): Promise { + const schema = await resolveMilkyEventSchema() + const result = await schema.safeParseAsync(obj) + + if (!result.success) { + throw new Error('milky: failed to resolve event', { cause: result.error }) + } + + return result.data +} + +export type { Event } from '@/gen/types' diff --git a/src/events/index.ts b/src/events/index.ts deleted file mode 100644 index 11c3a85..0000000 --- a/src/events/index.ts +++ /dev/null @@ -1,256 +0,0 @@ -import type { MilkyEventSource, MilkyEventSourceController, MilkyEventSourceEventMap } from '@/events/internal' -import type { - MilkyEventSourceConnection, - MilkyEventSourceConnectionKind, - MilkyEventSourceTransport, - MilkyResolvedEventSourceConnectionKind, -} from '@/events/source' -import type { Awaitable } from '@/utils' -import { MilkyEventSourceImpl } from '@/events/internal' -import { connectEventTransport } from '@/events/source' -import { joinURL, raceWithAbort, sleepWithAbort, withTimeout } from '@/utils' - -export interface MilkyEventSourceOptions { - token?: string - timeout?: number - reconnect?: false | { - interval: number - attempts: 'always' | number - } -} - -export interface MilkyEventSourceCreateOptions extends MilkyEventSourceOptions { - baseURL: string | URL -} - -export type MilkyEventSourceTransportFactory = ( - options: MilkyEventSourceOptions, - signal?: AbortSignal, -) => Awaitable - -let eventSourceConstructorPromise: Promise | undefined - -async function resolveEventSourceConstructor(): Promise { - if (globalThis.EventSource) { - return globalThis.EventSource - } - - eventSourceConstructorPromise ??= import('eventsource') - .then(module => module.EventSource as unknown as typeof EventSource) - .catch((error) => { - eventSourceConstructorPromise = undefined - throw new TypeError('milky: EventSource is not available in current runtime, install optional peer dependency "eventsource" to enable sse', { cause: error }) - }) - - return eventSourceConstructorPromise -} - -async function createTransportByKind( - kind: MilkyResolvedEventSourceConnectionKind, - options: MilkyEventSourceCreateOptions, -): Promise { - const url = joinURL(options.baseURL, '/event') - - if (options.token) { - url.searchParams.set('access_token', options.token) - } - - switch (kind) { - case 'sse': - return new (await resolveEventSourceConstructor())(url) - case 'websocket': - return new WebSocket(url) - default: - throw new TypeError(`milky: unknown event source kind: ${String(kind)}`) - } -} - -async function createConnectionByKind( - kind: MilkyEventSourceConnectionKind, - options: MilkyEventSourceCreateOptions, -): Promise { - // @ts-expect-error auto is no longer supported - if (kind === 'auto') { - throw new TypeError('milky: event source kind "auto" is no longer supported; use "websocket" or "sse"') - } - - return connectEventTransport(await createTransportByKind(kind, options)) -} - -function createDisconnectError(): Error { - return new Error('milky: event source disconnected') -} - -type MilkyEventSourceRunResult = 'retry' | 'stop' - -export function createMilkyEventSource( - factory: MilkyEventSourceTransportFactory, - options?: MilkyEventSourceOptions, -): MilkyEventSource -export function createMilkyEventSource( - kind: MilkyEventSourceConnectionKind, - options: MilkyEventSourceCreateOptions, -): MilkyEventSource -export function createMilkyEventSource( - kindOrFactory: MilkyEventSourceConnectionKind | MilkyEventSourceTransportFactory, - options?: MilkyEventSourceCreateOptions | MilkyEventSourceOptions, -): MilkyEventSource { - if (typeof kindOrFactory !== 'function' && (options == null || !('baseURL' in options))) { - throw new TypeError('milky: baseURL is required when creating event sources by kind') - } - - const timeout = options?.timeout ?? 15000 - const reconnect = options?.reconnect ?? false - const eventOptions = options ?? {} - - async function connect(signal?: AbortSignal): Promise { - const controller = new AbortController() - let shouldCloseTransport = false - - try { - const connectionPromise = typeof kindOrFactory === 'function' - ? Promise.resolve(kindOrFactory(eventOptions, controller.signal)).then(connectEventTransport) - : createConnectionByKind(kindOrFactory, options as MilkyEventSourceCreateOptions) - connectionPromise.then((connection) => { - if (shouldCloseTransport || controller.signal.aborted) { - connection.source.close() - } - }, () => {}) - - const pending = withTimeout(connectionPromise, timeout, () => { - shouldCloseTransport = true - controller.abort() - }) - - return await raceWithAbort(signal, pending, () => { - shouldCloseTransport = true - controller.abort() - }) - } - catch (error) { - shouldCloseTransport = true - controller.abort() - throw error - } - } - - const controller = new AbortController() - const { signal } = controller - let currentConnection: MilkyEventSourceConnection | undefined - let controllerState!: MilkyEventSourceController - const emitter = new MilkyEventSourceImpl((state) => { - controllerState = state - controllerState.setCloseHandler(() => { - controller.abort() - currentConnection?.source.close() - }) - }) - - async function forwardConnection(connection: MilkyEventSourceConnection) { - const stopForwarding = controllerState.forwardFrom(connection.source) - - try { - return await connection.termination - } - finally { - stopForwarding() - } - } - - async function runConnection(): Promise { - let connection: MilkyEventSourceConnection | undefined - - try { - connection = await connect(signal) - currentConnection = connection - - if (signal.aborted) { - connection.source.close() - return 'stop' - } - - const termination = await forwardConnection(connection) - currentConnection = undefined - - if (signal.aborted || termination.type === 'closed') { - return 'stop' - } - - if (connection.kind === 'sse') { - if (termination.type === 'error' && !termination.reported) { - controllerState.dispatchError(termination.error) - } - - return 'stop' - } - - if (reconnect) { - controllerState.markConnecting() - } - - if (termination.type === 'ended') { - controllerState.dispatchError(createDisconnectError()) - } - else if (!termination.reported) { - controllerState.dispatchError(termination.error) - } - - return reconnect ? 'retry' : 'stop' - } - catch (error) { - currentConnection = undefined - - if (signal.aborted) { - return 'stop' - } - - if (reconnect) { - controllerState.markConnecting() - } - - controllerState.dispatchError(error) - return reconnect ? 'retry' : 'stop' - } - } - - void (async () => { - try { - if (!reconnect) { - await runConnection() - return - } - - let attempts = 0 - - while (!signal.aborted) { - if (await runConnection() === 'stop') { - break - } - - if (signal.aborted) { - break - } - - if (reconnect.attempts !== 'always' && attempts >= reconnect.attempts) { - break - } - - attempts += 1 - - try { - await sleepWithAbort(signal, reconnect.interval) - } - catch { - break - } - } - } - finally { - controllerState.markClosed() - } - })() - - return emitter -} - -export type { MilkyEventSource, MilkyEventSourceEventMap, MilkyEventSourceTransport } diff --git a/src/events/internal.ts b/src/events/internal.ts deleted file mode 100644 index 899a3e4..0000000 --- a/src/events/internal.ts +++ /dev/null @@ -1,331 +0,0 @@ -/* eslint-disable ts/method-signature-style */ - -import type { Event as MilkyEvent } from '@/gen/zod' -import mitt from 'mitt' - -const subscribeClose = Symbol('MilkyEventSourceImpl.subscribeClose') -const finishAsyncIteration = Symbol('MilkyEventSourceImpl.finishAsyncIteration') - -function freezeShallow(value: T): T { - if (value != null && typeof value === 'object') { - return Object.freeze(value) - } - return value -} - -interface MilkyEventConstraint {} - -export type ReadonlyMilkyEvent = Readonly> & MilkyEventConstraint - -export type MilkyEventSourceEventMap = { - error: any - push: ReadonlyMilkyEvent - open: void -} & { - [P in MilkyEvent['event_type']]: ReadonlyMilkyEvent

-} - -type MilkyEventSourceEventKey = keyof MilkyEventSourceEventMap - -export interface MilkyEventSource extends AsyncIterable { - on( - type: K, - listener: (ev: MilkyEventSourceEventMap[K]) => any, - ): void - on( - type: string, - listener: (ev: unknown) => any, - ): void - off( - type: K, - listener: (ev: MilkyEventSourceEventMap[K]) => any, - ): void - off( - type: string, - listener: (ev: unknown) => any, - ): void - - readonly readyState: number - readonly CONNECTING: 0 - readonly OPEN: 1 - readonly CLOSED: 2 - - close(): void - [Symbol.dispose](): void - [Symbol.asyncIterator](): AsyncIterableIterator -} - -export interface MilkyEventSourceTerminate { - (result: Result): void - readonly promise: Promise -} - -export class MilkyEventSourceController { - private _closeHandler: () => void = () => {} - private _closed = false - - constructor(readonly source: MilkyEventSourceImpl) {} - - createTerminate(): MilkyEventSourceTerminate { - const deferred = Promise.withResolvers() - let settled = false - const finish = ((result: Result) => { - if (settled) { - return - } - - settled = true - this.markClosed() - deferred.resolve(result) - }) as MilkyEventSourceTerminate - - return Object.assign(finish, { - promise: deferred.promise, - }) - } - - setCloseHandler(closeHandler: () => void): void { - this._closeHandler = closeHandler - } - - markConnecting(): void { - if (!this._closed) { - this.source.readyState = this.source.CONNECTING - } - } - - markClosed(): void { - if (this._closed) { - return - } - - this._closed = true - this.source.readyState = this.source.CLOSED - this.source[finishAsyncIteration]() - } - - dispatchOpen(): void { - if (this._closed) { - return - } - - this.source.readyState = this.source.OPEN - this.source.emit('open', void 0) - } - - dispatchMessage(message: MilkyEvent): void { - if (this._closed) { - return - } - - const frozenMessage = freezeShallow(message) - this.source.emit('push', frozenMessage as never) - - if (this._closed) { - return - } - - this.source.emit(frozenMessage.event_type, frozenMessage as never) - } - - dispatchError(error: unknown): void { - if (this._closed) { - return - } - - this.source.emit('error', error) - } - - forwardFrom(source: MilkyEventSource): () => void { - const onOpen = () => { - this.dispatchOpen() - } - - const onPush = (event: MilkyEventSourceEventMap['push']) => { - this.dispatchMessage(event as MilkyEvent) - } - - const onError = (event: MilkyEventSourceEventMap['error']) => { - this.dispatchError(event.error ?? event) - } - - source.on('open', onOpen) - source.on('push', onPush) - source.on('error', onError) - - return () => { - source.off('open', onOpen) - source.off('push', onPush) - source.off('error', onError) - } - } - - close(): void { - if (this._closed) { - return - } - - this.markClosed() - this._closeHandler() - } -} - -export class MilkyEventSourceImpl implements MilkyEventSource { - private _closeListeners = new Set<() => void>() - private _emitter = mitt() - - readonly CONNECTING = 0 - readonly OPEN = 1 - readonly CLOSED = 2 - static readonly CONNECTING = 0 - static readonly OPEN = 1 - static readonly CLOSED = 2 - - readyState = this.CONNECTING - readonly controller: MilkyEventSourceController - - constructor(setup?: (controller: MilkyEventSourceController) => void) { - this.controller = new MilkyEventSourceController(this) - setup?.(this.controller) - } - - close(): void { - this.controller.close() - } - - [Symbol.dispose](): void { - this.close() - } - - on( - type: K, - listener: (ev: MilkyEventSourceEventMap[K]) => any, - ): void - on( - type: string, - listener: (ev: unknown) => any, - ): void - on(type: string, listener: (ev: unknown) => any): void { - this._emitter.on(type as never, listener as never) - } - - off( - type: K, - listener: (ev: MilkyEventSourceEventMap[K]) => any, - ): void - off( - type: string, - listener: (ev: unknown) => any, - ): void - off(type: string, listener: (ev: unknown) => any): void { - this._emitter.off(type as never, listener as never) - } - - emit( - type: K, - event: MilkyEventSourceEventMap[K], - ): void { - this._emitter.emit(type, event) - } - - [subscribeClose](listener: () => void): () => void { - if (this.readyState === this.CLOSED) { - listener() - return () => {} - } - - this._closeListeners.add(listener) - return () => { - this._closeListeners.delete(listener) - } - } - - [finishAsyncIteration](): void { - const listeners = [...this._closeListeners] - this._closeListeners.clear() - for (const listener of listeners) { - listener() - } - } - - [Symbol.asyncIterator](): AsyncIterableIterator { - const queue: ReadonlyMilkyEvent[] = [] - const deferreds: Array<(result: IteratorResult) => void> = [] - let done = this.readyState === this.CLOSED - let unsubscribeClose = () => {} - - const cleanup = () => { - // eslint-disable-next-line ts/no-use-before-define - this.off('push', onPush) - unsubscribeClose() - } - - const finish = () => { - if (done) { - return - } - - done = true - cleanup() - while (deferreds.length > 0) { - deferreds.shift()!({ - done: true, - value: undefined, - }) - } - } - - const onPush = (event: MilkyEventSourceEventMap['push']) => { - if (done) { - return - } - - const deferred = deferreds.shift() - if (deferred) { - deferred({ - done: false, - value: event, - }) - return - } - - queue.push(event) - } - - if (!done) { - this.on('push', onPush) - unsubscribeClose = this[subscribeClose](finish) - } - - return { - next: async () => { - const message = queue.shift() - if (message) { - return { - done: false, - value: message, - } - } - - if (done) { - return { - done: true, - value: undefined, - } - } - - return await new Promise(resolve => deferreds.push(resolve)) - }, - return: async () => { - finish() - return { - done: true, - value: undefined, - } - }, - [Symbol.asyncIterator]() { - return this - }, - } - } -} diff --git a/src/events/source.ts b/src/events/source.ts deleted file mode 100644 index 479dde2..0000000 --- a/src/events/source.ts +++ /dev/null @@ -1,176 +0,0 @@ -/* eslint-disable ts/no-use-before-define */ -import type { MilkyEventSource, MilkyEventSourceController, MilkyEventSourceTerminate } from '@/events/internal' -import { MilkyEventSourceImpl } from '@/events/internal' - -export type MilkyEventSourceConnectionKind = 'sse' | 'websocket' -export type MilkyResolvedEventSourceConnectionKind = Exclude -export type MilkyEventSourceTransport = EventSource | WebSocket - -export type MilkyEventSourceTermination - = | { type: 'closed' } - | { type: 'ended' } - | { type: 'error', error: unknown, reported: boolean } - -export interface MilkyEventSourceConnection { - readonly kind: MilkyResolvedEventSourceConnectionKind - readonly source: MilkyEventSource - readonly termination: Promise -} - -function dispatchDeferredOpen(controller: MilkyEventSourceController): void { - controller.source.readyState = controller.source.OPEN - - // Defer synthetic open so callers awaiting source creation can still subscribe. - setTimeout(() => { - controller.dispatchOpen() - }, 0) -} - -function createTransportConnection( - kind: MilkyResolvedEventSourceConnectionKind, - setup: (controller: MilkyEventSourceController, finish: MilkyEventSourceTerminate) => void, -): MilkyEventSourceConnection { - let termination!: Promise - - const source = new MilkyEventSourceImpl((controller) => { - const finish = controller.createTerminate() - termination = finish.promise - setup(controller, finish) - }) - - return { - kind, - source, - termination, - } -} - -function isEventSourceTransport(source: MilkyEventSourceTransport): source is EventSource { - return !!globalThis.EventSource && source instanceof globalThis.EventSource -} - -export async function connectWebSocket(source: WebSocket): Promise { - return createTransportConnection('websocket', (controller, finish) => { - const cleanup = () => { - source.removeEventListener('open', onOpen) - source.removeEventListener('message', onMessage) - source.removeEventListener('error', onError) - source.removeEventListener('close', onClose) - } - - const onOpen = () => { - controller.dispatchOpen() - } - - const onMessage = (event: MessageEvent) => { - try { - controller.dispatchMessage(JSON.parse(event.data.toString())) - } - catch (error) { - controller.dispatchError(error) - } - } - - const onError = (event: Event) => { - controller.dispatchError(event) - } - - const onClose = () => { - cleanup() - finish({ type: 'ended' }) - } - - controller.setCloseHandler(() => { - cleanup() - if (source.readyState === source.CLOSING || source.readyState === source.CLOSED) { - finish({ type: 'closed' }) - return - } - - source.close() - }) - - source.addEventListener('open', onOpen) - source.addEventListener('message', onMessage) - source.addEventListener('error', onError) - source.addEventListener('close', onClose, { once: true }) - - if (source.readyState === source.OPEN) { - dispatchDeferredOpen(controller) - } - - if (source.readyState === source.CLOSED) { - queueMicrotask(() => { - cleanup() - finish({ type: 'ended' }) - }) - } - }) -} - -export async function connectEventSource(source: EventSource): Promise { - return createTransportConnection('sse', (controller, finish) => { - let closedByUser = false - - const cleanup = () => { - source.removeEventListener('open', onOpen) - source.removeEventListener('milky_event', onMessage) - source.removeEventListener('error', onError) - } - - const onOpen = () => { - controller.dispatchOpen() - } - - const onMessage = (event: Event) => { - const messageEvent = event as MessageEvent - try { - controller.dispatchMessage(JSON.parse(messageEvent.data)) - } - catch (error) { - controller.dispatchError(error) - } - } - - const onError = (event: Event) => { - if (closedByUser) { - return - } - - controller.markConnecting() - controller.dispatchError(event) - - if (source.readyState === source.CLOSED) { - cleanup() - finish({ type: 'error', error: event, reported: true }) - } - } - - controller.setCloseHandler(() => { - closedByUser = true - cleanup() - source.close() - finish({ type: 'closed' }) - }) - - source.addEventListener('open', onOpen) - source.addEventListener('milky_event', onMessage) - source.addEventListener('error', onError) - - if (source.readyState === source.OPEN) { - dispatchDeferredOpen(controller) - } - }) -} - -export async function connectEventTransport(source: MilkyEventSourceTransport): Promise { - if (globalThis.WebSocket && source instanceof WebSocket) { - return connectWebSocket(source) - } - - if (isEventSourceTransport(source) || !globalThis.EventSource) { - return connectEventSource(source as EventSource) - } - - throw new TypeError('milky: unknown event source type') -} diff --git a/src/index.ts b/src/index.ts index 5a1f0f3..79c745c 100644 --- a/src/index.ts +++ b/src/index.ts @@ -1,5 +1,5 @@ export * from './client' -export * from './events/index' +export { resolveMilkyEvent } from './event' export { milkyPackageVersion, milkyVersion } from './gen/types' export type * from './gen/types' export type { MilkyRawEndpointName, MilkyRawEndpoints } from './types' diff --git a/src/types.ts b/src/types.ts index fa7574e..0e9bef4 100644 --- a/src/types.ts +++ b/src/types.ts @@ -1,4 +1,4 @@ -import type { zodApiCategories } from '@/gen/zod' +import type { zodApiCategories } from '@/gen/zod-api' import { endpointNamesByCategory } from '@/gen/meta' type ApiCategories = typeof zodApiCategories diff --git a/tests/client.test.ts b/tests/client.test.ts index ee25815..a2271e8 100644 --- a/tests/client.test.ts +++ b/tests/client.test.ts @@ -1,43 +1,7 @@ import type { MilkyFetchOptions } from '@/client/fetch' -import type { MilkyEventSourceOptions } from '@/events' -import type { MilkyEventSourceConnectionKind } from '@/events/source' import type { QuitGroupInput } from '@/index' -import { afterEach, expect, expectTypeOf, it, vi } from 'vitest' +import { expect, expectTypeOf, it, vi } from 'vitest' import { createMilkyClient } from '@/client/endpoint' -import { onceEvent, waitFor } from './helpers/async' -import { FakeEventSource, FakeWebSocket } from './helpers/transports' - -const { fallbackEventSourceUrls } = vi.hoisted(() => ({ - fallbackEventSourceUrls: [] as string[], -})) - -vi.mock('eventsource', () => ({ - EventSource: class extends EventTarget { - readonly CONNECTING = 0 - readonly OPEN = 1 - readonly CLOSED = 2 - - readyState = this.CONNECTING - - constructor(readonly url: string | URL) { - super() - fallbackEventSourceUrls.push(String(url)) - } - - close(): void { - this.readyState = this.CLOSED - } - }, -})) - -const originalEventSource = globalThis.EventSource -const originalWebSocket = globalThis.WebSocket - -afterEach(() => { - fallbackEventSourceUrls.length = 0 - globalThis.EventSource = originalEventSource - globalThis.WebSocket = originalWebSocket -}) it('proxies grouped client methods to API endpoints', async () => { const fetchMock = vi.fn(async (request: Request) => { @@ -127,175 +91,6 @@ it('forwards params and per-request overrides through grouped client methods', a expect(overrideFetch).toHaveBeenCalledOnce() }) -it('adds token query when creating sse event urls', async () => { - const urls: string[] = [] - - globalThis.EventSource = class extends FakeEventSource { - constructor(url: string | URL) { - super(url) - urls.push(String(url)) - } - } as unknown as typeof EventSource - - const client = createMilkyClient({ - baseURL: 'https://example.com/base', - token: 'root-token', - fetch: vi.fn(), - }) - - const source = await client.event('sse') - - await waitFor(() => urls.length === 1) - expect(urls).toEqual(['https://example.com/event?access_token=root-token']) - - source.close() -}) - -it('allows per-event token overrides when creating sse event urls', async () => { - const urls: string[] = [] - - globalThis.EventSource = class extends FakeEventSource { - constructor(url: string | URL) { - super(url) - urls.push(String(url)) - } - } as unknown as typeof EventSource - - const client = createMilkyClient({ - baseURL: 'https://example.com/base', - token: 'root-token', - fetch: vi.fn(), - }) - - const source = await client.event('sse', { - token: 'event-token', - }) - - await waitFor(() => urls.length === 1) - expect(urls).toEqual(['https://example.com/event?access_token=event-token']) - - source.close() -}) - -it('adds token query when creating websocket event urls', async () => { - const urls: string[] = [] - - globalThis.WebSocket = class extends FakeWebSocket { - constructor(url: string | URL) { - super(url) - urls.push(String(url)) - } - } as unknown as typeof WebSocket - - const client = createMilkyClient({ - baseURL: 'https://example.com/base', - token: 'root-token', - fetch: vi.fn(), - }) - - const source = await client.event('websocket') - - await waitFor(() => urls.length === 1) - expect(urls).toEqual(['https://example.com/event?access_token=root-token']) - - source.close() -}) - -it('reports unsupported auto event kind through the client wrapper', async () => { - const websocketUrls: string[] = [] - const sseUrls: string[] = [] - - globalThis.WebSocket = class extends FakeWebSocket { - constructor(url: string | URL) { - super(url) - websocketUrls.push(String(url)) - throw new Error('unexpected websocket') - } - } as unknown as typeof WebSocket - - globalThis.EventSource = class extends FakeEventSource { - constructor(url: string | URL) { - super(url) - sseUrls.push(String(url)) - } - } as unknown as typeof EventSource - - const client = createMilkyClient({ - baseURL: 'https://example.com/base', - token: 'root-token', - fetch: vi.fn(), - }) - - // @ts-expect-error auto is no longer supported - const source = await client.event('auto', { - token: 'event-token', - timeout: 25, - reconnect: false, - }) - - const error = await onceEvent(source, 'error') - expect(error.message).toBe('milky: event source kind "auto" is no longer supported; use "websocket" or "sse"') - await waitFor(() => source.readyState === source.CLOSED) - expect(websocketUrls).toEqual([]) - expect(sseUrls).toEqual([]) - - source.close() -}) - -it('defaults client event kind to websocket when omitted', async () => { - const websocketUrls: string[] = [] - const sseUrls: string[] = [] - - globalThis.WebSocket = class extends FakeWebSocket { - constructor(url: string | URL) { - super(url) - websocketUrls.push(String(url)) - queueMicrotask(() => { - this.dispatchEvent(new Event('error')) - this.close() - }) - } - } as unknown as typeof WebSocket - - globalThis.EventSource = class extends FakeEventSource { - constructor(url: string | URL) { - super(url) - sseUrls.push(String(url)) - } - } as unknown as typeof EventSource - - const client = createMilkyClient({ - baseURL: 'https://example.com/base', - token: 'root-token', - fetch: vi.fn(), - }) - - const source = client.event() - - await waitFor(() => websocketUrls.length === 1) - expect(websocketUrls).toEqual(['https://example.com/event?access_token=root-token']) - expect(sseUrls).toEqual([]) - - source.close() -}) - -it('falls back to peer eventsource when global EventSource is unavailable', async () => { - globalThis.EventSource = undefined as unknown as typeof EventSource - - const client = createMilkyClient({ - baseURL: 'https://example.com/base', - token: 'root-token', - fetch: vi.fn(), - }) - - const source = await client.event('sse') - - await waitFor(() => fallbackEventSourceUrls.length === 1) - expect(fallbackEventSourceUrls).toEqual(['https://example.com/event?access_token=root-token']) - - source.close() -}) - it('exposes grouped client methods with optional override options', () => { const client = createMilkyClient({ baseURL: 'https://example.com', @@ -308,5 +103,5 @@ it('exposes grouped client methods with optional override options', () => { expectTypeOf(client.system.getLoginInfo).parameters.toEqualTypeOf<[(undefined | null)?, MilkyFetchOptions?]>() expectTypeOf(client.group.quitGroup).parameters.toEqualTypeOf<[QuitGroupInput, MilkyFetchOptions?]>() - expectTypeOf(client.event).parameters.toEqualTypeOf<[(MilkyEventSourceConnectionKind | undefined)?, MilkyEventSourceOptions?]>() + expect(client).not.toHaveProperty('event') }) diff --git a/tests/event.test.ts b/tests/event.test.ts new file mode 100644 index 0000000..43ceb25 --- /dev/null +++ b/tests/event.test.ts @@ -0,0 +1,119 @@ +import type { Event as MilkyEvent } from '@/gen/types' +import { afterEach, expect, expectTypeOf, it, vi } from 'vitest' +import { resolveMilkyEvent } from '@/event' + +afterEach(() => { + vi.doUnmock('@/gen/zod-event') + vi.resetModules() +}) + +it('resolves valid events and strips unknown fields', async () => { + const event = await resolveMilkyEvent({ + event_type: 'bot_offline', + time: 1, + self_id: 10001, + data: { + reason: 'network error', + ignored: true, + }, + ignored: true, + }) + + expectTypeOf(event).toEqualTypeOf() + expect(event).toEqual({ + event_type: 'bot_offline', + time: 1, + self_id: 10001, + data: { + reason: 'network error', + }, + }) +}) + +it('resolves nested message events', async () => { + const event = await resolveMilkyEvent({ + event_type: 'message_receive', + time: 2, + self_id: 10001, + data: { + message_scene: 'temp', + peer_id: 10002, + message_seq: 3, + sender_id: 10002, + time: 2, + segments: [{ + type: 'text', + data: { text: 'hello' }, + }], + }, + }) + + expect(event.event_type).toBe('message_receive') + if (event.event_type === 'message_receive') { + expect(event.data.segments).toEqual([{ + type: 'text', + data: { text: 'hello' }, + }]) + } +}) + +it.each([ + ['unknown event type', { event_type: 'unknown' }], + ['missing fields', { event_type: 'bot_offline' }], + ['invalid nested fields', { + event_type: 'bot_offline', + time: 1, + self_id: 10001, + data: { reason: 42 }, + }], +])('rejects %s with the validation error as its cause', async (_name, value) => { + const pending = resolveMilkyEvent(value) + + await expect(pending).rejects.toMatchObject({ + message: 'milky: failed to resolve event', + cause: expect.objectContaining({ + name: 'ZodError', + }), + }) +}) + +it('reports schema loading failures and retries the import', async () => { + let moduleLoads = 0 + vi.doMock('@/gen/zod-event', () => { + moduleLoads += 1 + throw new Error('Cannot find package zod') + }) + + const { resolveMilkyEvent: resolveWithMissingZod } = await import('@/event') + + await expect(resolveWithMissingZod({})).rejects.toMatchObject({ + message: 'milky: zod is required to resolve events', + cause: expect.any(Error), + }) + await expect(resolveWithMissingZod({})).rejects.toThrow('milky: zod is required to resolve events') + expect(moduleLoads).toBe(2) +}) + +it('loads the event schema once after successful resolution', async () => { + let moduleLoads = 0 + const safeParseAsync = vi.fn(async (value: unknown) => ({ + success: true as const, + data: value as MilkyEvent, + })) + + vi.doMock('@/gen/zod-event', () => { + moduleLoads += 1 + return { + Event: { safeParseAsync }, + } + }) + + const { resolveMilkyEvent: resolveWithMock } = await import('@/event') + const value = { event_type: 'bot_offline' } as never + + await resolveWithMock(value) + await resolveWithMock(value) + + expect(moduleLoads).toBe(1) + expect(safeParseAsync).toHaveBeenCalledTimes(2) +}) diff --git a/tests/events.coverage.test.ts b/tests/events.coverage.test.ts deleted file mode 100644 index cfe6dec..0000000 --- a/tests/events.coverage.test.ts +++ /dev/null @@ -1,287 +0,0 @@ -import { afterEach, expect, it, vi } from 'vitest' -import { createMilkyEventSource } from '@/events/index' -import { createDeferred, sleep, waitFor } from './helpers/async' -import { FakeWebSocket } from './helpers/transports' - -const originalWebSocket = globalThis.WebSocket - -afterEach(() => { - globalThis.WebSocket = originalWebSocket - vi.doUnmock('@/events/source') - vi.resetModules() -}) - -it('closes websocket connections after transport failure', async () => { - globalThis.WebSocket = class extends FakeWebSocket { - constructor() { - super() - queueMicrotask(() => { - this.dispatchEvent(new Event('error')) - this.close() - }) - } - } as unknown as typeof WebSocket - - const source = createMilkyEventSource('websocket', { - baseURL: 'https://example.com', - timeout: 100, - }) - - await sleep(10) - source.close() - - expect(source.readyState).toBe(source.CLOSED) -}) - -it('handles signal abort during websocket connection', async () => { - globalThis.WebSocket = class extends FakeWebSocket { - constructor() { - super() - queueMicrotask(() => { - this.dispatchEvent(new Event('error')) - this.close() - }) - } - } as unknown as typeof WebSocket - - const source = createMilkyEventSource('websocket', { - baseURL: 'https://example.com', - timeout: 5000, - }) - - source.close() - await sleep(10) - - expect(source.readyState).toBe(source.CLOSED) -}) - -it('handles connection that resolves after abort during reconnect', async () => { - const { createMilkyEventSource } = await (async () => { - vi.resetModules() - const deferred = createDeferred() - vi.doMock('@/events/source', async () => { - const actual = await vi.importActual('@/events/source') - return { - ...actual, - connectEventTransport: vi.fn(async () => { - await sleep(100) - return deferred.promise - }), - } - }) - - const events = await import('@/events/index') - const internal = await import('@/events/internal') - - setTimeout(() => { - deferred.resolve({ - kind: 'websocket', - source: new internal.MilkyEventSourceImpl(), - termination: Promise.resolve({ type: 'closed' }), - }) - }, 50) - - return { - createMilkyEventSource: events.createMilkyEventSource, - MilkyEventSourceImpl: internal.MilkyEventSourceImpl, - } - })() - - globalThis.WebSocket = FakeWebSocket as unknown as typeof WebSocket - - const source = createMilkyEventSource(() => new FakeWebSocket() as never, { - reconnect: { - interval: 1, - attempts: 1, - }, - }) - - await sleep(20) - source.close() - await sleep(150) - - expect(source.readyState).toBe(source.CLOSED) -}) - -it('handles iterator on closed source', async () => { - const { MilkyEventSourceImpl } = await import('@/events/internal') - - const source = new MilkyEventSourceImpl() - source.close() - - const iterator = source[Symbol.asyncIterator]() - const result = await iterator.next() - - expect(result).toEqual({ - done: true, - value: undefined, - }) -}) - -it('handles async iterator return method', async () => { - const { MilkyEventSourceImpl } = await import('@/events/internal') - - const source = new MilkyEventSourceImpl() - const iterator = source[Symbol.asyncIterator]() - - const result = await iterator.return?.() - expect(result).toEqual({ - done: true, - value: undefined, - }) -}) - -it('handles async iterator with queued messages', async () => { - const { MilkyEventSourceImpl } = await import('@/events/internal') - - const source = new MilkyEventSourceImpl() - const iterator = source[Symbol.asyncIterator]() - - // Dispatch messages before consuming - source.controller.dispatchMessage({ event_type: 'test', id: 1 } as never) - source.controller.dispatchMessage({ event_type: 'test', id: 2 } as never) - - const first = await iterator.next() - expect(first.value).toMatchObject({ id: 1 }) - - const second = await iterator.next() - expect(second.value).toMatchObject({ id: 2 }) - - source.close() -}) - -it('handles reconnect interval sleep interruption', async () => { - const { createMilkyEventSource } = await (async () => { - vi.resetModules() - let attemptCount = 0 - vi.doMock('@/events/source', async () => { - const actual = await vi.importActual('@/events/source') - return { - ...actual, - connectEventTransport: vi.fn(async () => { - attemptCount++ - if (attemptCount === 1) { - return { - kind: 'websocket', - source: new (await import('@/events/internal')).MilkyEventSourceImpl(), - termination: Promise.resolve({ type: 'ended' }), - } - } - await sleep(1000) - throw new Error('should not reach') - }), - } - }) - - const events = await import('@/events/index') - const internal = await import('@/events/internal') - - return { - createMilkyEventSource: events.createMilkyEventSource, - MilkyEventSourceImpl: internal.MilkyEventSourceImpl, - } - })() - - globalThis.WebSocket = FakeWebSocket as unknown as typeof WebSocket - - const source = createMilkyEventSource(() => new FakeWebSocket() as never, { - reconnect: { - interval: 5000, - attempts: 'always', - }, - }) - - await sleep(50) - source.close() - await waitFor(() => source.readyState === source.CLOSED) -}) - -it('handles connection closed during signal abort check', async () => { - const { createMilkyEventSource } = await (async () => { - vi.resetModules() - const deferred = createDeferred() - vi.doMock('@/events/source', async () => { - const actual = await vi.importActual('@/events/source') - return { - ...actual, - connectEventTransport: vi.fn(async () => deferred.promise), - } - }) - - const events = await import('@/events/index') - const internal = await import('@/events/internal') - - setTimeout(() => { - const impl = new internal.MilkyEventSourceImpl() - deferred.resolve({ - kind: 'websocket', - source: impl, - termination: Promise.resolve({ type: 'closed' }), - }) - }, 10) - - return { - createMilkyEventSource: events.createMilkyEventSource, - MilkyEventSourceImpl: internal.MilkyEventSourceImpl, - } - })() - - globalThis.WebSocket = FakeWebSocket as unknown as typeof WebSocket - - const source = createMilkyEventSource(() => new FakeWebSocket() as never, { - reconnect: false, - }) - - await sleep(50) - expect(source.readyState).toBe(source.CLOSED) -}) - -it('handles websocket ended termination with reconnect', async () => { - const { createMilkyEventSource } = await (async () => { - vi.resetModules() - let callCount = 0 - vi.doMock('@/events/source', async () => { - const actual = await vi.importActual('@/events/source') - return { - ...actual, - connectEventTransport: vi.fn(async () => { - callCount++ - const impl = new (await import('@/events/internal')).MilkyEventSourceImpl() - if (callCount === 1) { - return { - kind: 'websocket', - source: impl, - termination: Promise.resolve({ type: 'ended' }), - } - } - return { - kind: 'websocket', - source: impl, - termination: new Promise(() => {}), // never resolves - } - }), - } - }) - - const events = await import('@/events/index') - const internal = await import('@/events/internal') - - return { - createMilkyEventSource: events.createMilkyEventSource, - MilkyEventSourceImpl: internal.MilkyEventSourceImpl, - } - })() - - globalThis.WebSocket = FakeWebSocket as unknown as typeof WebSocket - - const source = createMilkyEventSource(() => new FakeWebSocket() as never, { - reconnect: { - interval: 10, - attempts: 1, - }, - }) - - await sleep(50) - expect(source.readyState).toBe(source.CONNECTING) - source.close() -}) diff --git a/tests/events.logic.test.ts b/tests/events.logic.test.ts deleted file mode 100644 index 29d27b5..0000000 --- a/tests/events.logic.test.ts +++ /dev/null @@ -1,304 +0,0 @@ -import { afterEach, expect, it, vi } from 'vitest' -import { createDeferred, onceEvent, sleep, waitFor } from './helpers/async' -import { FakeWebSocket } from './helpers/transports' - -afterEach(() => { - vi.doUnmock('@/events/source') - vi.resetModules() -}) - -async function loadEventsWithMock( - connectImpl: (transport: unknown) => Promise, -): Promise<{ - createMilkyEventSource: typeof import('@/events/index').createMilkyEventSource - MilkyEventSourceImpl: typeof import('@/events/internal').MilkyEventSourceImpl -}> { - vi.resetModules() - vi.doMock('@/events/source', async () => { - const actual = await vi.importActual('@/events/source') - return { - ...actual, - connectEventTransport: vi.fn(connectImpl), - } - }) - - const events = await import('@/events/index') - const internal = await import('@/events/internal') - - return { - createMilkyEventSource: events.createMilkyEventSource, - MilkyEventSourceImpl: internal.MilkyEventSourceImpl, - } -} - -it('guards controller operations after closure and only settles terminations once', async () => { - const { MilkyEventSourceImpl } = await import('@/events/internal') - - const closedSource = new MilkyEventSourceImpl() - let openCount = 0 - let pushCount = 0 - let errorCount = 0 - - closedSource.on('open', () => { - openCount += 1 - }) - closedSource.on('push', () => { - pushCount += 1 - }) - closedSource.on('error', () => { - errorCount += 1 - }) - - closedSource.close() - closedSource.close() - closedSource.controller.dispatchOpen() - closedSource.controller.dispatchMessage({} as never) - closedSource.controller.dispatchError(new Error('late')) - - expect(openCount).toBe(0) - expect(pushCount).toBe(0) - expect(errorCount).toBe(0) - - const pendingSource = new MilkyEventSourceImpl() - const finish = pendingSource.controller.createTerminate() - - finish('first') - finish('second') - - await expect(finish.promise).resolves.toBe('first') -}) - -it('dispatches push events, typed events, and async iteration from the same payload', async () => { - const { MilkyEventSourceImpl } = await import('@/events/internal') - - const source = new MilkyEventSourceImpl() - const payload = { - event_type: 'private_message_created', - id: 1, - } as never - const iterator = source[Symbol.asyncIterator]() - - const pushEvent = onceEvent(source, 'push') - const typedEvent = onceEvent(source, 'private_message_created') - const firstMessage = iterator.next() - source.controller.dispatchMessage(payload) - - await expect(pushEvent).resolves.toEqual(payload) - await expect(typedEvent).resolves.toEqual(payload) - await expect(firstMessage).resolves.toMatchObject({ - done: false, - value: payload, - }) - - const finished = iterator.next() - source.close() - - await expect(finished).resolves.toMatchObject({ - done: true, - }) -}) - -it('exposes mitt-style subscriptions with readonly push payloads', async () => { - const { MilkyEventSourceImpl } = await import('@/events/internal') - - const source = new MilkyEventSourceImpl() as never as { - controller: InstanceType['controller'] - on: (type: string, handler: (event: unknown) => void) => void - } - const payload = { - event_type: 'private_message_created', - id: 1, - } as never - const received: unknown[] = [] - - expect(typeof source.on).toBe('function') - - source.on('push', (event) => { - received.push(event) - }) - - source.controller.dispatchMessage(payload) - - expect(received).toEqual([payload]) -}) - -it('exposes forwarded message events through async iteration', async () => { - const { createMilkyEventSource } = await import('@/events/index') - const originalWebSocket = globalThis.WebSocket - - globalThis.WebSocket = FakeWebSocket as unknown as typeof WebSocket - - try { - const socket = new FakeWebSocket() - const source = await createMilkyEventSource(() => socket as unknown as WebSocket) - const received: unknown[] = [] - const pushEvents: unknown[] = [] - - await sleep(0) - - const consume = (async () => { - for await (const message of source) { - received.push(message) - if (received.length === 2) { - break - } - } - })() - - source.on('push', (event) => { - pushEvents.push(event) - }) - - socket.open() - socket.sendMessage({ event_type: 'private_message_created', id: 1 }) - socket.sendMessage({ event_type: 'private_message_created', id: 2 }) - - await consume - - expect(received).toEqual([ - { event_type: 'private_message_created', id: 1 }, - { event_type: 'private_message_created', id: 2 }, - ]) - expect(pushEvents).toEqual([ - { event_type: 'private_message_created', id: 1 }, - { event_type: 'private_message_created', id: 2 }, - ]) - - source.close() - } - finally { - globalThis.WebSocket = originalWebSocket - } -}) - -it('dispatches unreported sse termination errors from reconnect loops', async () => { - const connected = createDeferred() - const termination = createDeferred<{ - type: 'error' - error: Error - reported: false - }>() - const { createMilkyEventSource, MilkyEventSourceImpl } = await loadEventsWithMock(async () => { - connected.resolve() - return { - kind: 'sse', - source: new MilkyEventSourceImpl(), - termination: termination.promise, - } - }) - - const source = await createMilkyEventSource(() => new FakeWebSocket() as never, { - reconnect: { - interval: 1, - attempts: 'always', - }, - }) - - await connected.promise - - let errorMessage: string | undefined - const onError = (event: any) => { - errorMessage = event instanceof Error ? event.message : String(event) - source.off('error', onError) - } - source.on('error', onError) - - termination.resolve({ - type: 'error', - error: new Error('sse failed'), - reported: false, - }) - - await waitFor(() => errorMessage != null) - expect(errorMessage).toBe('sse failed') - await waitFor(() => source.readyState === source.CLOSED) -}) - -it('stops reconnecting when websocket termination errors close the source from an error listener', async () => { - const termination = createDeferred<{ - type: 'error' - error: Error - reported: false - }>() - const { createMilkyEventSource, MilkyEventSourceImpl } = await loadEventsWithMock(async () => ({ - kind: 'websocket', - source: new MilkyEventSourceImpl(), - termination: termination.promise, - })) - - const source = await createMilkyEventSource(() => new FakeWebSocket() as never, { - reconnect: { - interval: 1, - attempts: 'always', - }, - }) - - const errorMessages: string[] = [] - const onError = (event: any) => { - errorMessages.push(event instanceof Error ? event.message : String(event)) - source.close() - source.off('error', onError) - } - source.on('error', onError) - - termination.resolve({ - type: 'error', - error: new Error('ws failed'), - reported: false, - }) - - await waitFor(() => source.readyState === source.CLOSED) - expect(errorMessages).toEqual(['ws failed']) -}) - -it('reports transport connection failures during reconnect loops', async () => { - const connected = createDeferred() - const { createMilkyEventSource } = await loadEventsWithMock(async () => { - connected.resolve() - await sleep(0) - throw new Error('connect failed') - }) - - const source = await createMilkyEventSource(() => new FakeWebSocket() as never, { - reconnect: { - interval: 1, - attempts: 0, - }, - }) - - await connected.promise - const error = onceEvent(source, 'error') - - const errorResult = await error as Error - expect(errorResult.message).toBe('connect failed') - await waitFor(() => source.readyState === source.CLOSED) -}) - -it('breaks reconnect loops when closed during an in-flight connect', async () => { - const deferred = createDeferred<{ - kind: 'websocket' - source: InstanceType - termination: Promise<{ type: 'closed' }> - }>() - const { createMilkyEventSource, MilkyEventSourceImpl } = await loadEventsWithMock(() => deferred.promise) - - const source = await createMilkyEventSource(() => new FakeWebSocket() as never, { - reconnect: { - interval: 20, - attempts: 'always', - }, - }) - - source.close() - deferred.resolve({ - kind: 'websocket', - source: new MilkyEventSourceImpl(), - termination: Promise.resolve({ - type: 'closed', - }), - }) - - await sleep(10) - - expect(source.readyState).toBe(source.CLOSED) -}) diff --git a/tests/events.test.ts b/tests/events.test.ts deleted file mode 100644 index c03c6b0..0000000 --- a/tests/events.test.ts +++ /dev/null @@ -1,458 +0,0 @@ -import { afterEach, expect, it, vi } from 'vitest' -import { createMilkyEventSource } from '@/events/index' -import { connectEventSource, connectEventTransport, connectWebSocket } from '@/events/source' -import { createDeferred, onceEvent, sleep, waitFor } from './helpers/async' -import { FakeEventSource, FakeWebSocket } from './helpers/transports' - -const originalEventSource = globalThis.EventSource -const originalWebSocket = globalThis.WebSocket - -afterEach(() => { - globalThis.EventSource = originalEventSource - globalThis.WebSocket = originalWebSocket - vi.doUnmock('eventsource') -}) - -it('connects websocket transports and forwards push and typed events plus parse errors', async () => { - globalThis.WebSocket = FakeWebSocket as unknown as typeof WebSocket - - const socket = new FakeWebSocket() - const connection = await connectWebSocket(socket as unknown as WebSocket) - const payload = { - event_type: 'private_message_created', - id: 1, - } as never - - const openEvent = onceEvent(connection.source, 'open') - socket.open() - await expect(openEvent).resolves.toBeUndefined() - - const pushEvent = onceEvent(connection.source, 'push') - const typedEvent = onceEvent(connection.source, 'private_message_created') - socket.sendMessage(payload) - await expect(pushEvent).resolves.toEqual(payload) - await expect(typedEvent).resolves.toEqual(payload) - - const errorEvent = onceEvent(connection.source, 'error') - socket.sendRawMessage('{') - expect(await errorEvent).toBeInstanceOf(Error) - - socket.close() - await expect(connection.termination).resolves.toEqual({ - type: 'ended', - }) -}) - -it('dispatches open for already-open websocket transports', async () => { - globalThis.WebSocket = FakeWebSocket as unknown as typeof WebSocket - - const socket = new FakeWebSocket() - socket.readyState = socket.OPEN - - const connection = await connectWebSocket(socket as unknown as WebSocket) - - await expect(Promise.race([ - onceEvent(connection.source, 'open'), - sleep(20).then(() => null), - ])).resolves.toBeUndefined() - expect(connection.source.readyState).toBe(connection.source.OPEN) - - connection.source.close() -}) - -it('dispatches open for already-open websocket sources created by milky', async () => { - globalThis.WebSocket = FakeWebSocket as unknown as typeof WebSocket - - const socket = new FakeWebSocket() - socket.readyState = socket.OPEN - - const source = await createMilkyEventSource(() => socket as unknown as WebSocket) - - await expect(Promise.race([ - onceEvent(source, 'open'), - sleep(20).then(() => null), - ])).resolves.toBeUndefined() - expect(source.readyState).toBe(source.OPEN) - - source.close() -}) - -it('does not close websocket transports that are already closing', async () => { - globalThis.WebSocket = FakeWebSocket as unknown as typeof WebSocket - - const socket = new FakeWebSocket() - const connection = await connectWebSocket(socket as unknown as WebSocket) - socket.readyState = socket.CLOSING - - connection.source.close() - - expect(socket.closeCalls).toBe(0) - await expect(connection.termination).resolves.toEqual({ - type: 'closed', - }) -}) - -it('connects event sources and reports terminal errors', async () => { - globalThis.EventSource = FakeEventSource as unknown as typeof EventSource - - const source = new FakeEventSource() - const connection = await connectEventSource(source as unknown as EventSource) - const payload = { - event_type: 'private_message_created', - id: 2, - } as never - - const openEvent = onceEvent(connection.source, 'open') - source.open() - await expect(openEvent).resolves.toBeUndefined() - - const pushEvent = onceEvent(connection.source, 'push') - const typedEvent = onceEvent(connection.source, 'private_message_created') - source.sendMessage(payload) - await expect(pushEvent).resolves.toEqual(payload) - await expect(typedEvent).resolves.toEqual(payload) - - const parseError = onceEvent(connection.source, 'error') - source.sendRawMessage('{') - expect(await parseError).toBeInstanceOf(Error) - - const terminalError = onceEvent(connection.source, 'error') - source.fail({ closed: true }) - - expect(await terminalError).toBeInstanceOf(Event) - await expect(connection.termination).resolves.toMatchObject({ - type: 'error', - reported: true, - }) -}) - -it('ignores event source errors after the consumer closes the source', async () => { - globalThis.EventSource = FakeEventSource as unknown as typeof EventSource - - const source = new FakeEventSource() - const connection = await connectEventSource(source as unknown as EventSource) - let errorCount = 0 - - connection.source.on('error', () => { - errorCount += 1 - }) - - connection.source.close() - source.fail({ closed: true }) - - expect(source.closeCalls).toBe(1) - expect(errorCount).toBe(0) - await expect(connection.termination).resolves.toEqual({ - type: 'closed', - }) -}) - -it('rejects unknown event transport values', async () => { - globalThis.EventSource = FakeEventSource as unknown as typeof EventSource - globalThis.WebSocket = FakeWebSocket as unknown as typeof WebSocket - - await expect(connectEventTransport(new EventTarget() as never)).rejects.toThrow('milky: unknown event source type') -}) - -it('requires baseURL when creating event sources by kind', async () => { - expect(() => createMilkyEventSource('sse', undefined as never)).toThrow('milky: baseURL is required when creating event sources by kind') -}) - -it('aborts timed out connection attempts and closes transports that resolve later', async () => { - globalThis.WebSocket = FakeWebSocket as unknown as typeof WebSocket - - const deferred = createDeferred() - const socket = new FakeWebSocket() - let aborted = false - - const pending = createMilkyEventSource(async (_options, signal) => { - signal?.addEventListener('abort', () => { - aborted = true - }, { once: true }) - - return deferred.promise - }, { - timeout: 5, - }) - - const error = onceEvent(pending, 'error') - const errorResult = await error as Error - expect(errorResult).toBeInstanceOf(Error) - expect(errorResult.message).toBe('milky: timed out after 5ms') - await waitFor(() => pending.readyState === pending.CLOSED) - - deferred.resolve(socket as unknown as WebSocket) - await sleep(10) - - expect(aborted).toBe(true) - expect(socket.closeCalls).toBe(1) -}) - -it('reconnects websocket transports with native EventSource semantics', async () => { - globalThis.WebSocket = FakeWebSocket as unknown as typeof WebSocket - - const sockets: FakeWebSocket[] = [] - const source = await createMilkyEventSource(() => { - const socket = new FakeWebSocket() - sockets.push(socket) - return socket as unknown as WebSocket - }, { - reconnect: { - interval: 1, - attempts: 1, - }, - }) - - let openCount = 0 - let errorCount = 0 - - source.on('open', () => { - openCount += 1 - }) - source.on('error', () => { - errorCount += 1 - }) - - await waitFor(() => sockets.length === 1) - await sleep(1) - sockets[0]!.open() - await sleep(1) - expect(source.readyState).toBe(source.OPEN) - expect(openCount).toBe(1) - - sockets[0]!.close() - await sleep(10) - - expect(errorCount).toBe(1) - expect(source.readyState).toBe(source.CONNECTING) - expect(sockets).toHaveLength(2) - - sockets[1]!.open() - await sleep(1) - expect(source.readyState).toBe(source.OPEN) - expect(openCount).toBe(2) - - source.close() - expect(source.readyState).toBe(source.CLOSED) - expect(sockets[1]!.closeCalls).toBe(1) -}) - -it('honors websocket reconnect attempt limits', async () => { - globalThis.WebSocket = FakeWebSocket as unknown as typeof WebSocket - - const sockets: FakeWebSocket[] = [] - const source = await createMilkyEventSource(() => { - const socket = new FakeWebSocket() - sockets.push(socket) - return socket as unknown as WebSocket - }, { - reconnect: { - interval: 1, - attempts: 0, - }, - }) - - await waitFor(() => sockets.length === 1) - sockets[0]!.open() - await sleep(0) - sockets[0]!.close() - await sleep(10) - - expect(sockets).toHaveLength(1) - expect(source.readyState).toBe(source.CLOSED) -}) - -it('stops reconnecting when closed during reconnect backoff', async () => { - globalThis.WebSocket = FakeWebSocket as unknown as typeof WebSocket - - const sockets: FakeWebSocket[] = [] - const source = await createMilkyEventSource(() => { - const socket = new FakeWebSocket() - sockets.push(socket) - return socket as unknown as WebSocket - }, { - reconnect: { - interval: 20, - attempts: 'always', - }, - }) - - await waitFor(() => sockets.length === 1) - sockets[0]!.open() - await sleep(1) - sockets[0]!.close() - source.close() - await sleep(30) - - expect(sockets).toHaveLength(1) - expect(source.readyState).toBe(source.CLOSED) -}) - -it('creates websocket transports from kind and options overload', async () => { - const urls: string[] = [] - - globalThis.WebSocket = class extends FakeWebSocket { - constructor(url: string | URL) { - super(url) - urls.push(String(url)) - } - } as unknown as typeof WebSocket - - const source = await createMilkyEventSource('websocket', { - baseURL: 'https://example.com/base', - token: 'event-token', - }) - - await waitFor(() => urls.length === 1) - expect(urls).toEqual(['https://example.com/event?access_token=event-token']) - - source.close() - expect(source.readyState).toBe(source.CLOSED) -}) - -it('closes sources through Symbol.dispose', async () => { - globalThis.WebSocket = FakeWebSocket as unknown as typeof WebSocket - - const socket = new FakeWebSocket() - const source = await createMilkyEventSource(() => socket as unknown as WebSocket) - - source[Symbol.dispose]() - - expect(source.readyState).toBe(source.CLOSED) - await waitFor(() => socket.closeCalls === 1) -}) - -it('reports unsupported auto event kind without creating transports', async () => { - const websocketUrls: string[] = [] - const sseUrls: string[] = [] - - globalThis.WebSocket = class extends FakeWebSocket { - constructor(url: string | URL) { - super(url) - websocketUrls.push(String(url)) - throw new Error('unexpected websocket') - } - } as unknown as typeof WebSocket - - globalThis.EventSource = class extends FakeEventSource { - constructor(url: string | URL) { - super(url) - sseUrls.push(String(url)) - throw new Error('unexpected sse fallback') - } - } as unknown as typeof EventSource - - // @ts-expect-error auto is no longer supported - const source = createMilkyEventSource('auto', { - baseURL: 'https://example.com/base', - token: 'event-token', - }) - - const error = await onceEvent(source, 'error') as Error - expect(error.message).toBe('milky: event source kind "auto" is no longer supported; use "websocket" or "sse"') - await waitFor(() => source.readyState === source.CLOSED) - expect(websocketUrls).toEqual([]) - expect(sseUrls).toEqual([]) - - source.close() - expect(source.readyState).toBe(source.CLOSED) -}) - -it('reports unsupported auto before websocket can open', async () => { - globalThis.WebSocket = class extends FakeWebSocket { - constructor(url: string | URL) { - super(url) - setTimeout(() => { - this.open() - }, 0) - } - } as unknown as typeof WebSocket - - // @ts-expect-error auto is no longer supported - const source = createMilkyEventSource('auto', { - baseURL: 'https://example.com/base', - }) - - const error = await onceEvent(source, 'error') as Error - expect(error.message).toBe('milky: event source kind "auto" is no longer supported; use "websocket" or "sse"') - await waitFor(() => source.readyState === source.CLOSED) - - source.close() -}) - -it('reports unsupported auto before websocket construction errors', async () => { - const sseUrls: string[] = [] - - globalThis.WebSocket = class { - constructor() { - throw new Error('boom') - } - } as unknown as typeof WebSocket - - globalThis.EventSource = class extends FakeEventSource { - constructor(url: string | URL) { - super(url) - sseUrls.push(String(url)) - throw new Error('unexpected sse fallback') - } - } as unknown as typeof EventSource - - // @ts-expect-error auto is no longer supported - const source = createMilkyEventSource('auto', { - baseURL: 'https://example.com/base', - token: 'event-token', - }) - - const error = await onceEvent(source, 'error') as Error - expect(error.message).toBe('milky: event source kind "auto" is no longer supported; use "websocket" or "sse"') - await waitFor(() => source.readyState === source.CLOSED) - expect(sseUrls).toEqual([]) - - source.close() -}) - -it('uses native EventSource reconnection behavior directly', async () => { - globalThis.EventSource = FakeEventSource as unknown as typeof EventSource - - const nativeSource = new FakeEventSource() - const source = await createMilkyEventSource(() => nativeSource as unknown as EventSource, { - reconnect: { - interval: 1, - attempts: 1, - }, - }) - - let openCount = 0 - let errorCount = 0 - source.on('open', () => { - openCount += 1 - }) - source.on('error', () => { - errorCount += 1 - }) - - await sleep(1) - nativeSource.open() - await sleep(1) - nativeSource.fail() - await sleep(1) - nativeSource.open() - await sleep(1) - expect(openCount).toBe(2) - expect(errorCount).toBe(1) - expect(source.readyState).toBe(source.OPEN) - - const payload = { - event_type: 'private_message_created', - id: 2, - } as never - const pushEvent = onceEvent(source, 'push') - const typedEvent = onceEvent(source, 'private_message_created') - nativeSource.sendMessage(payload) - await expect(pushEvent).resolves.toEqual(payload) - await expect(typedEvent).resolves.toEqual(payload) - - source.close() - expect(source.readyState).toBe(source.CLOSED) - expect(nativeSource.readyState).toBe(nativeSource.CLOSED) -}) diff --git a/tests/fetch.test.ts b/tests/fetch.test.ts index af9160d..73c2e9c 100644 --- a/tests/fetch.test.ts +++ b/tests/fetch.test.ts @@ -15,7 +15,7 @@ function createJsonResponse(body: unknown, init?: ResponseInit): Response { afterEach(() => { globalThis.fetch = originalFetch - vi.doUnmock('@/gen/zod') + vi.doUnmock('@/gen/zod-api') }) it('uses the global fetch implementation when no local fetch is provided', async () => { @@ -261,7 +261,7 @@ it('skips request and response validation when strict is disabled on the client' it('skips zod validation when zod schemas are unavailable', async () => { vi.resetModules() - vi.doMock('@/gen/zod', () => { + vi.doMock('@/gen/zod-api', () => { throw new Error('Cannot find package "zod"') }) const { createMilkyFetch: createFetchWithoutZod } = await import('@/client/fetch') diff --git a/tests/helpers/async.ts b/tests/helpers/async.ts index 9e2e8af..fd0c531 100644 --- a/tests/helpers/async.ts +++ b/tests/helpers/async.ts @@ -1,52 +1,3 @@ export function sleep(ms = 0): Promise { return new Promise(resolve => setTimeout(resolve, ms)) } - -export async function waitFor(predicate: () => boolean, timeout = 100): Promise { - const start = performance.now() - - while (!predicate()) { - if (performance.now() - start > timeout) { - throw new Error('timed out waiting for condition') - } - - await sleep(0) - } -} - -export function onceEvent( - target: { - on: (type: string, listener: (event: T) => void) => void - off: (type: string, listener: (event: T) => void) => void - }, - type: string, -): Promise { - return new Promise((resolve) => { - const listener = (event: T) => { - target.off(type, listener) - resolve(event) - } - - target.on(type, listener) - }) -} - -export function createDeferred(): { - readonly promise: Promise - resolve: (value: T | PromiseLike) => void - reject: (reason?: unknown) => void -} { - let resolve!: (value: T | PromiseLike) => void - let reject!: (reason?: unknown) => void - - const promise = new Promise((innerResolve, innerReject) => { - resolve = innerResolve - reject = innerReject - }) - - return { - promise, - resolve, - reject, - } -} diff --git a/tests/helpers/transports.ts b/tests/helpers/transports.ts deleted file mode 100644 index 9483f3f..0000000 --- a/tests/helpers/transports.ts +++ /dev/null @@ -1,89 +0,0 @@ -export class FakeWebSocket extends EventTarget { - readonly CONNECTING = 0 - readonly OPEN = 1 - readonly CLOSING = 2 - readonly CLOSED = 3 - - readyState = this.CONNECTING - closeCalls = 0 - - constructor(readonly url?: string | URL) { - super() - } - - open(): void { - if (this.readyState === this.CLOSED) { - return - } - - this.readyState = this.OPEN - this.dispatchEvent(new Event('open')) - } - - sendMessage(data: unknown): void { - this.sendRawMessage(JSON.stringify(data)) - } - - sendRawMessage(data: string): void { - this.dispatchEvent(new MessageEvent('message', { - data, - })) - } - - fail(event: Event = new Event('error')): void { - this.dispatchEvent(event) - } - - close(): void { - this.closeCalls += 1 - - if (this.readyState === this.CLOSED) { - return - } - - this.readyState = this.CLOSED - this.dispatchEvent(new Event('close')) - } -} - -export class FakeEventSource extends EventTarget { - readonly CONNECTING = 0 - readonly OPEN = 1 - readonly CLOSED = 2 - - readyState = this.CONNECTING - closeCalls = 0 - - constructor(readonly url?: string | URL) { - super() - } - - open(): void { - if (this.readyState === this.CLOSED) { - return - } - - this.readyState = this.OPEN - this.dispatchEvent(new Event('open')) - } - - sendMessage(data: unknown): void { - this.sendRawMessage(JSON.stringify(data)) - } - - sendRawMessage(data: string): void { - this.dispatchEvent(new MessageEvent('milky_event', { - data, - })) - } - - fail({ closed = false, event = new Event('error') }: { closed?: boolean, event?: Event } = {}): void { - this.readyState = closed ? this.CLOSED : this.CONNECTING - this.dispatchEvent(event) - } - - close(): void { - this.closeCalls += 1 - this.readyState = this.CLOSED - } -} diff --git a/tests/schema-split.test.ts b/tests/schema-split.test.ts new file mode 100644 index 0000000..1fb4103 --- /dev/null +++ b/tests/schema-split.test.ts @@ -0,0 +1,21 @@ +import { readFileSync } from 'node:fs' +import { expect, it } from 'vitest' + +const commonSchema = readFileSync('./src/gen/zod-common.ts', 'utf8') +const eventSchema = readFileSync('./src/gen/zod-event.ts', 'utf8') +const apiSchema = readFileSync('./src/gen/zod-api.ts', 'utf8') + +it('partitions generated schemas by runtime entry point', () => { + expect(eventSchema).toContain('export const Event =') + expect(eventSchema).not.toContain('export const zodApiCategories =') + expect(apiSchema).toContain('export const zodApiCategories =') + expect(apiSchema).not.toContain('export const Event =') +}) + +it('moves shared declarations into a one-way common module', () => { + expect(commonSchema).toContain('export const IncomingMessage =') + expect(eventSchema).toContain('from \'./zod-common\'') + expect(apiSchema).toContain('from \'./zod-common\'') + expect(commonSchema).not.toContain('from \'./zod-event\'') + expect(commonSchema).not.toContain('from \'./zod-api\'') +}) diff --git a/tsdown.config.ts b/tsdown.config.ts index 45cabf2..4ad78d5 100644 --- a/tsdown.config.ts +++ b/tsdown.config.ts @@ -1,11 +1,13 @@ import { defineConfig } from 'tsdown' export default defineConfig({ - entry: ['src/index.ts', - ], + entry: { + index: 'src/index.ts', + event: 'src/event.ts', + }, dts: { tsgo: true, }, - exports: true, + exports: false, // ...config options }) From 3ff760af7519ffe508e1e05a910cf8757e180075 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E6=B0=B4=E5=8C=96?= <168127639+suisanka@users.noreply.github.com> Date: Sun, 2 Aug 2026 17:34:21 +0800 Subject: [PATCH 2/6] Improve client endpoint type readability --- src/client/endpoint.ts | 21 +++++++++++---------- tests/client.test.ts | 12 ++++++++++-- 2 files changed, 21 insertions(+), 12 deletions(-) diff --git a/src/client/endpoint.ts b/src/client/endpoint.ts index 41784fd..f8988c0 100644 --- a/src/client/endpoint.ts +++ b/src/client/endpoint.ts @@ -1,5 +1,6 @@ import type { MilkyFetch, MilkyFetchCreateOptions, MilkyFetchOptions } from '@/client/fetch' -import type { MilkyApiCategories, MilkyCamelCase, MilkyClientEndpointNames, MilkyRawEndpoints } from '@/types' +import type { ApiCategories, ApiEndpoints } from '@/gen/types' +import type { MilkyCamelCase, MilkyClientEndpointNames, MilkyRawEndpoints } from '@/types' import { createMilkyFetch } from '@/client/fetch' import { clientEndpointNames } from '@/types' @@ -57,21 +58,21 @@ function createProxy(options: MilkyFetchCreateOptions): any { export type MilkyClient = { readonly fetch: MilkyFetch } & { - readonly [K in keyof MilkyApiCategories]: { - readonly [E in keyof MilkyApiCategories[K]['apis'] as MilkyCamelCase]: - E extends keyof MilkyRawEndpoints - ? (...params: MilkyClientMethodParameters) => Promise> + readonly [K in keyof ApiCategories]: { + readonly [E in keyof ApiCategories[K] as MilkyCamelCase]: + E extends keyof MilkyRawEndpoints & keyof ApiEndpoints + ? (...params: MilkyClientMethodParameters) => Promise : never } & { readonly name: K } & {} } & {} -type MilkyClientMethodParameters - = Parameters extends [param: infer P] - ? [param: P, override?: MilkyFetchOptions] - : Parameters extends [param?: infer P] - ? [param?: P, override?: MilkyFetchOptions] +type MilkyClientMethodParameters + = Parameters extends [param: unknown] + ? [param: ApiEndpoints[T]['request_ZodInput'], override?: MilkyFetchOptions] + : Parameters extends [param?: unknown] + ? [param?: null | undefined, override?: MilkyFetchOptions] : never export function createMilkyClient(options: MilkyFetchCreateOptions): MilkyClient { diff --git a/tests/client.test.ts b/tests/client.test.ts index a2271e8..90de9cf 100644 --- a/tests/client.test.ts +++ b/tests/client.test.ts @@ -1,5 +1,10 @@ import type { MilkyFetchOptions } from '@/client/fetch' -import type { QuitGroupInput } from '@/index' +import type { + QuitGroupInput_ZodInput, + SendPrivateMessageInput_ZodInput, + SendPrivateMessageOutput, + SetGroupEssenceMessageInput_ZodInput, +} from '@/index' import { expect, expectTypeOf, it, vi } from 'vitest' import { createMilkyClient } from '@/client/endpoint' @@ -102,6 +107,9 @@ it('exposes grouped client methods with optional override options', () => { }) expectTypeOf(client.system.getLoginInfo).parameters.toEqualTypeOf<[(undefined | null)?, MilkyFetchOptions?]>() - expectTypeOf(client.group.quitGroup).parameters.toEqualTypeOf<[QuitGroupInput, MilkyFetchOptions?]>() + expectTypeOf(client.group.quitGroup).parameters.toEqualTypeOf<[QuitGroupInput_ZodInput, MilkyFetchOptions?]>() + expectTypeOf(client.group.setGroupEssenceMessage).parameters.toEqualTypeOf<[SetGroupEssenceMessageInput_ZodInput, MilkyFetchOptions?]>() + expectTypeOf(client.message.sendPrivateMessage).parameters.toEqualTypeOf<[SendPrivateMessageInput_ZodInput, MilkyFetchOptions?]>() + expectTypeOf(client.message.sendPrivateMessage).returns.toEqualTypeOf>() expect(client).not.toHaveProperty('event') }) From af2ca3528feccce78cd55a839edb9cc53ed85f30 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E6=B0=B4=E5=8C=96?= <168127639+suisanka@users.noreply.github.com> Date: Sun, 2 Aug 2026 17:58:20 +0800 Subject: [PATCH 3/6] refactor(client)!: align fetch types and validation options --- README.md | 10 ++++++-- examples/client.ts | 62 +++++++++++++++++++++++++++++++++++++++++++++ examples/event.ts | 28 ++++++++++++++++++++ examples/fetch.ts | 23 +++++++++++++++++ package.json | 1 + pnpm-lock.yaml | 17 +++++++++++++ src/client/fetch.ts | 37 ++++++++++++++------------- tests/fetch.test.ts | 36 ++++++++++++++++++-------- 8 files changed, 183 insertions(+), 31 deletions(-) create mode 100644 examples/client.ts create mode 100644 examples/event.ts create mode 100644 examples/fetch.ts diff --git a/README.md b/README.md index e632474..7bb47e4 100644 --- a/README.md +++ b/README.md @@ -72,14 +72,20 @@ import { createMilkyFetch } from '@saltify/milky-tea' const milkyFetch = createMilkyFetch({ baseURL: 'https://milky.example.com', - strict: false, + zod: false, }) const login = await milkyFetch('get_login_info', undefined) console.log(login.uin) ``` -`strict` 默认为 `true`。关闭后会跳过请求参数和响应数据的 zod 校验;也可以在单次请求的 override 里单独设置。 +`zod` 默认为 `true`。关闭后会跳过请求参数和响应数据的 Zod 校验;也可以在单次请求的 override 里单独设置。 + +## 示例 + +- [`examples/client.ts`](./examples/client.ts):收到好友私聊事件后,将其中的文本消息 echo 给发送者 +- [`examples/fetch.ts`](./examples/fetch.ts):使用底层 `createMilkyFetch` 调用原始 endpoint +- [`examples/event.ts`](./examples/event.ts):解析事件并通过 `event_type` 缩窄事件数据类型 ## 开发 diff --git a/examples/client.ts b/examples/client.ts new file mode 100644 index 0000000..83650ff --- /dev/null +++ b/examples/client.ts @@ -0,0 +1,62 @@ +import type { OutgoingSegment } from '@saltify/milky-tea' +import process from 'node:process' +import { createMilkyClient } from '@saltify/milky-tea' +import { resolveMilkyEvent } from '@saltify/milky-tea/event' +import { EventSource } from 'eventsource' + +const baseURL = process.env.MILKY_BASE_URL ?? 'https://milky.example.com' +const token = process.env.MILKY_TOKEN + +const client = createMilkyClient({ + baseURL, + token, +}) + +export async function handleEventPayload(payload: string): Promise { + const rawEvent: unknown = JSON.parse(payload) + const event = await resolveMilkyEvent(rawEvent) + + if ( + event.event_type !== 'message_receive' + || event.data.message_scene !== 'friend' + ) { + return + } + + // Incoming and outgoing segment unions are intentionally different. This + // example echoes the text segments that are valid in both directions. + const message: OutgoingSegment[] = event.data.segments + .filter(segment => segment.type === 'text') + + if (message.length === 0) { + return + } + + await client.message.sendPrivateMessage({ + user_id: event.data.sender_id, + message, + }) +} + +const eventURL = new URL('/event', baseURL) +if (token) { + eventURL.searchParams.set('access_token', token) +} + +const eventSource = new EventSource(eventURL) + +eventSource.onmessage = (event) => { + handleEventPayload(String(event.data)).catch(reportError) +} + +// EventSource reconnects automatically after recoverable connection failures. +eventSource.onerror = (event) => { + reportError(event.message ?? `EventSource error${event.code ? ` (${event.code})` : ''}`) +} + +process.once('SIGINT', () => eventSource.close()) +process.once('SIGTERM', () => eventSource.close()) + +function reportError(error: unknown): void { + process.stderr.write(`${String(error)}\n`) +} diff --git a/examples/event.ts b/examples/event.ts new file mode 100644 index 0000000..172e79d --- /dev/null +++ b/examples/event.ts @@ -0,0 +1,28 @@ +import { resolveMilkyEvent } from '@saltify/milky-tea/event' + +type EmitEvent = (summary: string, details?: unknown) => void + +export async function handleEventPayload( + payload: string, + emit: EmitEvent, +): Promise { + const rawEvent: unknown = JSON.parse(payload) + const event = await resolveMilkyEvent(rawEvent) + + // event_type narrows both the event and its data payload. + switch (event.event_type) { + case 'message_receive': + emit( + `Message ${event.data.message_seq} from ${event.data.sender_id}`, + event.data.segments, + ) + break + + case 'bot_offline': + emit(`Bot ${event.self_id} went offline: ${event.data.reason}`) + break + + default: + emit(`Received ${event.event_type}`) + } +} diff --git a/examples/fetch.ts b/examples/fetch.ts new file mode 100644 index 0000000..f93ad12 --- /dev/null +++ b/examples/fetch.ts @@ -0,0 +1,23 @@ +import type { GetFriendInfoOutput } from '@saltify/milky-tea' +import process from 'node:process' +import { createMilkyFetch } from '@saltify/milky-tea' + +async function main(): Promise { + const milkyFetch = createMilkyFetch({ + baseURL: 'https://milky.example.com', + token: process.env.MILKY_TOKEN, + }) + + // Use the raw snake_case endpoint name when grouped client methods are not + // suitable. Request and response types are still inferred from the endpoint. + const result: GetFriendInfoOutput = await milkyFetch('get_friend_info', { + user_id: 10001, + }) + + process.stdout.write(`${JSON.stringify(result.friend, null, 2)}\n`) +} + +main().catch((error: unknown) => { + process.stderr.write(`${String(error)}\n`) + process.exitCode = 1 +}) diff --git a/package.json b/package.json index 0d590b8..b4d97ee 100644 --- a/package.json +++ b/package.json @@ -55,6 +55,7 @@ "bumpp": "^10.4.1", "change-case": "^5.4.4", "eslint": "^10.8.0", + "eventsource": "^4.1.0", "jiti": "^2.7.0", "tsdown": "^0.21.10", "typescript": "^5.9.3", diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index bde9ae5..f6ef04e 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -29,6 +29,9 @@ importers: eslint: specifier: ^10.8.0 version: 10.8.0(jiti@2.7.0) + eventsource: + specifier: ^4.1.0 + version: 4.1.0 jiti: specifier: ^2.7.0 version: 2.7.0 @@ -1220,6 +1223,14 @@ packages: resolution: {integrity: sha512-kVscqXk4OCp68SZ0dkgEKVi6/8ij300KBWTJq32P/dYeWTSwK41WyTxalN1eRmA5Z9UU/LX9D7FWSmV9SAYx6g==} engines: {node: '>=0.10.0'} + eventsource-parser@3.1.0: + resolution: {integrity: sha512-kJezFj9YFAMLeORyi7aCLxLbD5/qWMQnoMVlVPyHIll7lgRJCc3JVln9Vgl9nwQi0YkMnhdGTMNn7CkRRAptMg==} + engines: {node: '>=18.0.0'} + + eventsource@4.1.0: + resolution: {integrity: sha512-2GuF51iuHX6A9xdTccMTsNb7VO0lHZihApxhvQzJB5A03DvHDd2FQepodbMaztPBmBcE/ox7o2gqaxGhYB9LhQ==} + engines: {node: '>=20.0.0'} + expect-type@1.4.0: resolution: {integrity: sha512-KfYbmpRm0VbLjEvVa9yGwCi9GI34xvi7A/HXYWQO65CSD2u3MczUJSuwXKFIxlGsgBQizV9q5J9NHj4VG0n+pA==} engines: {node: '>=12.0.0'} @@ -3363,6 +3374,12 @@ snapshots: esutils@2.0.3: {} + eventsource-parser@3.1.0: {} + + eventsource@4.1.0: + dependencies: + eventsource-parser: 3.1.0 + expect-type@1.4.0: {} exsolve@1.1.1: {} diff --git a/src/client/fetch.ts b/src/client/fetch.ts index f0a83e2..9abb8a1 100644 --- a/src/client/fetch.ts +++ b/src/client/fetch.ts @@ -1,10 +1,11 @@ +import type { ApiEndpoints } from '@/gen/types' import type { MilkyProto, MilkyProtoStruct, MilkyRawEndpoints } from '@/types' import { createMilkyProto, rawEndpointNames } from '@/types' import { joinURL, withTimeout } from '@/utils' export interface MilkyFetchOptions { readonly baseURL?: string | URL - readonly strict?: boolean + readonly zod?: boolean readonly token?: string readonly timeout?: number | false readonly request?: Omit @@ -16,17 +17,17 @@ export type MilkyFetchCreateOptions = Omit & { } export interface MilkyFetch { - ( + ( name: T, ...args: MilkyFetchParameters - ): Promise> + ): Promise } -type MilkyFetchParameters - = Parameters extends [param: infer P] - ? [param: P, override?: MilkyFetchOptions] - : Parameters extends [param?: infer P] - ? [param?: P, override?: MilkyFetchOptions] +type MilkyFetchParameters + = Parameters extends [param: unknown] + ? [param: ApiEndpoints[T]['request_ZodInput'], override?: MilkyFetchOptions] + : Parameters extends [param?: unknown] + ? [param?: null | undefined, override?: MilkyFetchOptions] : never interface MilkyApiResponse { @@ -82,16 +83,16 @@ export function createMilkyFetch(options: MilkyFetchCreateOptions): MilkyFetch { const defaultFetch = options.fetch ?? globalThis.fetch.bind(globalThis) - return async function fetch( + return async function fetch( name: T, ...args: MilkyFetchParameters - ): Promise> { + ): Promise { let [params, override] = args as [unknown, MilkyFetchOptions | undefined] - const strict = override?.strict ?? options.strict ?? true + const zod = override?.zod ?? options.zod ?? true let paramStruct: MilkyProtoStruct | null | undefined let responseStruct: MilkyProtoStruct | null | undefined - if (strict) { + if (zod) { if (!rawEndpointNames.has(String(name))) { throw new Error(`milky: unknown endpoint ${String(name)}`) } @@ -102,7 +103,7 @@ export function createMilkyFetch(options: MilkyFetchCreateOptions): MilkyFetch { } } - if (strict && paramStruct != null) { + if (zod && paramStruct != null) { const paramParseResult = await paramStruct.safeParseAsync(params) if (!paramParseResult.success) { @@ -167,10 +168,10 @@ export function createMilkyFetch(options: MilkyFetchCreateOptions): MilkyFetch { }, ) - let payload: MilkyApiResponse> + let payload: MilkyApiResponse try { - payload = await response.json() as MilkyApiResponse> + payload = await response.json() as MilkyApiResponse } catch (error) { throw new Error(`milky: failed to parse response for ${String(name)}`, { cause: error }) @@ -180,8 +181,8 @@ export function createMilkyFetch(options: MilkyFetchCreateOptions): MilkyFetch { throw new Error(payload.message ?? `milky: invoke ${String(name)} failed: ${payload.message} (${payload.retcode})`) } - if (!strict || responseStruct == null) { - return payload.data as ReturnType + if (!zod || responseStruct == null) { + return payload.data as ApiEndpoints[T]['response'] } const responseParseResult = await responseStruct.safeParseAsync(payload.data) @@ -190,6 +191,6 @@ export function createMilkyFetch(options: MilkyFetchCreateOptions): MilkyFetch { throw new Error(`milky: failed to parse response for ${String(name)}: ${responseParseResult.error.message}`) } - return responseParseResult.data as ReturnType + return responseParseResult.data as ApiEndpoints[T]['response'] } } diff --git a/tests/fetch.test.ts b/tests/fetch.test.ts index 73c2e9c..33b6d20 100644 --- a/tests/fetch.test.ts +++ b/tests/fetch.test.ts @@ -1,4 +1,9 @@ -import { afterEach, expect, it, vi } from 'vitest' +import type { + GetFriendInfoInput_ZodInput, + GetFriendInfoOutput, + GetLoginInfoOutput, +} from '@/index' +import { afterEach, expect, expectTypeOf, it, vi } from 'vitest' import { createMilkyFetch } from '@/client/fetch' import { sleep } from './helpers/async' @@ -38,7 +43,10 @@ it('uses the global fetch implementation when no local fetch is provided', async baseURL: 'https://example.com', }) - await expect(milkyFetch('get_login_info', undefined)).resolves.toEqual({ + const pending = milkyFetch('get_login_info') + expectTypeOf(pending).toEqualTypeOf>() + + await expect(pending).resolves.toEqual({ uin: 10001, nickname: 'bot', }) @@ -123,7 +131,10 @@ it('allows per-request overrides for baseURL, token, fetch and headers', async ( }, }) - await expect(milkyFetch('get_friend_info', { user_id: 10001 } as never, { + const input = { + user_id: 10001, + } satisfies GetFriendInfoInput_ZodInput + const pending = milkyFetch('get_friend_info', input, { baseURL: 'https://override.example.com/', token: 'override-token', fetch: overrideFetch, @@ -132,7 +143,10 @@ it('allows per-request overrides for baseURL, token, fetch and headers', async ( 'x-sdk': 'milky', }, }, - })).resolves.toEqual({ + }) + expectTypeOf(pending).toEqualTypeOf>() + + await expect(pending).resolves.toEqual({ friend: { user_id: 10001, nickname: 'friend', @@ -196,7 +210,7 @@ it('throws when the endpoint name is unknown', async () => { expect(fetchMock).not.toHaveBeenCalled() }) -it('skips unknown endpoint checks when strict is disabled', async () => { +it('skips unknown endpoint checks when zod is disabled', async () => { const fetchMock = vi.fn(async (request: Request) => { expect(request.url).toBe('https://example.com/api/unknown_endpoint') @@ -211,7 +225,7 @@ it('skips unknown endpoint checks when strict is disabled', async () => { const milkyFetch = createMilkyFetch({ baseURL: 'https://example.com', - strict: false, + zod: false, fetch: fetchMock, }) @@ -232,7 +246,7 @@ it('validates request params before issuing the request', async () => { expect(fetchMock).not.toHaveBeenCalled() }) -it('skips request and response validation when strict is disabled on the client', async () => { +it('skips request and response validation when zod is disabled on the client', async () => { const fetchMock = vi.fn(async (request: Request) => { expect(await request.text()).toBe('{}') @@ -248,7 +262,7 @@ it('skips request and response validation when strict is disabled on the client' const milkyFetch = createMilkyFetch({ baseURL: 'https://example.com', - strict: false, + zod: false, fetch: fetchMock, }) @@ -291,7 +305,7 @@ it('skips zod validation when zod schemas are unavailable', async () => { expect(fetchMock).toHaveBeenCalledOnce() }) -it('allows overriding strict per request', async () => { +it('allows overriding zod per request', async () => { const fetchMock = vi.fn(async () => createJsonResponse({ status: 'ok', retcode: 0, @@ -303,7 +317,7 @@ it('allows overriding strict per request', async () => { const milkyFetch = createMilkyFetch({ baseURL: 'https://example.com', - strict: false, + zod: false, fetch: fetchMock, }) @@ -312,7 +326,7 @@ it('allows overriding strict per request', async () => { nickname: 'bot', }) await expect(milkyFetch('get_login_info', undefined, { - strict: true, + zod: true, })).rejects.toThrow('milky: failed to parse response for get_login_info') expect(fetchMock).toHaveBeenCalledTimes(2) }) From b3e08ae5ec12b83a65187d3c878badb5ab687b9b Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E6=B0=B4=E5=8C=96?= <168127639+suisanka@users.noreply.github.com> Date: Sun, 2 Aug 2026 18:05:19 +0800 Subject: [PATCH 4/6] docs(api): document public interfaces --- src/client/endpoint.ts | 12 ++++++++++++ src/client/fetch.ts | 32 ++++++++++++++++++++++++++++++++ src/event.ts | 7 +++++++ src/index.ts | 12 +++++++++++- src/types.ts | 2 ++ 5 files changed, 64 insertions(+), 1 deletion(-) diff --git a/src/client/endpoint.ts b/src/client/endpoint.ts index f8988c0..1bd79d4 100644 --- a/src/client/endpoint.ts +++ b/src/client/endpoint.ts @@ -55,7 +55,9 @@ function createProxy(options: MilkyFetchCreateOptions): any { }) } +/** A category-based, camel-case client for all Milky API endpoints. */ export type MilkyClient = { + /** Low-level endpoint caller configured with the same defaults as this client. */ readonly fetch: MilkyFetch } & { readonly [K in keyof ApiCategories]: { @@ -64,6 +66,7 @@ export type MilkyClient = { ? (...params: MilkyClientMethodParameters) => Promise : never } & { + /** Snake-case Milky API category name. */ readonly name: K } & {} } & {} @@ -75,6 +78,15 @@ type MilkyClientMethodParameters + /** Custom Fetch API implementation. Defaults to `globalThis.fetch`. */ readonly fetch?: (request: Request) => Promise } +/** Options used to create a {@link MilkyFetch} instance. */ export type MilkyFetchCreateOptions = Omit & { + /** Base URL of the Milky implementation, including any path prefix. */ readonly baseURL: string | URL } +/** A typed function for invoking Milky API endpoints by their protocol names. */ export interface MilkyFetch { + /** + * Invokes a Milky API endpoint. + * + * @param name - Snake-case endpoint name from the Milky protocol. + * @param args - Endpoint parameters followed by optional per-request overrides. + * @returns The endpoint response data. + */ ( name: T, ...args: MilkyFetchParameters @@ -76,6 +101,13 @@ async function resolveMilkyProto(): Promise { return milkyProtoPromise } +/** + * Creates a typed, low-level Milky API caller. + * + * @param options - Default connection, validation, and request options. + * @returns A function that invokes endpoints by their snake-case protocol names. + * @throws If no Fetch API implementation is available. + */ export function createMilkyFetch(options: MilkyFetchCreateOptions): MilkyFetch { if (options.fetch == null && globalThis.fetch == null) { throw new Error('milky: fetch is not provided') diff --git a/src/event.ts b/src/event.ts index 35084ff..fb4c16b 100644 --- a/src/event.ts +++ b/src/event.ts @@ -46,6 +46,13 @@ async function resolveMilkyEventSchema(): Promise { return milkyEventSchemaPromise } +/** + * Validates and resolves an unknown value as a Milky event. + * + * @param obj - Deserialized event payload received from a Milky implementation. + * @returns A validated Milky event with its discriminated union type preserved. + * @throws If Zod is unavailable or the payload is not a valid Milky event. + */ export async function resolveMilkyEvent(obj: unknown): Promise { const schema = await resolveMilkyEventSchema() const result = await schema.safeParseAsync(obj) diff --git a/src/index.ts b/src/index.ts index 79c745c..3b63b74 100644 --- a/src/index.ts +++ b/src/index.ts @@ -1,5 +1,15 @@ +import { + milkyPackageVersion as generatedMilkyPackageVersion, + milkyVersion as generatedMilkyVersion, +} from './gen/types' + export * from './client' export { resolveMilkyEvent } from './event' -export { milkyPackageVersion, milkyVersion } from './gen/types' export type * from './gen/types' export type { MilkyRawEndpointName, MilkyRawEndpoints } from './types' + +/** Milky protocol version targeted by this SDK. */ +export const milkyVersion = generatedMilkyVersion + +/** Full Milky protocol package version used to generate the SDK types. */ +export const milkyPackageVersion = generatedMilkyPackageVersion diff --git a/src/types.ts b/src/types.ts index 0e9bef4..21fdc89 100644 --- a/src/types.ts +++ b/src/types.ts @@ -38,10 +38,12 @@ type UnionToIntersection = (T extends unknown ? (value: T) => void : never) e ? R : never +/** Map of snake-case Milky endpoint names to their request and response signatures. */ export type MilkyRawEndpoints = UnionToIntersection<{ [C in ApiCategoryName]: RawEndpointsForCategory }[ApiCategoryName]> +/** Name of any endpoint defined by the supported Milky protocol version. */ export type MilkyRawEndpointName = keyof MilkyRawEndpoints export type MilkyApiCategories = ApiCategories From 5ad7a06c4aaf4c7da4229f709e38dae2901ad758 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E6=B0=B4=E5=8C=96?= <168127639+suisanka@users.noreply.github.com> Date: Sun, 2 Aug 2026 18:13:56 +0800 Subject: [PATCH 5/6] chore: tsconfig --- tsconfig.json | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/tsconfig.json b/tsconfig.json index 36492be..216b5fa 100644 --- a/tsconfig.json +++ b/tsconfig.json @@ -19,5 +19,6 @@ "isolatedModules": true, "verbatimModuleSyntax": true, "skipLibCheck": true - } + }, + "exclude": ["examples/**/*.ts"] } From e941f6ab40058caf60eeae27896fe169fc7140b7 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E6=B0=B4=E5=8C=96?= <168127639+suisanka@users.noreply.github.com> Date: Sun, 2 Aug 2026 18:22:30 +0800 Subject: [PATCH 6/6] fix(client): preserve schema-less response types --- examples/client.ts | 4 ++-- src/client/endpoint.ts | 4 ++-- src/client/fetch.ts | 14 +++++++------- src/types.ts | 2 ++ tests/client.test.ts | 1 + tests/fetch.test.ts | 14 ++++++++++++++ 6 files changed, 28 insertions(+), 11 deletions(-) diff --git a/examples/client.ts b/examples/client.ts index 83650ff..ab274f2 100644 --- a/examples/client.ts +++ b/examples/client.ts @@ -45,9 +45,9 @@ if (token) { const eventSource = new EventSource(eventURL) -eventSource.onmessage = (event) => { +eventSource.addEventListener('milky_event', (event) => { handleEventPayload(String(event.data)).catch(reportError) -} +}) // EventSource reconnects automatically after recoverable connection failures. eventSource.onerror = (event) => { diff --git a/src/client/endpoint.ts b/src/client/endpoint.ts index 1bd79d4..33178c2 100644 --- a/src/client/endpoint.ts +++ b/src/client/endpoint.ts @@ -1,6 +1,6 @@ import type { MilkyFetch, MilkyFetchCreateOptions, MilkyFetchOptions } from '@/client/fetch' import type { ApiCategories, ApiEndpoints } from '@/gen/types' -import type { MilkyCamelCase, MilkyClientEndpointNames, MilkyRawEndpoints } from '@/types' +import type { MilkyCamelCase, MilkyClientEndpointNames, MilkyEndpointResponse, MilkyRawEndpoints } from '@/types' import { createMilkyFetch } from '@/client/fetch' import { clientEndpointNames } from '@/types' @@ -63,7 +63,7 @@ export type MilkyClient = { readonly [K in keyof ApiCategories]: { readonly [E in keyof ApiCategories[K] as MilkyCamelCase]: E extends keyof MilkyRawEndpoints & keyof ApiEndpoints - ? (...params: MilkyClientMethodParameters) => Promise + ? (...params: MilkyClientMethodParameters) => Promise> : never } & { /** Snake-case Milky API category name. */ diff --git a/src/client/fetch.ts b/src/client/fetch.ts index 1cbbf8d..1e005c2 100644 --- a/src/client/fetch.ts +++ b/src/client/fetch.ts @@ -1,5 +1,5 @@ import type { ApiEndpoints } from '@/gen/types' -import type { MilkyProto, MilkyProtoStruct, MilkyRawEndpoints } from '@/types' +import type { MilkyEndpointResponse, MilkyProto, MilkyProtoStruct, MilkyRawEndpoints } from '@/types' import { createMilkyProto, rawEndpointNames } from '@/types' import { joinURL, withTimeout } from '@/utils' @@ -45,7 +45,7 @@ export interface MilkyFetch { ( name: T, ...args: MilkyFetchParameters - ): Promise + ): Promise> } type MilkyFetchParameters @@ -118,7 +118,7 @@ export function createMilkyFetch(options: MilkyFetchCreateOptions): MilkyFetch { return async function fetch( name: T, ...args: MilkyFetchParameters - ): Promise { + ): Promise> { let [params, override] = args as [unknown, MilkyFetchOptions | undefined] const zod = override?.zod ?? options.zod ?? true let paramStruct: MilkyProtoStruct | null | undefined @@ -200,10 +200,10 @@ export function createMilkyFetch(options: MilkyFetchCreateOptions): MilkyFetch { }, ) - let payload: MilkyApiResponse + let payload: MilkyApiResponse> try { - payload = await response.json() as MilkyApiResponse + payload = await response.json() as MilkyApiResponse> } catch (error) { throw new Error(`milky: failed to parse response for ${String(name)}`, { cause: error }) @@ -214,7 +214,7 @@ export function createMilkyFetch(options: MilkyFetchCreateOptions): MilkyFetch { } if (!zod || responseStruct == null) { - return payload.data as ApiEndpoints[T]['response'] + return payload.data as MilkyEndpointResponse } const responseParseResult = await responseStruct.safeParseAsync(payload.data) @@ -223,6 +223,6 @@ export function createMilkyFetch(options: MilkyFetchCreateOptions): MilkyFetch { throw new Error(`milky: failed to parse response for ${String(name)}: ${responseParseResult.error.message}`) } - return responseParseResult.data as ApiEndpoints[T]['response'] + return responseParseResult.data as MilkyEndpointResponse } } diff --git a/src/types.ts b/src/types.ts index 21fdc89..f1b380d 100644 --- a/src/types.ts +++ b/src/types.ts @@ -46,6 +46,8 @@ export type MilkyRawEndpoints = UnionToIntersection<{ /** Name of any endpoint defined by the supported Milky protocol version. */ export type MilkyRawEndpointName = keyof MilkyRawEndpoints +export type MilkyEndpointResponse = ReturnType + export type MilkyApiCategories = ApiCategories export type MilkyCamelCase = CamelCase diff --git a/tests/client.test.ts b/tests/client.test.ts index 90de9cf..aabdfca 100644 --- a/tests/client.test.ts +++ b/tests/client.test.ts @@ -110,6 +110,7 @@ it('exposes grouped client methods with optional override options', () => { expectTypeOf(client.group.quitGroup).parameters.toEqualTypeOf<[QuitGroupInput_ZodInput, MilkyFetchOptions?]>() expectTypeOf(client.group.setGroupEssenceMessage).parameters.toEqualTypeOf<[SetGroupEssenceMessageInput_ZodInput, MilkyFetchOptions?]>() expectTypeOf(client.message.sendPrivateMessage).parameters.toEqualTypeOf<[SendPrivateMessageInput_ZodInput, MilkyFetchOptions?]>() + expectTypeOf(client.group.quitGroup).returns.toEqualTypeOf>() expectTypeOf(client.message.sendPrivateMessage).returns.toEqualTypeOf>() expect(client).not.toHaveProperty('event') }) diff --git a/tests/fetch.test.ts b/tests/fetch.test.ts index 33b6d20..7227d7f 100644 --- a/tests/fetch.test.ts +++ b/tests/fetch.test.ts @@ -92,6 +92,20 @@ it('posts JSON requests to the API endpoint and returns data payloads', async () expect(fetchMock).toHaveBeenCalledOnce() }) +it('types schema-less endpoint responses as void', async () => { + const milkyFetch = createMilkyFetch({ + baseURL: 'https://example.com', + fetch: async () => createJsonResponse({ + status: 'ok', + retcode: 0, + }), + }) + + const pending = milkyFetch('quit_group', { group_id: 10001 }) + expectTypeOf(pending).toEqualTypeOf>() + await expect(pending).resolves.toBeUndefined() +}) + it('allows per-request overrides for baseURL, token, fetch and headers', async () => { const defaultFetch = vi.fn() const overrideFetch = vi.fn(async (request: Request) => {