/** * LLM service: adapter registry with a waterfall-interceptable streaming call * surface. Exports the `LlmService` default, the abstract `LlmAdapter` for * provider backends, and `BlockAssembler` for chunk assembly. * * @module @deepseek-ai/dsh-llm */ import { Context, Service } from 'cordis' import type { GenerateOptions, LlmConfigurableProvider, LlmFailure, LlmModelContext, LlmModelInfo, LlmResolvedModelInfo, LlmProviderInfo, StreamChunk, } from './types.ts' import { freezeMessage, type Message } from './message.ts' import { resolveRetryPolicy } from './retry-policy.ts' import type { ResolvedRetryPolicy } from './retry-policy.ts' import type { ProviderRequestId } from './brand.ts' import { callConfigEquals, deepFreeze } from './call-config.ts' import type { LlmCallConfig, LlmCallConfigAdapterDefaults } from './call-config.ts' import { HarnessError } from './error.ts' import { bindAdapterFailureScope, markLlmAdapterFailure } from './adapter-failure.ts' import type { AdapterFailureScope } from './adapter-failure.ts' export * from './attribution.ts' export * from './brand.ts' export * from './never.ts' export * from './error.ts' export * from './types.ts' export * from './message.ts' export * from './retry-policy.ts' export { BlockAssembler } from './assembler.ts' export { callConfigEquals, deepFreeze, isAgentLoopRequest, markAgentLoopRequest } from './call-config.ts' export type { LlmCallConfig, LlmCallConfigAdapterDefaults } from './call-config.ts' export { isLlmAdapterFailure, llmFailureOf, llmRetryPolicyOf } from './adapter-failure.ts' declare module 'cordis' { interface Context { llm: LlmService } interface Events { /** * Waterfall around every streaming model call (retry, replay, routing). * Bound to the {@link LlmService}; call `next()` to reach the resolved * adapter's stream, or yield your own chunks to short-circuit. * @param options - the full request. A LOOP-built request carries the * process-local {@link markAgentLoopRequest} identity and arrives deep-frozen * (mutation throws): its content is a pure function of the session log (the * reconstructability Agent Note), so listeners read it, never rewrite it. * Hand-built calls do not carry that marker; their messages already obey * the immutable creation contract. * @mode waterfall */ 'llm/stream'(this: LlmService, options: GenerateOptions, next: () => AsyncIterable): AsyncIterable /** * The provider topology changed: an adapter registered or unregistered * routes, or the configurable-provider directory gained or lost entries. * This is a payload-free registry notification fired at each commit point * (including registration disposal); consumers re-read `listProviders()`, * `listModels()`, or `listConfigurableProviders()` for the new state. * Observer failures are contained and cannot veto the registry mutation. * @mode emit */ 'llm/adapters-updated'(): void } } /** Structured provider facts and cause accepted by {@link LlmError}. */ export interface LlmErrorOptions extends ErrorOptions { /** Valid HTTP status observed at the provider boundary. */ status?: number /** Positive finite provider-requested delay in milliseconds. */ providerRetryAfterMs?: number /** Non-empty opaque provider request id. */ requestId?: ProviderRequestId } /** * Typed error for LLM-related failures. Extends {@link HarnessError}, so the * `code` string (e.g. `AUTH`, `RATE_LIMIT`, `NO_ADAPTER`) is shared taxonomy. */ export class LlmError extends HarnessError { /** Serializable facts retained beside this live Error. */ readonly failure: LlmFailure /** * @param message - non-empty human-readable failure summary. * @param code - non-empty stable provider-neutral machine code. * @param options - optional cause and validated serializable provider facts. */ constructor(message: string, code: string, options?: LlmErrorOptions) { if (typeof message !== 'string' || message.length === 0) throw new Error('LlmError message must be a non-empty string') if (typeof code !== 'string' || code.length === 0) throw new Error('LlmError code must be a non-empty string') if (options?.status !== undefined && (!Number.isInteger(options.status) || options.status < 100 || options.status > 599)) { throw new Error('LlmError status must be an integer from 100 through 599') } if (options?.providerRetryAfterMs !== undefined && (!Number.isFinite(options.providerRetryAfterMs) || options.providerRetryAfterMs <= 0)) { throw new Error('LlmError providerRetryAfterMs must be a positive finite number') } if (options?.requestId !== undefined && (typeof options.requestId !== 'string' || options.requestId.length === 0)) { throw new Error('LlmError requestId must be a non-empty string') } super(message, code, options) this.name = 'LlmError' this.failure = Object.freeze({ message, code, ...options?.status === undefined ? {} : { status: options.status }, ...options?.providerRetryAfterMs === undefined ? {} : { providerRetryAfterMs: options.providerRetryAfterMs }, ...options?.requestId === undefined ? {} : { requestId: options.requestId }, }) } } /** One model call whose config and adapter registration were resolved together. */ export interface PreparedLlmCall { /** Detached, deep-frozen config with any adapter-owned default materialized. */ readonly config: LlmCallConfig /** Detached context metadata resolved with the registration-bound call. */ readonly context?: LlmModelContext /** Config fields materialized by the captured adapter rather than proposed by the caller. */ readonly adapterDefaults: LlmCallConfigAdapterDefaults /** * Dispatch this call once through the registration captured during * preparation. The request's call-config fields must match {@link config}; * reuse or mismatch fails with `INVALID_PREPARED_CALL`. * @param options - fully assembled request carrying the prepared config. * @returns the chunk stream, including the `llm/stream` waterfall. */ stream(options: GenerateOptions): AsyncIterable } /** * Provider-wire adapter for the harness message and stream vocabulary. Register implementations * with `ctx.llm.registerAdapter(providers, adapter)`. Every provider HTTP request must include * `attributionHeaders()`; prove that at the wire or library header-hook boundary. The direct-fetch * DeepSeek and library-backed pi-ai adapters intentionally exercise this contract through different internals. */ export abstract class LlmAdapter { /** * Describe one provider route owned by this adapter. * @param provider - a route passed to `registerAdapter()` for this instance. * @returns detached display metadata whose id must equal `provider`. */ providerInfo(provider: string): LlmProviderInfo { return { id: provider, name: provider } } /** * Return the provider-owned retry policy captured with this route. * @param _provider - a route passed to `registerAdapter()` for this instance. * @returns a resolved policy, or `undefined` to use the normal defaults. */ providerRetryPolicy(_provider: string): ResolvedRetryPolicy | undefined { return undefined } /** * List models this adapter can currently advertise for one owned provider. * The result is advisory: an adapter may accept unlisted model ids, and * consumers must not turn absence into request rejection. * @param _provider - one provider route owned by this adapter. * @returns discoverable models in adapter-preferred order. */ listModels(_provider: string): Promise { return Promise.resolve([]) } /** * Resolve all metadata available for one exact model. This query is * independent of the advisory catalog and does not validate request routing. * @param provider - one provider route owned by this adapter. * @param model - exact model id passed to {@link GenerateOptions.model}. * @param _signal - cancellation for this exact-model lookup; asynchronous * implementations must settle promptly after it aborts. * @returns provider/model identity plus any context, call-default, and reasoning metadata. */ resolveModel( provider: string, model: string, _signal?: AbortSignal, ): Promise { return Promise.resolve({ provider, id: model, name: model }) } /** * Stream one model call as raw chunks. The only required method. * @param options - the fully-assembled request; implementations must honor `options.signal`. * @returns the chunk stream, obeying the adapter contract documented on `StreamChunk`. */ abstract stream(options: GenerateOptions): AsyncIterable } /** * What {@link LlmService.registerAdapter} returns: the disposer, plus an * atomic route replacement for the same adapter instance. */ export interface AdapterRegistrationHandle { /** Release every route this registration currently holds. */ (): void /** * Replace this registration's routes with `providers`, keeping the same * adapter instance. The candidate set is validated in full first — a * conflict with another adapter, an invalid name, or bad provider metadata * throws and leaves the current routes untouched — and the swap itself is * one synchronous section, so no request can observe a gap. An empty array * is legal here (a settings section that emptied holds zero routes while * staying registered), unlike an empty initial registration. * * Throws `LlmError` with code `REGISTRATION_DISPOSED` once the registration * has been released: its routes are gone and its disposer has already run, * so anything registered afterwards would have no owner left to release it. * @param providers - the complete next route set for this registration. */ replace(providers: string[]): void } /** * The abstract `llm` service: an adapter registry plus a streaming model-call * surface, interceptable via the `llm/stream` waterfall. */ export class LlmService extends Service { private adapters = new Map() private directory = new Map() constructor(ctx: Context) { super(ctx, 'llm') } /** Notify topology observers without letting one broken listener veto the commit. */ private emitAdaptersUpdated(): void { // Cordis emit uses Array.map: one synchronous throw starves later // listeners. Registry notifications are non-vetoing, so contain each // callback independently; INVARIANT-coded failures still surface. let invariantFailure: unknown for (const listener of this.ctx.events.dispatch('emit', ['llm/adapters-updated']) as Array<() => unknown>) { try { const returned = listener() if (returned != null && typeof (returned as PromiseLike).then === 'function') { // An emit listener may still be an async function; its rejection // cannot reach the synchronous INVARIANT rethrow below, so it is // contained here instead of becoming an unhandled rejection. void Promise.resolve(returned as PromiseLike).then(undefined, (error: unknown) => { this.warnAdaptersListenerFailure(error) }) } } catch (error) { if ((error as { code?: unknown } | null)?.code === 'INVARIANT') { invariantFailure ??= error continue } this.warnAdaptersListenerFailure(error) } } if (invariantFailure !== undefined) throw invariantFailure as Error } /** Contained-listener diagnostic shared by the sync and async failure paths. */ private warnAdaptersListenerFailure(error: unknown): void { this.ctx.logger.warn('llm: an llm/adapters-updated listener failed') this.ctx.logger.warn(error) } /** * Register an adapter for the given provider routes. Throws `LlmError` with code * `DUPLICATE_ADAPTER` if any provider already has an adapter (all-or-nothing). * Disposed with the fiber. * @param providers - every provider route this adapter should serve. * @param adapter - the adapter that streams calls for those providers. * @returns the disposer, carrying {@link AdapterRegistrationHandle.replace}. */ registerAdapter(providers: string[], adapter: LlmAdapter): AdapterRegistrationHandle { // The routes this registration currently holds; `replace` rewrites it, and // the disposer releases whatever it holds at disposal time. const owned = new Set() // The disposer has run: `owned` being empty cannot say so on its own, // because `replace([])` legally leaves a live registration holding none. let released = false const dispose = this.ctx.effect(function* (this: LlmService) { if (providers.length === 0) throw new LlmError('an adapter must register at least one provider', 'INVALID_ADAPTER') this.commitRoutes(owned, this.prepareRoutes(providers, adapter, owned)) yield () => { released = true for (const provider of owned) this.adapters.delete(provider) owned.clear() this.emitAdaptersUpdated() } }.bind(this), 'llm.registerAdapter()') // ctx.effect's disposer returns Promise; our disposer API is // synchronous fire-and-forget — discard the (always-resolved) promise. const handle = (() => void dispose()) as AdapterRegistrationHandle handle.replace = (next: string[]): void => { // Registering here would leak: the effect's disposer already ran, so // nothing remains to release whatever this call would put in the map. if (released) { throw new LlmError('a disposed adapter registration cannot replace its routes', 'REGISTRATION_DISPOSED') } this.commitRoutes(owned, this.prepareRoutes(next, adapter, owned)) } return handle } /** * Validate one candidate route set for `adapter`, treating routes this * registration already holds as available. Nothing is mutated: a rejected * candidate leaves the registry exactly as it was. */ private prepareRoutes(providers: string[], adapter: LlmAdapter, owned: ReadonlySet): AdapterRegistration[] { const unique = new Set() const registrations: AdapterRegistration[] = [] for (const provider of providers) { if (provider.length === 0) throw new LlmError('adapter provider names must be non-empty', 'INVALID_ADAPTER') if (unique.has(provider) || (this.adapters.has(provider) && !owned.has(provider))) { throw new LlmError(`an adapter for provider "${provider}" is already registered`, 'DUPLICATE_ADAPTER') } const info = adapter.providerInfo(provider) if (typeof info.id !== 'string' || info.id !== provider || typeof info.name !== 'string' || info.name.length === 0) { throw new LlmError(`adapter metadata for provider "${provider}" must preserve its id and have a non-empty name`, 'INVALID_ADAPTER') } unique.add(provider) const retryPolicy = adapter.providerRetryPolicy(provider) ?? resolveRetryPolicy(undefined, `llm: provider "${provider}" retryPolicy`) registrations.push({ adapter, provider: { id: info.id, name: info.name }, retryPolicy, }) } return registrations } /** * Swap this registration's routes for the prepared ones in one synchronous * section, so no observer can see the registry between the release and the * re-registration. The route set's one mutation point is also where * `llm/adapters-updated` is published, so a `replace` announces itself * exactly like a first registration. */ private commitRoutes(owned: Set, registrations: readonly AdapterRegistration[]): void { for (const provider of owned) this.adapters.delete(provider) owned.clear() for (const registration of registrations) { this.adapters.set(registration.provider.id, registration) owned.add(registration.provider.id) } this.emitAdaptersUpdated() } /** * Describe provider routes with a registered adapter. * @returns detached provider metadata in registration order. */ listProviders(): LlmProviderInfo[] { return [...this.adapters.values()].map(({ provider }) => ({ ...provider })) } /** * Declare provider routes an adapter plugin can activate through * configuration. Registration is all-or-nothing: an empty list, invalid * entry, or a provider already declared by any registration throws * `LlmError` without registering the rest. Disposed with the fiber. * @param entries - every configurable provider this plugin owns. * @returns the disposer that withdraws all of them. */ registerConfigurableProviders(entries: readonly LlmConfigurableProvider[]): () => void { const dispose = this.ctx.effect(function* (this: LlmService) { if (entries.length === 0) { throw new LlmError('a configurable-provider registration must declare at least one provider', 'INVALID_DIRECTORY') } const detached: LlmConfigurableProvider[] = [] for (const entry of entries) { if (entry.provider.length === 0 || entry.displayName.length === 0 || entry.settingsNs.length === 0) { throw new LlmError('configurable providers need a non-empty provider, displayName, and settingsNs', 'INVALID_DIRECTORY') } if (entry.settingsPath.some(segment => segment.length === 0)) { throw new LlmError(`configurable provider "${entry.provider}" has an empty settingsPath segment`, 'INVALID_DIRECTORY') } if (this.directory.has(entry.provider) || detached.some(seen => seen.provider === entry.provider)) { throw new LlmError(`configurable provider "${entry.provider}" is already declared`, 'DUPLICATE_DIRECTORY') } detached.push({ ...entry, settingsPath: [...entry.settingsPath] }) } for (const entry of detached) this.directory.set(entry.provider, entry) this.emitAdaptersUpdated() yield () => { for (const entry of detached) this.directory.delete(entry.provider) this.emitAdaptersUpdated() } }.bind(this), 'llm.registerConfigurableProviders()') return () => void dispose() } /** * List every declared configurable provider, registered or dormant. * @returns detached directory entries in declaration order. */ listConfigurableProviders(): LlmConfigurableProvider[] { return [...this.directory.values()].map(entry => ({ ...entry, settingsPath: [...entry.settingsPath] })) } /** * Resolve the retry policy captured when one provider route was registered. * @param provider - registered provider route to inspect. * @returns the provider-owned policy, with normal defaults already resolved. */ providerRetryPolicy(provider: string): ResolvedRetryPolicy { return this.registration(provider).retryPolicy } /** * Discover models advertised by one registered provider. Catalog membership * is advisory and never changes routing or request validation. * @param provider - registered provider route to inspect. * @returns detached model metadata in adapter-preferred order. */ async listModels(provider: string): Promise { const adapter = this.registration(provider).adapter const models = await adapter.listModels(provider) const seen = new Set() return models.map((model) => { if ( typeof model.provider !== 'string' || model.provider !== provider || typeof model.id !== 'string' || model.id.length === 0 || typeof model.name !== 'string' || model.name.length === 0 || (model.description !== undefined && typeof model.description !== 'string') || seen.has(model.id) ) { throw new LlmError(`adapter returned invalid or duplicate model metadata for provider "${provider}"`, 'INVALID_CATALOG') } seen.add(model.id) return { provider: model.provider, id: model.id, name: model.name, ...model.description === undefined ? {} : { description: model.description }, } }) } /** * Resolve and validate all metadata from the adapter that owns one exact * route. The result is detached from adapter-owned objects; catalog * membership remains advisory and does not control request routing. * @param provider - registered provider route to inspect. * @param model - exact model id passed to the adapter. * @param signal - optional cancellation for adapter-owned asynchronous lookup. * @returns exact model identity plus available context and reasoning metadata. */ async resolveModelInfo( provider: string, model: string, signal?: AbortSignal, ): Promise { return this.resolveModelInfoFor(this.registration(provider), model, signal) } private async resolveModelInfoFor( registration: AdapterRegistration, model: string, signal?: AbortSignal, ): Promise { const provider = registration.provider.id const resolved = await registration.adapter.resolveModel(provider, model, signal) if ( typeof resolved.provider !== 'string' || resolved.provider !== provider || typeof resolved.id !== 'string' || resolved.id !== model || typeof resolved.name !== 'string' || resolved.name.length === 0 || (resolved.description !== undefined && typeof resolved.description !== 'string') ) { throw new LlmError( `adapter returned invalid exact model metadata for provider "${provider}" model "${model}"`, 'INVALID_MODEL_INFO', ) } const context = resolved.context if (context !== undefined && (!Number.isInteger(context.contextWindow) || context.contextWindow <= 0)) { throw new LlmError( `adapter returned invalid context metadata for provider "${provider}" model "${model}"`, 'INVALID_MODEL_CONTEXT', ) } const defaultMaxTokens = resolved.defaultMaxTokens if (defaultMaxTokens !== undefined && (!Number.isSafeInteger(defaultMaxTokens) || defaultMaxTokens <= 0)) { throw new LlmError( `adapter returned invalid default maxTokens for provider "${provider}" model "${model}"`, 'INVALID_MODEL_MAX_TOKENS', ) } const info: LlmResolvedModelInfo = { provider, id: model, name: resolved.name, ...resolved.description === undefined ? {} : { description: resolved.description }, ...context === undefined ? {} : { context: { contextWindow: context.contextWindow } }, ...defaultMaxTokens === undefined ? {} : { defaultMaxTokens }, } const reasoning = resolved.reasoning if (reasoning === undefined) return info if (reasoning.efforts.length === 0) { throw new LlmError( `adapter returned invalid reasoning metadata for provider "${provider}" model "${model}"`, 'INVALID_MODEL_REASONING', ) } const seen = new Set() const efforts = reasoning.efforts.map((effort) => { if ( typeof effort.id !== 'string' || effort.id.length === 0 || typeof effort.name !== 'string' || effort.name.length === 0 || (effort.description !== undefined && typeof effort.description !== 'string') || seen.has(effort.id) ) { throw new LlmError( `adapter returned invalid or duplicate reasoning effort metadata for provider "${provider}" model "${model}"`, 'INVALID_MODEL_REASONING', ) } seen.add(effort.id) return { id: effort.id, name: effort.name, ...effort.description === undefined ? {} : { description: effort.description }, } }) if (reasoning.defaultEffort !== undefined && !seen.has(reasoning.defaultEffort)) { throw new LlmError( `adapter returned an unknown default reasoning effort for provider "${provider}" model "${model}"`, 'INVALID_MODEL_REASONING', ) } return { ...info, reasoning: { efforts, ...reasoning.defaultEffort === undefined ? {} : { defaultEffort: reasoning.defaultEffort }, }, } } /** * Validate a conversation call config against its exact model capability and * materialize adapter-configured defaults. Unsupported explicit efforts * reject before provider I/O; no clamping or aliasing is performed. This * standalone query does not bind a later dispatch; use {@link prepareCall} * when logging and streaming must share one adapter registration. * @param config - provider/model route and optional request controls. * @param signal - optional cancellation for adapter-owned capability lookup. * @returns a detached config only when a default must be materialized. */ async resolveCallConfig(config: LlmCallConfig, signal?: AbortSignal): Promise { return (await this.resolveCallFor(this.registration(config.provider), config, signal)).config } private async resolveCallFor( registration: AdapterRegistration, config: LlmCallConfig, signal?: AbortSignal, ): Promise<{ config: LlmCallConfig; context?: LlmModelContext }> { const info = await this.resolveModelInfoFor(registration, config.model, signal) const defaulted = config.maxTokens === undefined && info.defaultMaxTokens !== undefined ? { ...config, maxTokens: info.defaultMaxTokens } : config const reasoning = info.reasoning const requested = defaulted.reasoningEffort let resolvedConfig = defaulted if (reasoning === undefined) { if (requested !== undefined) { throw new LlmError( `provider "${config.provider}" model "${config.model}" does not support reasoning effort "${requested}"`, 'UNSUPPORTED_REASONING_EFFORT', ) } } else { const effective = requested ?? reasoning.defaultEffort if (effective !== undefined) { if (!reasoning.efforts.some(effort => effort.id === effective)) { throw new LlmError( `provider "${config.provider}" model "${config.model}" does not support reasoning effort "${effective}"`, 'UNSUPPORTED_REASONING_EFFORT', ) } if (requested !== effective) resolvedConfig = { ...defaulted, reasoningEffort: effective } } } return { config: resolvedConfig, ...info.context === undefined ? {} : { context: info.context }, } } /** * Resolve one call under its current adapter registration. The returned * one-shot handle keeps that registration across header logging and dispatch, * so HMR cannot combine one adapter's capability result with another adapter. * @param config - provider/model route and optional request controls. * @param signal - optional cancellation for adapter-owned capability lookup. * @returns a prepared config and its registration-bound stream entry point. */ async prepareCall(config: LlmCallConfig, signal?: AbortSignal): Promise { const registration = this.registration(config.provider) const resolved = await this.resolveCallFor(registration, config, signal) const resolvedConfig = deepFreeze(structuredClone(resolved.config)) const context = resolved.context === undefined ? undefined : deepFreeze(structuredClone(resolved.context)) const adapterDefaults = deepFreeze({ ...config.reasoningEffort === undefined && resolvedConfig.reasoningEffort !== undefined ? { reasoningEffort: true } : {}, ...config.maxTokens === undefined && resolvedConfig.maxTokens !== undefined ? { maxTokens: true } : {}, }) let dispatched = false return Object.freeze({ config: resolvedConfig, adapterDefaults, ...context === undefined ? {} : { context }, stream: (options: GenerateOptions): AsyncIterable => { if (dispatched) { throw new LlmError('a prepared LLM call can only be dispatched once', 'INVALID_PREPARED_CALL') } dispatched = true return this.streamWithRegistration(options, { registration, config: resolvedConfig }) }, }) } private registration(provider: string): AdapterRegistration { const registration = this.adapters.get(provider) if (!registration) throw new LlmError(`no adapter registered for provider "${provider}"`, 'NO_ADAPTER') return registration } /** Remove replay state whose historical route is owned by another adapter. */ private forAdapter(options: GenerateOptions, adapter: LlmAdapter): GenerateOptions { const messages: Message[] = options.messages.map((message) => { const source = message.source if (message.role !== 'assistant' || source.kind !== 'model' || source.replayState === undefined) return message if (this.adapters.get(source.provider)?.adapter === adapter) return message return freezeMessage({ ...message, source: { kind: 'model', provider: source.provider, model: source.model }, }) }) if (messages.every((message, index) => message === options.messages[index])) return options const filtered = { ...options, messages } return Object.isFrozen(options) ? deepFreeze(filtered) : filtered } /** * Final adapter boundary. It tags only failures from adapter selection, * synchronous dispatch, iterator construction, or iteration while preserving * the original Error object. Middleware outside this generator remains * distinguishable as plugin work. An iteration failure skips adapter cleanup * so it cannot suppress the primary provider error. A downstream close awaits * adapter cleanup, whose failures remain ordinary untagged work. */ private async * adapterStream( options: GenerateOptions, failures: AdapterFailureScope, prepared?: { registration: AdapterRegistration; config: LlmCallConfig }, ): AsyncGenerator { let iterator: AsyncIterator try { const registration = prepared?.registration ?? this.registration(options.provider) failures.retryPolicy = registration.retryPolicy const resolvedConfig = prepared === undefined ? (await this.resolveCallFor(registration, options, options.signal)).config : prepared.config if (prepared !== undefined && !callConfigEquals(options, resolvedConfig)) { throw new LlmError( 'prepared LLM call config changed before adapter dispatch', 'INVALID_PREPARED_CALL', ) } const resolvedOptions = callConfigEquals(options, resolvedConfig) ? options : Object.isFrozen(options) ? deepFreeze({ ...options, ...resolvedConfig }) : { ...options, ...resolvedConfig } const adapter = registration.adapter const stream = adapter.stream(this.forAdapter(resolvedOptions, adapter)) iterator = stream[Symbol.asyncIterator]() } catch (error: unknown) { throw markLlmAdapterFailure(failures, error) } let completed = false let iterationFailed = false try { while (true) { let value: StreamChunk try { const item = await iterator.next() if (item.done) { completed = true return } value = item.value } catch (error: unknown) { iterationFailed = true throw markLlmAdapterFailure(failures, error) } // End the adapter-owned try before yielding: consumer/middleware // failures resumed into this generator must remain untagged. yield value } } finally { // oxlint-disable-next-line typescript/no-unnecessary-condition -- the iteration catch sets its latch before entering finally. if (!completed && !iterationFailed) { const close = iterator.return?.bind(iterator) if (close) await close() } } } /** * Stream one model call as raw chunks (token-level deltas). Throws * `LlmError` with code `NO_ADAPTER` if no adapter is registered for * `options.provider`. Replay state is retained only when the same adapter * instance owns its historical provider and the target provider. Final * adapter selection remains fixed through asynchronous exact-model resolution * and dispatch. Selection, dispatch, and iteration failures retain their * original Error identity and are tagged in a call-local scope for narrow * agent-loop request recovery; middleware and nested-call failures remain * untagged for the outer call. * @param options - the full request; `options.provider` selects the adapter. * @returns the chunk stream, possibly wrapped by `llm/stream` listeners. */ stream(options: GenerateOptions): AsyncIterable { return this.streamWithRegistration(options) } private streamWithRegistration( options: GenerateOptions, prepared?: { registration: AdapterRegistration; config: LlmCallConfig }, ): AsyncIterable { const failures: AdapterFailureScope = { failures: new WeakMap() } const stream = this.ctx.waterfall( this, 'llm/stream', options, () => this.adapterStream(options, failures, prepared), ) return bindAdapterFailureScope(stream, failures) } } interface AdapterRegistration { readonly adapter: LlmAdapter readonly provider: LlmProviderInfo readonly retryPolicy: ResolvedRetryPolicy } export default LlmService