feat(llm): route adapters by provider
This commit is contained in:
@@ -1,181 +1,98 @@
|
||||
/**
|
||||
* `PiAiAdapter`: the `@earendil-works/pi-ai`-backed implementation of the
|
||||
* harness LLM seam, pointed at a DeepSeek (OpenAI-compatible) endpoint.
|
||||
*
|
||||
* This adapter exists as a design-verification twin of
|
||||
* `@deepseek-ai/dsh-llm-deepseek`: same models, same wire protocol,
|
||||
* completely different internals (a unified LLM library with its own event
|
||||
* vocabulary vs hand-rolled fetch/SSE). Anything the StreamChunk protocol
|
||||
* cannot express for BOTH implementations is a core-vocabulary bug.
|
||||
* Generic pi-ai-backed implementation of the Harness LLM seam.
|
||||
*
|
||||
* @module dsh-llm-pi-ai/adapter
|
||||
*/
|
||||
|
||||
import { stream as piStream } from '@earendil-works/pi-ai'
|
||||
import type { Model } from '@earendil-works/pi-ai'
|
||||
import { attributionHeaders, LlmAdapter } from '@deepseek-ai/dsh-llm'
|
||||
import { CallId } from '@deepseek-ai/dsh-llm'
|
||||
import {
|
||||
getModels,
|
||||
streamSimple,
|
||||
} from '@earendil-works/pi-ai'
|
||||
import type {
|
||||
Api,
|
||||
KnownProvider,
|
||||
Model,
|
||||
SimpleStreamOptions,
|
||||
} from '@earendil-works/pi-ai'
|
||||
import { attributionHeaders, LlmAdapter, LlmError } from '@deepseek-ai/dsh-llm'
|
||||
import type { GenerateOptions, StreamChunk } from '@deepseek-ai/dsh-llm'
|
||||
import { toPiContext, toStreamChunks } from './convert.ts'
|
||||
import type { PiAiProviderProfile } from './config.ts'
|
||||
import { toPiContext } from './context.ts'
|
||||
import { toStreamChunks } from './stream.ts'
|
||||
|
||||
/** Reasoning levels surfaced by this adapter (DeepSeek wire: high|max). */
|
||||
export type PiAiReasoning = 'off' | 'high' | 'xhigh'
|
||||
|
||||
/** Constructor options for {@link PiAiAdapter}; the plugin's `apply` resolves them from Config + environment. */
|
||||
/** Constructor options for {@link PiAiAdapter}. */
|
||||
export interface PiAiAdapterOptions {
|
||||
/** Bearer token pi-ai sends on every request. */
|
||||
apiKey: string
|
||||
/** Endpoint base; `/chat/completions` is appended. */
|
||||
baseURL: string
|
||||
/** Thinking level applied to every request ('off' disables thinking). */
|
||||
reasoning?: PiAiReasoning | undefined
|
||||
/** Validated provider profiles this adapter instance owns. */
|
||||
profiles: readonly PiAiProviderProfile[]
|
||||
}
|
||||
|
||||
/**
|
||||
* Build the inline pi-ai model descriptor for one DeepSeek model name.
|
||||
* @param modelId - harness model name; sent verbatim on the wire.
|
||||
* @param options - adapter options; only `baseURL` is read here (key and reasoning apply per request, not per descriptor).
|
||||
* @returns a descriptor with every DeepSeek compat flag explicit — pi-ai's URL-based auto-detection is never relied on.
|
||||
* Resolve a catalog model dynamically and apply only the configured endpoint
|
||||
* override, preserving the catalog's API/capability/compatibility metadata.
|
||||
*/
|
||||
export function buildModel(modelId: string, options: PiAiAdapterOptions): Model<'openai-completions'> {
|
||||
function resolveModel(profile: PiAiProviderProfile, modelId: string): Model<Api> {
|
||||
const model = getModels(profile.provider as KnownProvider).find(candidate => candidate.id === modelId) as Model<Api> | undefined
|
||||
if (model === undefined) {
|
||||
throw new LlmError(`pi-ai provider "${profile.provider}" has no catalog model "${modelId}"`, 'UNKNOWN_MODEL')
|
||||
}
|
||||
return profile.baseURL === undefined ? model : { ...model, baseUrl: profile.baseURL }
|
||||
}
|
||||
|
||||
/** Copy profile stream knobs into pi-ai's common option vocabulary. */
|
||||
function profileOptions(profile: PiAiProviderProfile): SimpleStreamOptions {
|
||||
return {
|
||||
id: modelId,
|
||||
name: modelId,
|
||||
api: 'openai-completions',
|
||||
provider: 'deepseek',
|
||||
baseUrl: options.baseURL,
|
||||
// Always true: pi-ai only emits the DeepSeek `thinking` field for
|
||||
// reasoning-capable models, deriving enabled/disabled from whether a
|
||||
// reasoningEffort option is passed. DeepSeek's provider default is
|
||||
// ENABLED, so 'off' must send an explicit {type: 'disabled'} — which
|
||||
// requires this flag to stay on.
|
||||
reasoning: true,
|
||||
// DeepSeek's official effort levels: high|max (xhigh maps to max).
|
||||
thinkingLevelMap: { minimal: null, low: null, medium: null, high: 'high', xhigh: 'max' },
|
||||
input: ['text'],
|
||||
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 },
|
||||
contextWindow: 128_000,
|
||||
maxTokens: 64_000,
|
||||
compat: {
|
||||
// Auto-detection only fires for *.deepseek.com base URLs; the internal
|
||||
// endpoint (and test mocks) need these set explicitly.
|
||||
thinkingFormat: 'deepseek',
|
||||
requiresReasoningContentOnAssistantMessages: true,
|
||||
supportsReasoningEffort: true,
|
||||
// DeepSeek documents max_tokens (not OpenAI's max_completion_tokens).
|
||||
maxTokensField: 'max_tokens',
|
||||
},
|
||||
...profile.apiKey === undefined ? {} : { apiKey: profile.apiKey },
|
||||
...profile.reasoning === undefined ? {} : { reasoning: profile.reasoning },
|
||||
...profile.thinkingBudgets === undefined ? {} : { thinkingBudgets: profile.thinkingBudgets },
|
||||
...profile.cacheRetention === undefined ? {} : { cacheRetention: profile.cacheRetention },
|
||||
...profile.transport === undefined ? {} : { transport: profile.transport },
|
||||
...profile.timeoutMs === undefined ? {} : { timeoutMs: profile.timeoutMs },
|
||||
...profile.websocketConnectTimeoutMs === undefined ? {} : { websocketConnectTimeoutMs: profile.websocketConnectTimeoutMs },
|
||||
...profile.maxRetries === undefined ? {} : { maxRetries: profile.maxRetries },
|
||||
...profile.maxRetryDelayMs === undefined ? {} : { maxRetryDelayMs: profile.maxRetryDelayMs },
|
||||
}
|
||||
}
|
||||
|
||||
type Payload = {
|
||||
tools?: { function?: { strict?: unknown } }[]
|
||||
messages?: {
|
||||
role?: unknown
|
||||
tool_calls?: { id?: unknown; function?: { arguments?: unknown } }[]
|
||||
}[]
|
||||
reasoning_effort?: unknown
|
||||
stop?: unknown
|
||||
}
|
||||
|
||||
function rawToolArguments(options: GenerateOptions): Map<CallId, string> {
|
||||
const raw = new Map<CallId, string>()
|
||||
for (const message of options.messages) {
|
||||
if (message.role !== 'assistant') continue
|
||||
for (const block of message.content) {
|
||||
if (block.type === 'tool-call') raw.set(block.id, block.arguments)
|
||||
}
|
||||
}
|
||||
return raw
|
||||
}
|
||||
|
||||
function patchPayload(payload: unknown, options: GenerateOptions, reasoning: PiAiReasoning | undefined): unknown {
|
||||
/* v8 ignore next -- pi-ai onPayload always receives an object; tolerate unusual future hooks defensively */
|
||||
if (typeof payload !== 'object' || payload === null) return payload
|
||||
const body = payload as Payload
|
||||
|
||||
if (reasoning === undefined) {
|
||||
delete body.reasoning_effort
|
||||
}
|
||||
if (options.stop !== undefined) {
|
||||
body.stop = options.stop
|
||||
}
|
||||
|
||||
// pi-ai stamps its own `strict` default on every serialized tool; the
|
||||
// harness tool contract has no strict field and the hand-rolled twin sends
|
||||
// none, so scrub it for wire parity.
|
||||
for (const tool of body.tools ?? []) {
|
||||
/* v8 ignore next -- malformed pi-ai payload guard: real tool entries always carry function */
|
||||
if (tool.function === undefined) continue
|
||||
delete tool.function.strict
|
||||
}
|
||||
|
||||
const rawById = rawToolArguments(options)
|
||||
/* v8 ignore next -- defensive for non-chat payloads; OpenAI chat payloads always carry messages */
|
||||
for (const message of body.messages ?? []) {
|
||||
if (message.role !== 'assistant') continue
|
||||
/* v8 ignore next -- assistant messages without tool_calls need no raw-argument patch */
|
||||
for (const call of message.tool_calls ?? []) {
|
||||
/* v8 ignore next -- malformed pi-ai payload guard: real tool calls always carry a string id */
|
||||
if (typeof call.id !== 'string') continue
|
||||
const raw = rawById.get(CallId(call.id))
|
||||
/* v8 ignore next -- pi-ai always emits a function object for assistant tool_calls; guard malformed payloads defensively */
|
||||
if (raw !== undefined && call.function !== undefined) call.function.arguments = raw
|
||||
}
|
||||
}
|
||||
|
||||
return body
|
||||
}
|
||||
|
||||
/**
|
||||
* pi-ai-backed adapter. One instance serves every registered model name.
|
||||
*
|
||||
* Implementation notes:
|
||||
* - `onPayload` patches provider payload details pi-ai cannot express directly:
|
||||
* stop sequences, scrubbing pi-ai's own per-tool `strict` default (the
|
||||
* hand-rolled twin sends no such field), omitted reasoning effort, and raw
|
||||
* replayed tool-call arguments.
|
||||
* - pi-ai reports request failures as in-stream error events; convert.ts
|
||||
* maps them to `finish {kind:'error'|'aborted'}` chunks rather than
|
||||
* throwing — both are sanctioned StreamChunk error paths.
|
||||
* pi-ai-backed multi-provider adapter. Model descriptors are resolved for each
|
||||
* request, so models need not be registered during the Cordis lifecycle.
|
||||
*/
|
||||
export class PiAiAdapter extends LlmAdapter {
|
||||
constructor(private readonly options: PiAiAdapterOptions) {
|
||||
private readonly profiles: ReadonlyMap<string, PiAiProviderProfile>
|
||||
|
||||
constructor(options: PiAiAdapterOptions) {
|
||||
super()
|
||||
this.profiles = new Map(options.profiles.map(profile => [profile.provider, profile]))
|
||||
}
|
||||
|
||||
async * stream(options: GenerateOptions): AsyncIterable<StreamChunk> {
|
||||
const model = buildModel(options.model, this.options)
|
||||
// Undefined config means "provider default" (DeepSeek: thinking ENABLED),
|
||||
// matching llm-deepseek's omission semantics. pi-ai derives the wire
|
||||
// thinking toggle from whether reasoningEffort is passed, so undefined maps
|
||||
// internally to 'high' to get `thinking: enabled`; patchPayload then removes
|
||||
// `reasoning_effort` so the provider chooses its default effort.
|
||||
const reasoning = this.options.reasoning ?? 'high'
|
||||
if (options.stop !== undefined) {
|
||||
throw new LlmError('llm-pi-ai does not support GenerateOptions.stop', 'UNSUPPORTED_OPTION')
|
||||
}
|
||||
const profile = this.profiles.get(options.provider)
|
||||
if (profile === undefined) {
|
||||
throw new LlmError(`pi-ai adapter does not own provider "${options.provider}"`, 'NO_ADAPTER')
|
||||
}
|
||||
const model = resolveModel(profile, options.model)
|
||||
|
||||
// pi-ai's event stream has no iterator-return cancellation hook: if our
|
||||
// consumer stops early (break / loop abort), the underlying HTTP stream
|
||||
// would keep draining. Chain an internal controller onto the caller's
|
||||
// signal and abort it when this generator exits for any reason.
|
||||
// pi-ai's event stream has no iterator-return cancellation hook: abort its
|
||||
// provider stream when our consumer exits early as well as on caller abort.
|
||||
const controller = new AbortController()
|
||||
const onCallerAbort = (): void => { controller.abort(options.signal?.reason) }
|
||||
if (options.signal?.aborted) controller.abort(options.signal.reason)
|
||||
else options.signal?.addEventListener('abort', onCallerAbort, { once: true })
|
||||
|
||||
try {
|
||||
const events = piStream(model, toPiContext(options), {
|
||||
apiKey: this.options.apiKey,
|
||||
// pi-ai merges caller headers last over its provider defaults, so the
|
||||
// harness attribution always reaches the wire.
|
||||
headers: attributionHeaders(),
|
||||
...options.temperature !== undefined ? { temperature: options.temperature } : {},
|
||||
...options.maxTokens !== undefined ? { maxTokens: options.maxTokens } : {},
|
||||
const events = streamSimple(model, toPiContext(options), {
|
||||
...profileOptions(profile),
|
||||
...options.temperature === undefined ? {} : { temperature: options.temperature },
|
||||
...options.maxTokens === undefined ? {} : { maxTokens: options.maxTokens },
|
||||
...options.sessionId === undefined ? {} : { sessionId: String(options.sessionId) },
|
||||
signal: controller.signal,
|
||||
...reasoning !== 'off' ? { reasoningEffort: reasoning } : {},
|
||||
onPayload: payload => patchPayload(payload, options, this.options.reasoning),
|
||||
maxRetries: 0,
|
||||
// Profile headers are deployment-owned; attribution names are
|
||||
// Harness-owned and therefore win collisions.
|
||||
headers: { ...profile.headers, ...attributionHeaders() },
|
||||
})
|
||||
|
||||
yield* toStreamChunks(events)
|
||||
} finally {
|
||||
options.signal?.removeEventListener('abort', onCallerAbort)
|
||||
|
||||
99
packages/llm/llm-pi-ai/src/config.ts
Normal file
99
packages/llm/llm-pi-ai/src/config.ts
Normal file
@@ -0,0 +1,99 @@
|
||||
/**
|
||||
* Configuration schema and provider-profile validation for the pi-ai adapter.
|
||||
*
|
||||
* @module dsh-llm-pi-ai/config
|
||||
*/
|
||||
|
||||
import { getProviders } from '@earendil-works/pi-ai'
|
||||
import type { CacheRetention, ThinkingBudgets, ThinkingLevel, Transport } from '@earendil-works/pi-ai'
|
||||
import z from 'schemastery'
|
||||
|
||||
/** Configuration for one pi-ai provider route. */
|
||||
export interface PiAiProviderProfile {
|
||||
/** pi-ai provider catalog name and Harness route key. */
|
||||
provider: string
|
||||
/** Provider credential; when absent pi-ai uses its provider-native ambient discovery. */
|
||||
apiKey?: string
|
||||
/** Override the selected catalog model's endpoint without changing its protocol metadata. */
|
||||
baseURL?: string
|
||||
/** Provider request headers; Harness attribution wins reserved names. */
|
||||
headers?: Record<string, string>
|
||||
/** Provider-neutral pi-ai reasoning level. */
|
||||
reasoning?: ThinkingLevel
|
||||
/** Token budgets used by reasoning providers that support them. */
|
||||
thinkingBudgets?: ThinkingBudgets
|
||||
/** Prompt-cache retention preference. */
|
||||
cacheRetention?: CacheRetention
|
||||
/** Streaming transport preference. */
|
||||
transport?: Transport
|
||||
/** HTTP/provider SDK timeout in milliseconds. */
|
||||
timeoutMs?: number
|
||||
/** WebSocket connection timeout in milliseconds. */
|
||||
websocketConnectTimeoutMs?: number
|
||||
/** Provider SDK retry count. */
|
||||
maxRetries?: number
|
||||
/** Maximum provider-requested retry delay in milliseconds. */
|
||||
maxRetryDelayMs?: number
|
||||
}
|
||||
|
||||
/** Plugin configuration: the non-empty provider profiles this instance owns. */
|
||||
export interface Config {
|
||||
/** Non-empty set of pi-ai provider routes this adapter instance owns. */
|
||||
providers: PiAiProviderProfile[]
|
||||
}
|
||||
|
||||
const thinkingBudgets = z.object({
|
||||
minimal: z.number(),
|
||||
low: z.number(),
|
||||
medium: z.number(),
|
||||
high: z.number(),
|
||||
})
|
||||
|
||||
const profile = z.object({
|
||||
provider: z.string().required(),
|
||||
apiKey: z.string(),
|
||||
baseURL: z.string(),
|
||||
headers: z.dict(z.string()),
|
||||
reasoning: z.union(['minimal', 'low', 'medium', 'high', 'xhigh']),
|
||||
thinkingBudgets,
|
||||
cacheRetention: z.union(['none', 'short', 'long']),
|
||||
transport: z.union(['sse', 'websocket', 'websocket-cached', 'auto']),
|
||||
timeoutMs: z.number(),
|
||||
websocketConnectTimeoutMs: z.number(),
|
||||
maxRetries: z.number(),
|
||||
maxRetryDelayMs: z.number(),
|
||||
})
|
||||
|
||||
/** Runtime schema for {@link Config}. */
|
||||
export const Config: z<Config> = z.object({
|
||||
providers: z.array(profile).required(),
|
||||
})
|
||||
|
||||
/**
|
||||
* Validate profiles against the installed pi-ai catalog and return a detached
|
||||
* shallow copy suitable for adapter construction.
|
||||
* @param profiles - configured provider profiles.
|
||||
* @returns validated profiles in configuration order.
|
||||
*/
|
||||
export function resolveProfiles(profiles: readonly PiAiProviderProfile[]): PiAiProviderProfile[] {
|
||||
if (profiles.length === 0) throw new Error('llm-pi-ai: providers must contain at least one profile')
|
||||
const supported = new Set<string>(getProviders())
|
||||
const seen = new Set<string>()
|
||||
return profiles.map((source) => {
|
||||
if (source.provider.length === 0) throw new Error('llm-pi-ai: provider names must be non-empty')
|
||||
if (!supported.has(source.provider)) throw new Error(`llm-pi-ai: unknown pi-ai provider "${source.provider}"`)
|
||||
if (seen.has(source.provider)) throw new Error(`llm-pi-ai: duplicate provider profile "${source.provider}"`)
|
||||
if (source.apiKey !== undefined && source.apiKey.length === 0) {
|
||||
throw new Error(`llm-pi-ai: provider "${source.provider}" has an empty apiKey; omit it to use ambient authentication`)
|
||||
}
|
||||
if (source.baseURL !== undefined && source.baseURL.length === 0) {
|
||||
throw new Error(`llm-pi-ai: provider "${source.provider}" has an empty baseURL`)
|
||||
}
|
||||
seen.add(source.provider)
|
||||
return {
|
||||
...source,
|
||||
...source.headers === undefined ? {} : { headers: { ...source.headers } },
|
||||
...source.thinkingBudgets === undefined ? {} : { thinkingBudgets: { ...source.thinkingBudgets } },
|
||||
}
|
||||
})
|
||||
}
|
||||
85
packages/llm/llm-pi-ai/src/context.ts
Normal file
85
packages/llm/llm-pi-ai/src/context.ts
Normal file
@@ -0,0 +1,85 @@
|
||||
/**
|
||||
* Harness request-history conversion into pi-ai's Context vocabulary.
|
||||
*
|
||||
* @module dsh-llm-pi-ai/context
|
||||
*/
|
||||
|
||||
import { CallId } from '@deepseek-ai/dsh-llm'
|
||||
import type { GenerateOptions, Message } from '@deepseek-ai/dsh-llm'
|
||||
import type { Context as PiContext, Message as PiMessage, Tool as PiTool } from '@earendil-works/pi-ai'
|
||||
import { toPiAssistant } from './replay.ts'
|
||||
|
||||
/** Join the text blocks of a harness message. */
|
||||
function flattenText(message: Message): string {
|
||||
return message.content
|
||||
.filter(block => block.type === 'text')
|
||||
.map(block => block.text)
|
||||
.join('')
|
||||
}
|
||||
|
||||
/**
|
||||
* Convert harness history to a pi-ai Context. Tool results need the tool
|
||||
* NAME (pi-ai's `toolName`), which the harness doesn't carry on the result
|
||||
* block — it is recovered from the preceding assistant tool-call with the
|
||||
* same id.
|
||||
* @param options - the harness request; `options.system` maps to pi-ai's single `systemPrompt` slot.
|
||||
* @returns the pi-ai context; `tools` is omitted entirely when the request declares none.
|
||||
*/
|
||||
export function toPiContext(options: GenerateOptions): PiContext {
|
||||
const toolNames = new Map<CallId, string>()
|
||||
const messages: PiMessage[] = []
|
||||
|
||||
for (const message of options.messages) {
|
||||
if (message.role === 'system') {
|
||||
// pi-ai has a single systemPrompt slot; in-history system messages are
|
||||
// folded into user messages to preserve order (rare in practice — the
|
||||
// harness sends the system prompt via options.system).
|
||||
messages.push({ role: 'user', content: flattenText(message), timestamp: 0 })
|
||||
continue
|
||||
}
|
||||
if (message.role === 'assistant') {
|
||||
const assistant = toPiAssistant(message)
|
||||
for (const block of assistant.content) {
|
||||
if (block.type === 'toolCall') toolNames.set(CallId(block.id), block.name)
|
||||
}
|
||||
messages.push(assistant)
|
||||
continue
|
||||
}
|
||||
// user role: text + tool results (each result becomes its own message).
|
||||
const text = flattenText(message)
|
||||
const results = message.content.filter(block => block.type === 'tool-result')
|
||||
if (text.length > 0 || results.length === 0) {
|
||||
messages.push({ role: 'user', content: text, timestamp: 0 })
|
||||
}
|
||||
for (const result of results) {
|
||||
messages.push({
|
||||
role: 'toolResult',
|
||||
toolCallId: result.toolCallId,
|
||||
toolName: toolNames.get(result.toolCallId) ?? 'unknown',
|
||||
content: [{
|
||||
type: 'text',
|
||||
text: result.content
|
||||
.filter(block => block.type === 'text')
|
||||
.map(block => block.text)
|
||||
.join('') || '(no output)',
|
||||
}],
|
||||
isError: result.isError ?? false,
|
||||
timestamp: 0,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
const tools: PiTool[] | undefined = options.tools?.map(tool => ({
|
||||
name: tool.name,
|
||||
description: tool.description,
|
||||
// ToolSchema.parameters is a JSON Schema object; pi-ai's TSchema
|
||||
// (TypeBox) is structurally JSON Schema, so it assigns directly.
|
||||
parameters: tool.parameters,
|
||||
}))
|
||||
|
||||
return {
|
||||
...options.system !== undefined ? { systemPrompt: options.system } : {},
|
||||
messages,
|
||||
...tools !== undefined && tools.length > 0 ? { tools } : {},
|
||||
}
|
||||
}
|
||||
@@ -1,289 +0,0 @@
|
||||
/**
|
||||
* Bidirectional mapping between the harness vocabulary and pi-ai's:
|
||||
* `GenerateOptions`/`Message[]` → pi-ai `Context`, and pi-ai
|
||||
* `AssistantMessageEvent`s → harness `StreamChunk`s.
|
||||
*
|
||||
* Vocabulary differences worth knowing (they are exactly why this adapter
|
||||
* exists — an independent implementation stress-tests the StreamChunk
|
||||
* protocol):
|
||||
* - pi-ai tool-call `arguments` are PARSED OBJECTS; the harness keeps the
|
||||
* raw JSON string. We parse on the way into pi-ai, patch provider payloads
|
||||
* back to the original raw string in the adapter, and re-stringify on output.
|
||||
* - pi-ai reports errors as in-stream `error` events (it never throws
|
||||
* mid-stream); the harness expresses those as `finish {kind:'error'}` /
|
||||
* `{kind:'aborted'}` chunks.
|
||||
* - pi-ai folds reasoning tokens into `usage.output`; there is no separate
|
||||
* reasoning count to map.
|
||||
*
|
||||
* @module dsh-llm-pi-ai/convert
|
||||
*/
|
||||
|
||||
import { CallId, LlmError } from '@deepseek-ai/dsh-llm'
|
||||
import type { FinishReason, GenerateOptions, Message, StreamChunk, TokenUsage } from '@deepseek-ai/dsh-llm'
|
||||
import type {
|
||||
AssistantMessage,
|
||||
AssistantMessageEvent,
|
||||
Context as PiContext,
|
||||
Message as PiMessage,
|
||||
Tool as PiTool,
|
||||
Usage as PiUsage,
|
||||
} from '@earendil-works/pi-ai'
|
||||
|
||||
/** Join the text blocks of a harness message. */
|
||||
function flattenText(message: Message): string {
|
||||
return message.content
|
||||
.filter(block => block.type === 'text')
|
||||
.map(block => block.text)
|
||||
.join('')
|
||||
}
|
||||
|
||||
/** Parse tool-call argument JSON; tolerate model malformations with {}. */
|
||||
function parseArguments(raw: string): Record<string, unknown> {
|
||||
try {
|
||||
const parsed: unknown = JSON.parse(raw)
|
||||
if (typeof parsed === 'object' && parsed !== null && !Array.isArray(parsed)) {
|
||||
return parsed as Record<string, unknown>
|
||||
}
|
||||
} catch {
|
||||
// fall through
|
||||
}
|
||||
return {}
|
||||
}
|
||||
|
||||
/**
|
||||
* Convert harness history to a pi-ai Context. Tool results need the tool
|
||||
* NAME (pi-ai's `toolName`), which the harness doesn't carry on the result
|
||||
* block — it is recovered from the preceding assistant tool-call with the
|
||||
* same id.
|
||||
* @param options - the harness request; `options.system` maps to pi-ai's single `systemPrompt` slot.
|
||||
* @returns the pi-ai context; `tools` is omitted entirely when the request declares none.
|
||||
*/
|
||||
export function toPiContext(options: GenerateOptions): PiContext {
|
||||
const toolNames = new Map<CallId, string>()
|
||||
const messages: PiMessage[] = []
|
||||
|
||||
for (const message of options.messages) {
|
||||
if (message.role === 'system') {
|
||||
// pi-ai has a single systemPrompt slot; in-history system messages are
|
||||
// folded into user messages to preserve order (rare in practice — the
|
||||
// harness sends the system prompt via options.system).
|
||||
messages.push({ role: 'user', content: flattenText(message), timestamp: 0 })
|
||||
continue
|
||||
}
|
||||
if (message.role === 'assistant') {
|
||||
const content: AssistantMessage['content'] = []
|
||||
for (const block of message.content) {
|
||||
switch (block.type) {
|
||||
case 'text':
|
||||
content.push({ type: 'text', text: block.text })
|
||||
break
|
||||
case 'reasoning':
|
||||
// thinkingSignature names the wire field pi-ai replays the CoT
|
||||
// under. Without it pi-ai falls back to reasoning_content: ""
|
||||
// (its requiresReasoningContentOnAssistantMessages shim), which
|
||||
// violates DeepSeek's thinking-mode passback rule on tool-call
|
||||
// turns (guides/thinking_mode.mdx § Tool Calls).
|
||||
content.push({ type: 'thinking', thinking: block.text, thinkingSignature: 'reasoning_content' })
|
||||
break
|
||||
case 'tool-call':
|
||||
toolNames.set(block.id, block.name)
|
||||
content.push({
|
||||
type: 'toolCall',
|
||||
id: block.id,
|
||||
name: block.name,
|
||||
arguments: parseArguments(block.arguments),
|
||||
})
|
||||
break
|
||||
default:
|
||||
// plugin-added block types: not representable here.
|
||||
break
|
||||
}
|
||||
}
|
||||
messages.push({
|
||||
role: 'assistant',
|
||||
content,
|
||||
api: 'openai-completions',
|
||||
provider: 'deepseek',
|
||||
model: options.model,
|
||||
usage: emptyPiUsage(),
|
||||
stopReason: content.some(piece => piece.type === 'toolCall') ? 'toolUse' : 'stop',
|
||||
timestamp: 0,
|
||||
})
|
||||
continue
|
||||
}
|
||||
// user role: text + tool results (each result becomes its own message).
|
||||
const text = flattenText(message)
|
||||
const results = message.content.filter(block => block.type === 'tool-result')
|
||||
if (text.length > 0 || results.length === 0) {
|
||||
messages.push({ role: 'user', content: text, timestamp: 0 })
|
||||
}
|
||||
for (const result of results) {
|
||||
messages.push({
|
||||
role: 'toolResult',
|
||||
toolCallId: result.toolCallId,
|
||||
toolName: toolNames.get(result.toolCallId) ?? 'unknown',
|
||||
content: [{
|
||||
type: 'text',
|
||||
text: result.content
|
||||
.filter(block => block.type === 'text')
|
||||
.map(block => block.text)
|
||||
.join('') || '(no output)',
|
||||
}],
|
||||
isError: result.isError ?? false,
|
||||
timestamp: 0,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
const tools: PiTool[] | undefined = options.tools?.map(tool => ({
|
||||
name: tool.name,
|
||||
description: tool.description,
|
||||
// ToolSchema.parameters is a JSON Schema object; pi-ai's TSchema
|
||||
// (TypeBox) is structurally JSON Schema, so it assigns directly.
|
||||
parameters: tool.parameters,
|
||||
}))
|
||||
|
||||
return {
|
||||
...options.system !== undefined ? { systemPrompt: options.system } : {},
|
||||
messages,
|
||||
...tools !== undefined && tools.length > 0 ? { tools } : {},
|
||||
}
|
||||
}
|
||||
|
||||
function emptyPiUsage(): PiUsage {
|
||||
return {
|
||||
input: 0,
|
||||
output: 0,
|
||||
cacheRead: 0,
|
||||
cacheWrite: 0,
|
||||
totalTokens: 0,
|
||||
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Map pi-ai usage (reasoning folded into output by pi-ai).
|
||||
* @param usage - cumulative usage from the terminal pi-ai event.
|
||||
* @returns harness counts; cache fields appear only when non-zero (pi-ai reports zeros, not absence).
|
||||
*/
|
||||
export function mapUsage(usage: PiUsage): TokenUsage {
|
||||
return {
|
||||
inputTokens: usage.input,
|
||||
outputTokens: usage.output,
|
||||
...usage.cacheRead > 0 ? { cacheReadTokens: usage.cacheRead } : {},
|
||||
...usage.cacheWrite > 0 ? { cacheWriteTokens: usage.cacheWrite } : {},
|
||||
}
|
||||
}
|
||||
|
||||
function classifyPiAiError(message: string): string {
|
||||
if (/\b(?:401|403)\b/.test(message)) return 'AUTH'
|
||||
if (/\b429\b|rate.?limit/i.test(message)) return 'RATE_LIMIT'
|
||||
if (/\b400\b|invalid.?request/i.test(message)) return 'INVALID_REQUEST'
|
||||
if (/\b5\d\d\b/.test(message)) return 'SERVER'
|
||||
return 'PI_AI_ERROR'
|
||||
}
|
||||
|
||||
/**
|
||||
* Map a terminal pi-ai event to the harness finish reason.
|
||||
* @param message - the assistant message carried by the `done` or `error` event.
|
||||
* @returns the harness reason; `error` yields `{kind: 'error'}` with a code classified from the error text.
|
||||
*/
|
||||
export function mapStopReason(message: AssistantMessage): FinishReason {
|
||||
switch (message.stopReason) {
|
||||
case 'stop': return { kind: 'stop' }
|
||||
case 'length': return { kind: 'max-tokens' }
|
||||
case 'toolUse': return { kind: 'tool-calls' }
|
||||
case 'aborted': return { kind: 'aborted' }
|
||||
case 'error': {
|
||||
const text = message.errorMessage ?? 'pi-ai stream error'
|
||||
return { kind: 'error', message: text, code: classifyPiAiError(text) }
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Translate the pi-ai event stream into StreamChunks. pi-ai never throws
|
||||
* mid-stream — failures arrive as `error` events, which become error/aborted
|
||||
* `finish` chunks (the harness protocol's other error-delivery style).
|
||||
* @param events - one assistant turn's pi-ai event stream.
|
||||
* @returns the harness chunks, ending with `usage` then `finish`; throws
|
||||
* `LlmError` (`STREAM_CLOSED`) if the source ends without a terminal event.
|
||||
*/
|
||||
export async function* toStreamChunks(events: AsyncIterable<AssistantMessageEvent>): AsyncGenerator<StreamChunk> {
|
||||
// pi-ai contentIndex ↔ our block index map 1:1 (both count blocks from 0
|
||||
// in stream order), but we track ids per index for tool calls.
|
||||
const toolIds = new Map<number, { id: string; name: string }>()
|
||||
|
||||
for await (const event of events) {
|
||||
switch (event.type) {
|
||||
case 'start':
|
||||
break
|
||||
case 'text_start':
|
||||
yield { type: 'block-start', index: event.contentIndex, blockType: 'text' }
|
||||
break
|
||||
case 'text_delta':
|
||||
yield { type: 'text-delta', index: event.contentIndex, text: event.delta }
|
||||
break
|
||||
case 'text_end':
|
||||
yield { type: 'block-end', index: event.contentIndex, block: { type: 'text', text: event.content } }
|
||||
break
|
||||
case 'thinking_start':
|
||||
yield { type: 'block-start', index: event.contentIndex, blockType: 'reasoning' }
|
||||
break
|
||||
case 'thinking_delta':
|
||||
yield { type: 'reasoning-delta', index: event.contentIndex, text: event.delta }
|
||||
break
|
||||
case 'thinking_end':
|
||||
yield { type: 'block-end', index: event.contentIndex, block: { type: 'reasoning', text: event.content } }
|
||||
break
|
||||
case 'toolcall_start': {
|
||||
// The id/name live on the partial's content at this index.
|
||||
const partial = event.partial.content[event.contentIndex]
|
||||
const id = partial?.type === 'toolCall' ? partial.id : ''
|
||||
const name = partial?.type === 'toolCall' ? partial.name : ''
|
||||
toolIds.set(event.contentIndex, { id, name })
|
||||
yield { type: 'block-start', index: event.contentIndex, blockType: 'tool-call' }
|
||||
break
|
||||
}
|
||||
case 'toolcall_delta': {
|
||||
const known = toolIds.get(event.contentIndex)
|
||||
yield {
|
||||
type: 'tool-call-delta',
|
||||
index: event.contentIndex,
|
||||
id: CallId(known?.id ?? ''),
|
||||
...known?.name !== undefined && known.name.length > 0 ? { name: known.name } : {},
|
||||
argumentsDelta: event.delta,
|
||||
}
|
||||
break
|
||||
}
|
||||
case 'toolcall_end':
|
||||
yield {
|
||||
type: 'block-end',
|
||||
index: event.contentIndex,
|
||||
block: {
|
||||
type: 'tool-call',
|
||||
id: CallId(event.toolCall.id),
|
||||
name: event.toolCall.name,
|
||||
// pi-ai hands back the PARSED arguments; the harness vocabulary
|
||||
// keeps the raw string.
|
||||
arguments: JSON.stringify(event.toolCall.arguments),
|
||||
},
|
||||
}
|
||||
break
|
||||
case 'done':
|
||||
yield { type: 'usage', usage: mapUsage(event.message.usage) }
|
||||
yield { type: 'finish', reason: mapStopReason(event.message) }
|
||||
return
|
||||
case 'error':
|
||||
// In-stream error delivery (pi-ai's style) → error finish chunk
|
||||
// (the harness's other sanctioned error path besides throwing).
|
||||
yield { type: 'usage', usage: mapUsage(event.error.usage) }
|
||||
yield { type: 'finish', reason: mapStopReason(event.error) }
|
||||
return
|
||||
// no default: AssistantMessageEvent is pi-ai's closed union; a new
|
||||
// event type should fail compilation here via tsc's exhaustiveness
|
||||
// when one is added (switch covers all current variants).
|
||||
}
|
||||
}
|
||||
throw new LlmError('pi-ai event stream ended without done/error', 'STREAM_CLOSED')
|
||||
}
|
||||
@@ -1,76 +1,45 @@
|
||||
/**
|
||||
* pi-ai-backed DeepSeek adapter plugin. Same Config shape as
|
||||
* `@deepseek-ai/dsh-llm-deepseek` (one-line swap in cordis.yml), different
|
||||
* implementation underneath — see `./adapter.ts` for why both exist.
|
||||
* Generic pi-ai-backed LLM adapter plugin. One plugin instance registers an
|
||||
* explicit set of provider profiles; requests select a profile by provider and
|
||||
* resolve the model dynamically from pi-ai's installed catalog.
|
||||
*
|
||||
* ```yaml
|
||||
* - id: llm
|
||||
* name: '@deepseek-ai/dsh-llm-pi-ai'
|
||||
* config:
|
||||
* apiKey: !!js process.env.DEEPSEEK_API_KEY
|
||||
* baseURL: !!js process.env.DEEPSEEK_BASE_URL
|
||||
* models: [deepseek-v4-flash, deepseek-v4-pro]
|
||||
* reasoning: high
|
||||
* providers:
|
||||
* - provider: openai
|
||||
* apiKey: !!js process.env.OPENAI_API_KEY
|
||||
* - provider: anthropic
|
||||
* apiKey: !!js process.env.ANTHROPIC_API_KEY
|
||||
* - provider: openrouter
|
||||
* apiKey: !!js process.env.OPENROUTER_API_KEY
|
||||
* baseURL: https://proxy.example.com/v1
|
||||
* ```
|
||||
*
|
||||
* @module @deepseek-ai/dsh-llm-pi-ai
|
||||
*/
|
||||
|
||||
import type { Context } from 'cordis'
|
||||
import z from 'schemastery'
|
||||
import type {} from '@deepseek-ai/dsh-llm'
|
||||
import { PiAiAdapter } from './adapter.ts'
|
||||
import type { PiAiReasoning } from './adapter.ts'
|
||||
import { Config, resolveProfiles } from './config.ts'
|
||||
|
||||
export { buildModel, PiAiAdapter } from './adapter.ts'
|
||||
export type { PiAiAdapterOptions, PiAiReasoning } from './adapter.ts'
|
||||
export { mapStopReason, mapUsage, toPiContext, toStreamChunks } from './convert.ts'
|
||||
export { PiAiAdapter } from './adapter.ts'
|
||||
export type { PiAiAdapterOptions } from './adapter.ts'
|
||||
export { Config, resolveProfiles } from './config.ts'
|
||||
export type { PiAiProviderProfile } from './config.ts'
|
||||
export { toPiContext } from './context.ts'
|
||||
export { toPiReplayState } from './replay.ts'
|
||||
export type { PiAiReplayState } from './replay.ts'
|
||||
export { mapStopReason, mapUsage, toStreamChunks } from './stream.ts'
|
||||
|
||||
export const name = 'llm-pi-ai'
|
||||
export const inject = ['llm']
|
||||
|
||||
/**
|
||||
* Plugin config, validated by the same-named schemastery schema. Every field
|
||||
* is optional in yml: credentials/endpoint fall back to the environment (a
|
||||
* missing API key fails plugin load, not the first call).
|
||||
*/
|
||||
export interface Config {
|
||||
/** API key; falls back to $DEEPSEEK_API_KEY. Required one way or the other. */
|
||||
apiKey?: string
|
||||
/** Endpoint base; falls back to $DEEPSEEK_BASE_URL, then the public API. */
|
||||
baseURL?: string
|
||||
/** Model names to register (sent verbatim on the wire). */
|
||||
models?: string[]
|
||||
/**
|
||||
* Thinking level for every request: 'off' disables thinking mode; 'high'
|
||||
* and 'xhigh' (wire 'max') set the effort. Omitted = provider default
|
||||
* (thinking enabled), matching llm-deepseek's omission semantics.
|
||||
*/
|
||||
reasoning?: PiAiReasoning
|
||||
}
|
||||
|
||||
export const Config: z<Config> = z.object({
|
||||
apiKey: z.string(),
|
||||
baseURL: z.string(),
|
||||
models: z.array(z.string()).default(['deepseek-v4-flash', 'deepseek-v4-pro']),
|
||||
reasoning: z.union(['off', 'high', 'xhigh']),
|
||||
})
|
||||
|
||||
/** Public API default; the internal endpoint comes from $DEEPSEEK_BASE_URL. */
|
||||
export const PUBLIC_BASE_URL = 'https://api.deepseek.com'
|
||||
|
||||
/** Register one generic pi-ai adapter for all configured provider routes. */
|
||||
export function apply(ctx: Context, config: Config): void {
|
||||
const apiKey = config.apiKey ?? process.env.DEEPSEEK_API_KEY
|
||||
if (apiKey === undefined || apiKey.length === 0) {
|
||||
throw new Error('llm-pi-ai: an API key is required (Config.apiKey or $DEEPSEEK_API_KEY)')
|
||||
}
|
||||
const baseURL = config.baseURL ?? process.env.DEEPSEEK_BASE_URL ?? PUBLIC_BASE_URL
|
||||
// schemastery's .default() guarantees models is set after validation.
|
||||
const models = config.models as string[]
|
||||
|
||||
ctx.llm.registerAdapter(models, new PiAiAdapter({
|
||||
apiKey,
|
||||
baseURL,
|
||||
reasoning: config.reasoning,
|
||||
}))
|
||||
const profiles = resolveProfiles(config.providers)
|
||||
const adapter = new PiAiAdapter({ profiles })
|
||||
ctx.llm.registerAdapter(profiles.map(entry => entry.provider), adapter)
|
||||
}
|
||||
|
||||
208
packages/llm/llm-pi-ai/src/replay.ts
Normal file
208
packages/llm/llm-pi-ai/src/replay.ts
Normal file
@@ -0,0 +1,208 @@
|
||||
/**
|
||||
* Durable pi-ai replay metadata and assistant-history reconstruction.
|
||||
*
|
||||
* Harness content remains the durable source for text and tool calls. This
|
||||
* module stores only the provider-native metadata needed to reconstruct a
|
||||
* pi-ai assistant message on a later request.
|
||||
*
|
||||
* @module dsh-llm-pi-ai/replay
|
||||
*/
|
||||
|
||||
import { LlmError } from '@deepseek-ai/dsh-llm'
|
||||
import type { Message } from '@deepseek-ai/dsh-llm'
|
||||
import type { Api, AssistantMessage, Usage as PiUsage } from '@earendil-works/pi-ai'
|
||||
|
||||
type PiAiReplayBlock =
|
||||
| { type: 'text'; textSignature?: string }
|
||||
| { type: 'reasoning'; thinkingSignature?: string; redacted?: boolean }
|
||||
| { type: 'tool-call'; thoughtSignature?: string }
|
||||
|
||||
/** Versioned adapter-private projection required to replay a pi-ai response. */
|
||||
export interface PiAiReplayState {
|
||||
kind: 'pi-ai'
|
||||
version: 1
|
||||
api: Api
|
||||
provider: string
|
||||
model: string
|
||||
responseModel?: string
|
||||
responseId?: string
|
||||
stopReason: AssistantMessage['stopReason']
|
||||
blocks: PiAiReplayBlock[]
|
||||
}
|
||||
|
||||
/** Parse tool-call argument JSON; tolerate model malformations with {}. */
|
||||
function parseArguments(raw: string): Record<string, unknown> {
|
||||
try {
|
||||
const parsed: unknown = JSON.parse(raw)
|
||||
if (typeof parsed === 'object' && parsed !== null && !Array.isArray(parsed)) {
|
||||
return parsed as Record<string, unknown>
|
||||
}
|
||||
} catch {
|
||||
// fall through
|
||||
}
|
||||
return {}
|
||||
}
|
||||
|
||||
/** Construct the zero usage value required by historical pi-ai messages. */
|
||||
function emptyPiUsage(): PiUsage {
|
||||
return {
|
||||
input: 0,
|
||||
output: 0,
|
||||
cacheRead: 0,
|
||||
cacheWrite: 0,
|
||||
totalTokens: 0,
|
||||
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Project a successful pi-ai response into the minimal durable replay state.
|
||||
* @param message - completed native pi-ai assistant response.
|
||||
* @returns the versioned lossless-JSON replay projection.
|
||||
*/
|
||||
export function toPiReplayState(message: AssistantMessage): PiAiReplayState {
|
||||
return {
|
||||
kind: 'pi-ai',
|
||||
version: 1,
|
||||
api: message.api,
|
||||
provider: message.provider,
|
||||
model: message.model,
|
||||
...message.responseModel === undefined ? {} : { responseModel: message.responseModel },
|
||||
...message.responseId === undefined ? {} : { responseId: message.responseId },
|
||||
stopReason: message.stopReason,
|
||||
blocks: message.content.map((block): PiAiReplayBlock => {
|
||||
switch (block.type) {
|
||||
case 'text': return {
|
||||
type: 'text',
|
||||
...block.textSignature === undefined ? {} : { textSignature: block.textSignature },
|
||||
}
|
||||
case 'thinking': return {
|
||||
type: 'reasoning',
|
||||
...block.thinkingSignature === undefined ? {} : { thinkingSignature: block.thinkingSignature },
|
||||
...block.redacted === undefined ? {} : { redacted: block.redacted },
|
||||
}
|
||||
case 'toolCall': return {
|
||||
type: 'tool-call',
|
||||
...block.thoughtSignature === undefined ? {} : { thoughtSignature: block.thoughtSignature },
|
||||
}
|
||||
}
|
||||
}),
|
||||
}
|
||||
}
|
||||
|
||||
function invalidReplay(message: string): never {
|
||||
throw new LlmError(`invalid pi-ai replay state: ${message}`, 'INVALID_REPLAY_STATE')
|
||||
}
|
||||
|
||||
/** Validate the adapter-private state before it reaches pi-ai. */
|
||||
function readReplayState(value: unknown): PiAiReplayState {
|
||||
if (typeof value !== 'object' || value === null || Array.isArray(value)) return invalidReplay('expected an object')
|
||||
const state = value as Record<string, unknown>
|
||||
if (state['kind'] !== 'pi-ai') return invalidReplay('unknown state kind')
|
||||
if (state['version'] !== 1) return invalidReplay(`unsupported version ${String(state['version'])}`)
|
||||
for (const key of ['api', 'provider', 'model'] as const) {
|
||||
if (typeof state[key] !== 'string' || state[key].length === 0) return invalidReplay(`${key} must be a non-empty string`)
|
||||
}
|
||||
if (!['stop', 'length', 'toolUse', 'error', 'aborted'].includes(String(state['stopReason']))) {
|
||||
return invalidReplay('unknown stopReason')
|
||||
}
|
||||
if (state['responseModel'] !== undefined && typeof state['responseModel'] !== 'string') return invalidReplay('responseModel must be a string')
|
||||
if (state['responseId'] !== undefined && typeof state['responseId'] !== 'string') return invalidReplay('responseId must be a string')
|
||||
if (!Array.isArray(state['blocks'])) return invalidReplay('blocks must be an array')
|
||||
for (const [index, value] of state['blocks'].entries()) {
|
||||
if (typeof value !== 'object' || value === null || Array.isArray(value)) return invalidReplay(`block ${index} must be an object`)
|
||||
const block = value as Record<string, unknown>
|
||||
if (!['text', 'reasoning', 'tool-call'].includes(String(block['type']))) return invalidReplay(`block ${index} has an unknown type`)
|
||||
for (const signature of ['textSignature', 'thinkingSignature', 'thoughtSignature'] as const) {
|
||||
if (block[signature] !== undefined && typeof block[signature] !== 'string') return invalidReplay(`block ${index} ${signature} must be a string`)
|
||||
}
|
||||
if (block['redacted'] !== undefined && typeof block['redacted'] !== 'boolean') return invalidReplay(`block ${index} redacted must be boolean`)
|
||||
}
|
||||
return state as unknown as PiAiReplayState
|
||||
}
|
||||
|
||||
/** Convert provider-neutral blocks without trusting them as same-model replay. */
|
||||
function foreignAssistant(message: Message): AssistantMessage {
|
||||
const content: AssistantMessage['content'] = []
|
||||
for (const block of message.content) {
|
||||
switch (block.type) {
|
||||
case 'text': content.push({ type: 'text', text: block.text }); break
|
||||
case 'reasoning': content.push({ type: 'thinking', thinking: block.text }); break
|
||||
case 'tool-call': content.push({
|
||||
type: 'toolCall',
|
||||
id: block.id,
|
||||
name: block.name,
|
||||
arguments: parseArguments(block.arguments),
|
||||
}); break
|
||||
default:
|
||||
// plugin-added block types are not representable in pi-ai.
|
||||
break
|
||||
}
|
||||
}
|
||||
return {
|
||||
role: 'assistant',
|
||||
content,
|
||||
// Deliberately never equals a catalog API: absent replay state is foreign
|
||||
// even if provenance names the same provider/model as this request.
|
||||
api: 'dsh-foreign',
|
||||
provider: message.provenance?.provider ?? 'dsh-foreign',
|
||||
model: message.provenance?.model ?? 'dsh-foreign',
|
||||
usage: emptyPiUsage(),
|
||||
stopReason: content.some(piece => piece.type === 'toolCall') ? 'toolUse' : 'stop',
|
||||
timestamp: 0,
|
||||
}
|
||||
}
|
||||
|
||||
/** Recombine durable Harness content with validated pi-ai replay metadata. */
|
||||
function replayedAssistant(message: Message, rawState: unknown): AssistantMessage {
|
||||
const state = readReplayState(rawState)
|
||||
if (state.blocks.length !== message.content.length) return invalidReplay('block count does not match assistant content')
|
||||
const content: AssistantMessage['content'] = message.content.map((block, index) => {
|
||||
const replay = state.blocks[index]
|
||||
if (replay === undefined || replay.type !== block.type) return invalidReplay(`block ${index} does not match assistant content`)
|
||||
switch (block.type) {
|
||||
case 'text': return {
|
||||
type: 'text',
|
||||
text: block.text,
|
||||
...replay.type === 'text' && replay.textSignature !== undefined ? { textSignature: replay.textSignature } : {},
|
||||
}
|
||||
case 'reasoning': return {
|
||||
type: 'thinking',
|
||||
thinking: block.text,
|
||||
...replay.type === 'reasoning' && replay.thinkingSignature !== undefined ? { thinkingSignature: replay.thinkingSignature } : {},
|
||||
...replay.type === 'reasoning' && replay.redacted !== undefined ? { redacted: replay.redacted } : {},
|
||||
}
|
||||
case 'tool-call': return {
|
||||
type: 'toolCall',
|
||||
id: block.id,
|
||||
name: block.name,
|
||||
arguments: parseArguments(block.arguments),
|
||||
...replay.type === 'tool-call' && replay.thoughtSignature !== undefined ? { thoughtSignature: replay.thoughtSignature } : {},
|
||||
}
|
||||
/* v8 ignore next -- readReplayState rejects unknown replay tags, so an equal plugin-added Harness tag cannot reach this switch */
|
||||
default: return invalidReplay(`block ${index} has an unsupported Harness type`)
|
||||
}
|
||||
})
|
||||
return {
|
||||
role: 'assistant',
|
||||
content,
|
||||
api: state.api,
|
||||
provider: state.provider,
|
||||
model: state.model,
|
||||
...state.responseModel === undefined ? {} : { responseModel: state.responseModel },
|
||||
...state.responseId === undefined ? {} : { responseId: state.responseId },
|
||||
usage: emptyPiUsage(),
|
||||
stopReason: state.stopReason,
|
||||
timestamp: 0,
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Convert one durable Harness assistant message into pi-ai history.
|
||||
* @param message - assistant content with optional adapter-owned replay metadata.
|
||||
* @returns a native pi-ai assistant message reconstructed from durable content.
|
||||
*/
|
||||
export function toPiAssistant(message: Message): AssistantMessage {
|
||||
const replayState = message.provenance?.replayState
|
||||
return replayState === undefined ? foreignAssistant(message) : replayedAssistant(message, replayState)
|
||||
}
|
||||
141
packages/llm/llm-pi-ai/src/stream.ts
Normal file
141
packages/llm/llm-pi-ai/src/stream.ts
Normal file
@@ -0,0 +1,141 @@
|
||||
/**
|
||||
* pi-ai assistant event translation into the Harness streaming protocol.
|
||||
*
|
||||
* pi-ai tool-call arguments are parsed objects while the Harness keeps their
|
||||
* raw JSON representation. pi-ai also reports failures as terminal stream
|
||||
* events, which this module maps into Harness finish chunks.
|
||||
*
|
||||
* @module dsh-llm-pi-ai/stream
|
||||
*/
|
||||
|
||||
import { CallId, LlmError } from '@deepseek-ai/dsh-llm'
|
||||
import type { FinishReason, StreamChunk, TokenUsage } from '@deepseek-ai/dsh-llm'
|
||||
import type { AssistantMessage, AssistantMessageEvent, Usage as PiUsage } from '@earendil-works/pi-ai'
|
||||
import { toPiReplayState } from './replay.ts'
|
||||
|
||||
/**
|
||||
* Map pi-ai usage (reasoning folded into output by pi-ai).
|
||||
* @param usage - cumulative usage from the terminal pi-ai event.
|
||||
* @returns harness counts; cache fields appear only when non-zero (pi-ai reports zeros, not absence).
|
||||
*/
|
||||
export function mapUsage(usage: PiUsage): TokenUsage {
|
||||
return {
|
||||
inputTokens: usage.input,
|
||||
outputTokens: usage.output,
|
||||
...usage.cacheRead > 0 ? { cacheReadTokens: usage.cacheRead } : {},
|
||||
...usage.cacheWrite > 0 ? { cacheWriteTokens: usage.cacheWrite } : {},
|
||||
}
|
||||
}
|
||||
|
||||
function classifyPiAiError(message: string): string {
|
||||
if (/\b(?:401|403)\b/.test(message)) return 'AUTH'
|
||||
if (/\b429\b|rate.?limit/i.test(message)) return 'RATE_LIMIT'
|
||||
if (/\b400\b|invalid.?request/i.test(message)) return 'INVALID_REQUEST'
|
||||
if (/\b5\d\d\b/.test(message)) return 'SERVER'
|
||||
return 'PI_AI_ERROR'
|
||||
}
|
||||
|
||||
/**
|
||||
* Map a terminal pi-ai event to the harness finish reason.
|
||||
* @param message - the assistant message carried by the `done` or `error` event.
|
||||
* @returns the harness reason; `error` yields `{kind: 'error'}` with a code classified from the error text.
|
||||
*/
|
||||
export function mapStopReason(message: AssistantMessage): FinishReason {
|
||||
switch (message.stopReason) {
|
||||
case 'stop': return { kind: 'stop' }
|
||||
case 'length': return { kind: 'max-tokens' }
|
||||
case 'toolUse': return { kind: 'tool-calls' }
|
||||
case 'aborted': return { kind: 'aborted' }
|
||||
case 'error': {
|
||||
const text = message.errorMessage ?? 'pi-ai stream error'
|
||||
return { kind: 'error', message: text, code: classifyPiAiError(text) }
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Translate the pi-ai event stream into StreamChunks. pi-ai never throws
|
||||
* mid-stream — failures arrive as `error` events, which become error/aborted
|
||||
* `finish` chunks (the harness protocol's other error-delivery style).
|
||||
* @param events - one assistant turn's pi-ai event stream.
|
||||
* @returns the harness chunks, ending with `usage` then `finish`; throws
|
||||
* `LlmError` (`STREAM_CLOSED`) if the source ends without a terminal event.
|
||||
*/
|
||||
export async function* toStreamChunks(events: AsyncIterable<AssistantMessageEvent>): AsyncGenerator<StreamChunk> {
|
||||
// pi-ai contentIndex ↔ our block index map 1:1 (both count blocks from 0
|
||||
// in stream order), but we track ids per index for tool calls.
|
||||
const toolIds = new Map<number, { id: string; name: string }>()
|
||||
|
||||
for await (const event of events) {
|
||||
switch (event.type) {
|
||||
case 'start':
|
||||
break
|
||||
case 'text_start':
|
||||
yield { type: 'block-start', index: event.contentIndex, blockType: 'text' }
|
||||
break
|
||||
case 'text_delta':
|
||||
yield { type: 'text-delta', index: event.contentIndex, text: event.delta }
|
||||
break
|
||||
case 'text_end':
|
||||
yield { type: 'block-end', index: event.contentIndex, block: { type: 'text', text: event.content } }
|
||||
break
|
||||
case 'thinking_start':
|
||||
yield { type: 'block-start', index: event.contentIndex, blockType: 'reasoning' }
|
||||
break
|
||||
case 'thinking_delta':
|
||||
yield { type: 'reasoning-delta', index: event.contentIndex, text: event.delta }
|
||||
break
|
||||
case 'thinking_end':
|
||||
yield { type: 'block-end', index: event.contentIndex, block: { type: 'reasoning', text: event.content } }
|
||||
break
|
||||
case 'toolcall_start': {
|
||||
// The id/name live on the partial's content at this index.
|
||||
const partial = event.partial.content[event.contentIndex]
|
||||
const id = partial?.type === 'toolCall' ? partial.id : ''
|
||||
const name = partial?.type === 'toolCall' ? partial.name : ''
|
||||
toolIds.set(event.contentIndex, { id, name })
|
||||
yield { type: 'block-start', index: event.contentIndex, blockType: 'tool-call' }
|
||||
break
|
||||
}
|
||||
case 'toolcall_delta': {
|
||||
const known = toolIds.get(event.contentIndex)
|
||||
yield {
|
||||
type: 'tool-call-delta',
|
||||
index: event.contentIndex,
|
||||
id: CallId(known?.id ?? ''),
|
||||
...known?.name !== undefined && known.name.length > 0 ? { name: known.name } : {},
|
||||
argumentsDelta: event.delta,
|
||||
}
|
||||
break
|
||||
}
|
||||
case 'toolcall_end':
|
||||
yield {
|
||||
type: 'block-end',
|
||||
index: event.contentIndex,
|
||||
block: {
|
||||
type: 'tool-call',
|
||||
id: CallId(event.toolCall.id),
|
||||
name: event.toolCall.name,
|
||||
// pi-ai hands back the PARSED arguments; the harness vocabulary
|
||||
// keeps the raw string.
|
||||
arguments: JSON.stringify(event.toolCall.arguments),
|
||||
},
|
||||
}
|
||||
break
|
||||
case 'done':
|
||||
yield { type: 'usage', usage: mapUsage(event.message.usage) }
|
||||
yield { type: 'finish', reason: mapStopReason(event.message), replayState: toPiReplayState(event.message) }
|
||||
return
|
||||
case 'error':
|
||||
// In-stream error delivery (pi-ai's style) → error finish chunk
|
||||
// (the harness's other sanctioned error path besides throwing).
|
||||
yield { type: 'usage', usage: mapUsage(event.error.usage) }
|
||||
yield { type: 'finish', reason: mapStopReason(event.error) }
|
||||
return
|
||||
// no default: AssistantMessageEvent is pi-ai's closed union; a new
|
||||
// event type should fail compilation here via tsc's exhaustiveness
|
||||
// when one is added (switch covers all current variants).
|
||||
}
|
||||
}
|
||||
throw new LlmError('pi-ai event stream ended without done/error', 'STREAM_CLOSED')
|
||||
}
|
||||
Reference in New Issue
Block a user