| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728 |
- /**
- * Web session model-directory and selection behavior: dynamic provider grouping,
- * provider-local catalog failures, logged-target restoration, advisory unlisted
- * models, and the prompt-assembly boundary for a running selection change.
- */
- import { describe, expect, it, vi } from 'vitest'
- import { Context, FiberState } from 'cordis'
- import type { Fiber } from 'cordis'
- import AgentRegistry, { agentEvents, installAgentLlmTarget } from '@deepseek-ai/dsh-agent'
- import type { Agent, AgentLlmTargetRef } from '@deepseek-ai/dsh-agent'
- import LlmService, { LlmAdapter, ReasoningEffortId } from '@deepseek-ai/dsh-llm'
- import type {
- GenerateOptions, LlmCallConfig, LlmModelInfo, LlmModelReasoningInfo, LlmProviderInfo,
- LlmResolvedModelInfo, StreamChunk,
- } from '@deepseek-ai/dsh-llm'
- import SessionStore, { SessionId } from '@deepseek-ai/dsh-session'
- import type { Session } from '@deepseek-ai/dsh-session'
- import SystemPrompt from '@deepseek-ai/dsh-system-prompt'
- import UserInteractionService from '@deepseek-ai/dsh-user-interaction'
- import type { RpcRequest } from '@deepseek-ai/dsh-host-apiproxy/api/rpc'
- import { RpcId } from '@deepseek-ai/dsh-host-apiproxy/api/rpc'
- import type { MuxFrame } from '@deepseek-ai/dsh-host-apiproxy/api'
- import { createApiProxy } from '../src/api-proxy.ts'
- let nextRpc = 1
- function request<P>(payload: P): RpcRequest<P> {
- return { rpcId: RpcId(`models-${String(nextRpc++)}`), payload }
- }
- class CatalogAdapter extends LlmAdapter {
- constructor(
- private readonly name: string,
- private readonly models: readonly LlmModelInfo[] | Error,
- private readonly reasoning?: LlmModelReasoningInfo,
- private readonly exactError?: Error,
- ) {
- super()
- }
- override providerInfo(provider: string): LlmProviderInfo {
- return { id: provider, name: this.name }
- }
- override listModels(): Promise<readonly LlmModelInfo[]> {
- return this.models instanceof Error
- ? Promise.reject(this.models)
- : Promise.resolve(this.models)
- }
- override resolveModel(provider: string, model: string): Promise<LlmResolvedModelInfo> {
- if (this.exactError !== undefined) return Promise.reject(this.exactError)
- return Promise.resolve({
- provider,
- id: model,
- name: model,
- context: { contextWindow: model === 'private-preview' ? 128_000 : 64_000 },
- ...this.reasoning === undefined ? {} : { reasoning: this.reasoning },
- })
- }
- override async *stream(_options: GenerateOptions): AsyncIterable<StreamChunk> {
- // Catalog tests never enter provider streaming.
- }
- }
- class DeferredCatalogAdapter extends CatalogAdapter {
- readonly pending: PromiseWithResolvers<LlmResolvedModelInfo>[] = []
- constructor() {
- super('Deferred', [
- { provider: 'deferred', id: 'lifecycle-model', name: 'Lifecycle model' },
- ])
- }
- override resolveModel(_provider: string, _model: string): Promise<LlmResolvedModelInfo> {
- const result = Promise.withResolvers<LlmResolvedModelInfo>()
- this.pending.push(result)
- return result.promise
- }
- resolve(index: number, contextWindow: number): void {
- const pending = this.pending[index]
- if (pending === undefined) throw new Error(`no pending resolution at index ${String(index)}`)
- pending.resolve({
- provider: 'deferred',
- id: 'lifecycle-model',
- name: 'Lifecycle model',
- context: { contextWindow },
- })
- }
- }
- const REASONING: LlmModelReasoningInfo = {
- efforts: [
- { id: ReasoningEffortId('off'), name: 'Off' },
- { id: ReasoningEffortId('high'), name: 'High' },
- { id: ReasoningEffortId('max'), name: 'Max' },
- ],
- defaultEffort: ReasoningEffortId('high'),
- }
- async function hostContext(
- onSessions?: (fiber: Fiber) => void,
- onAgents?: (fiber: Fiber) => void,
- ): Promise<Context> {
- const ctx = new Context()
- const sessionsFiber = await ctx.plugin(SessionStore)
- onSessions?.(sessionsFiber)
- await ctx.plugin(SystemPrompt, { persona: '' })
- await ctx.plugin(LlmService)
- await ctx.plugin(UserInteractionService)
- const agentsFiber = await ctx.plugin(AgentRegistry)
- onAgents?.(agentsFiber)
- ctx.llm.registerAdapter(['deepseek'], new CatalogAdapter('DeepSeek', [
- { provider: 'deepseek', id: 'deepseek-chat', name: 'DeepSeek Chat' },
- { provider: 'deepseek', id: 'deepseek-reasoner', name: 'DeepSeek Reasoner', description: 'Reasoning model' },
- ], REASONING))
- ctx.llm.registerAdapter(['broken'], new CatalogAdapter('Broken Provider', new Error('catalog offline')))
- ctx.llm.registerAdapter(['metadata-broken'], new CatalogAdapter('Metadata Broken', [
- { provider: 'metadata-broken', id: 'listed', name: 'Listed' },
- ], undefined, new Error('reasoning metadata offline')))
- ctx.llm.registerAdapter(['empty'], new CatalogAdapter('Empty Provider', []))
- ctx.llm.registerAdapter(['duplicate'], new CatalogAdapter('Duplicate Provider', [
- { provider: 'duplicate', id: 'same', name: 'Same' },
- { provider: 'duplicate', id: 'same', name: 'Same Again' },
- ]))
- return ctx
- }
- async function harness(logged?: {
- provider: string
- model: string
- reasoningEffort?: ReasoningEffortId
- }): Promise<{
- ctx: Context
- agent: Agent
- sessionId: SessionId
- }> {
- const ctx = await hostContext()
- const session = ctx.sessions.create()
- if (logged !== undefined) {
- session.append('request/header', { header: { config: logged }, reason: 'initial' })
- }
- const agent = {
- id: session.id,
- session,
- status: 'running',
- ctx,
- } as Agent
- ctx.agents.register(agent)
- return { ctx, agent, sessionId: session.id }
- }
- function expectValue<T>(response: { result: { ok: true; value: T } | { ok: false } }): T {
- if (!response.result.ok) throw new Error('expected successful response')
- return response.result.value
- }
- async function nextMetrics(
- iterator: AsyncIterator<RpcRequest<MuxFrame>>,
- ): Promise<Extract<MuxFrame, { type: 'session/metrics' }>['metrics']> {
- for (;;) {
- const next = await iterator.next()
- if (next.done) throw new Error('mux ended before a metrics frame')
- if (next.value.payload.type === 'session/metrics') return next.value.payload.metrics
- }
- }
- function attachLifecycleSession(
- ctx: Context,
- sessionId: SessionId,
- withMarker = false,
- ): { session: Session; detach: () => void } {
- const session = ctx.sessions.prepare(sessionId)
- session.append('request/header', {
- header: { config: { provider: 'deferred', model: 'lifecycle-model' } },
- reason: 'initial',
- })
- if (withMarker) {
- session.append('user/message', {
- content: [{ type: 'text', text: 'replacement marker' }],
- source: { kind: 'plugin', plugin: 'test' },
- }, { surfaceOp: 'append' })
- }
- const detach = ctx.sessions.enter(session)
- ctx.sessions.announce(session)
- return { session, detach }
- }
- function attachLifecycleAgent(
- ctx: Context,
- session: Session,
- ): () => void {
- const agent = {
- id: session.id,
- session,
- status: 'running',
- ctx,
- } as Agent
- const detach = ctx.agents.enter(agent, undefined)
- ctx.agents.announce(agent)
- return detach
- }
- function settleCapacityCompletion(): Promise<void> {
- return new Promise<void>((resolve) => { setImmediate(resolve) })
- }
- function installDeferredAdapter(
- ctx: Context,
- adapter: DeferredCatalogAdapter,
- ): Fiber & PromiseLike<Fiber> {
- return ctx.plugin(Object.assign((inner: Context) => {
- inner.llm.registerAdapter(['deferred'], adapter)
- }, { inject: ['llm'] }))
- }
- describe('Web session model selection', () => {
- it('groups successful providers, isolates failures, and preserves an unlisted current model', async () => {
- const { ctx, sessionId } = await harness({
- provider: 'deepseek',
- model: 'private-preview',
- reasoningEffort: ReasoningEffortId('max'),
- })
- const api = createApiProxy(ctx, { provider: 'deepseek', model: 'deepseek-chat', cwd: '/tmp', workspaceRoot: '/tmp' })
- const catalog = expectValue(await api.sessions.models(request({ sessionId })))
- expect(catalog.current).toEqual({
- provider: 'deepseek',
- model: 'private-preview',
- reasoningEffort: 'max',
- })
- expect(catalog.groups).toEqual([{
- id: 'deepseek',
- name: 'DeepSeek',
- models: [
- { id: 'deepseek-chat', name: 'DeepSeek Chat', reasoning: REASONING },
- {
- id: 'deepseek-reasoner',
- name: 'DeepSeek Reasoner',
- description: 'Reasoning model',
- reasoning: REASONING,
- },
- {
- id: 'private-preview',
- name: 'private-preview',
- unlisted: true,
- reasoning: REASONING,
- },
- ],
- }])
- expect(catalog.failures).toEqual([
- { id: 'broken', name: 'Broken Provider', message: 'catalog offline' },
- { id: 'metadata-broken', name: 'Metadata Broken', message: 'reasoning metadata offline' },
- {
- id: 'duplicate',
- name: 'Duplicate Provider',
- message: 'adapter returned invalid or duplicate model metadata for provider "duplicate"',
- },
- ])
- await ctx.fiber.dispose()
- })
- it('accepts an advisory-unlisted model, rejects an unavailable provider, and switches only after the next assembly', async () => {
- const { ctx, agent, sessionId } = await harness()
- const api = createApiProxy(ctx, { provider: 'deepseek', model: 'deepseek-chat', cwd: '/tmp', workspaceRoot: '/tmp' })
- const seed: LlmCallConfig = { provider: 'seed', model: 'seed', temperature: 0.2 }
- const signal = new AbortController().signal
- expect(expectValue(await api.sessions.models(request({ sessionId }))).current)
- .toEqual({ provider: 'deepseek', model: 'deepseek-chat' })
- expect((await ctx.systemPrompt.assemble()).variables)
- .toMatchObject({ provider: 'deepseek', model: 'deepseek-chat' })
- const selected = expectValue(await api.sessions.selectModel(request({
- sessionId,
- provider: 'deepseek',
- model: 'private-preview',
- reasoningEffort: 'max',
- })))
- expect(selected.selected).toEqual({
- provider: 'deepseek',
- model: 'private-preview',
- reasoningEffort: 'max',
- })
- await expect(agentEvents(ctx, agent).waterfall(
- 'agent/request', 1, 0, signal, () => Promise.resolve(seed),
- )).resolves.toMatchObject({ provider: 'deepseek', model: 'deepseek-chat' })
- expect((await ctx.systemPrompt.assemble()).variables)
- .toMatchObject({ provider: 'deepseek', model: 'private-preview' })
- await expect(agentEvents(ctx, agent).waterfall(
- 'agent/request', 1, 1, signal, () => Promise.resolve(seed),
- )).resolves.toMatchObject({
- provider: 'deepseek',
- model: 'private-preview',
- reasoningEffort: 'max',
- })
- const unsupported = await api.sessions.selectModel(request({
- sessionId,
- provider: 'deepseek',
- model: 'private-preview',
- reasoningEffort: 'medium',
- }))
- expect(unsupported.result).toMatchObject({
- ok: false,
- error: {
- code: 'model-unavailable',
- message: 'provider "deepseek" model "private-preview" does not support reasoning effort "medium"',
- },
- })
- const rejected = await api.sessions.selectModel(request({
- sessionId,
- provider: 'missing',
- model: 'model',
- }))
- expect(rejected.result).toEqual({
- ok: false,
- error: {
- code: 'model-unavailable',
- message: 'no adapter registered for provider "missing"',
- details: { provider: 'missing', model: 'model' },
- },
- })
- expect(expectValue(await api.sessions.models(request({ sessionId }))).current)
- .toEqual({ provider: 'deepseek', model: 'private-preview', reasoningEffort: 'max' })
- await ctx.fiber.dispose()
- })
- it('publishes unknown capacity immediately on selection, then the exact selected route capacity', async () => {
- const { ctx, sessionId } = await harness()
- const api = createApiProxy(ctx, { provider: 'deepseek', model: 'deepseek-chat', cwd: '/tmp', workspaceRoot: '/tmp' })
- expectValue(await api.sessions.models(request({ sessionId })))
- const controller = new AbortController()
- const iterator = api.events.mux(request({}), controller.signal)[Symbol.asyncIterator]()
- expect((await nextMetrics(iterator)).contextWindow).toBeUndefined()
- expect((await nextMetrics(iterator)).contextWindow).toBe(64_000)
- expectValue(await api.sessions.selectModel(request({
- sessionId,
- provider: 'deepseek',
- model: 'private-preview',
- })))
- expect((await nextMetrics(iterator)).contextWindow).toBeUndefined()
- expect((await nextMetrics(iterator)).contextWindow).toBe(128_000)
- controller.abort()
- await iterator.return?.()
- await ctx.fiber.dispose()
- })
- it('uses logged capacity without installing Web routing while scheduling foreign metrics', async () => {
- const ctx = await hostContext()
- const api = createApiProxy(ctx, {
- provider: 'deepseek',
- model: 'deepseek-chat',
- cwd: '/tmp',
- workspaceRoot: '/tmp',
- })
- const controller = new AbortController()
- const iterator = api.events.mux(request({}), controller.signal)[Symbol.asyncIterator]()
- const initialMetrics = nextMetrics(iterator)
- const session = ctx.sessions.create()
- expect((await initialMetrics).contextWindow).toBeUndefined()
- session.append('request/header', {
- header: { config: { provider: 'deepseek', model: 'private-preview' } },
- reason: 'change',
- })
- const foreign = {
- id: session.id,
- session,
- status: 'running',
- ctx,
- } as Agent
- const foreignTarget: AgentLlmTargetRef = {
- current: { provider: 'foreign', model: 'foreign-model' },
- assembled: undefined,
- }
- const disposeForeignTarget = installAgentLlmTarget(foreign.ctx, foreignTarget)
- const scheduledMetrics = nextMetrics(iterator)
- ctx.agents.register(foreign)
- expect((await scheduledMetrics).contextWindow).toBeUndefined()
- expect((await nextMetrics(iterator)).contextWindow).toBe(128_000)
- expect((await ctx.systemPrompt.assemble()).variables)
- .toMatchObject({ provider: 'foreign', model: 'foreign-model' })
- const seed: LlmCallConfig = { provider: 'seed', model: 'seed', temperature: 0.2 }
- const signal = new AbortController().signal
- await expect(agentEvents(ctx, foreign).waterfall(
- 'agent/request', 1, 0, signal, () => Promise.resolve(seed),
- )).resolves.toMatchObject({ provider: 'foreign', model: 'foreign-model' })
- disposeForeignTarget()
- expect((await ctx.systemPrompt.assemble()).variables).not.toHaveProperty('provider')
- await expect(agentEvents(ctx, foreign).waterfall(
- 'agent/request', 1, 1, signal, () => Promise.resolve(seed),
- )).resolves.toBe(seed)
- controller.abort()
- await iterator.return?.()
- await ctx.fiber.dispose()
- })
- it('drops capacity completion from a replaced agent that retains the exact session', async () => {
- const ctx = await hostContext()
- const deferred = new DeferredCatalogAdapter()
- ctx.llm.registerAdapter(['deferred'], deferred)
- const lifecycle = attachLifecycleSession(ctx, SessionId('capacity-agent-lifecycle'))
- const retire = attachLifecycleAgent(ctx, lifecycle.session)
- const api = createApiProxy(ctx, { provider: 'deepseek', model: 'deepseek-chat', cwd: '/tmp', workspaceRoot: '/tmp' })
- const controller = new AbortController()
- const iterator = api.events.mux(request({}), controller.signal)[Symbol.asyncIterator]()
- expect((await nextMetrics(iterator)).contextWindow).toBeUndefined()
- await vi.waitFor(() => { expect(deferred.pending).toHaveLength(1) })
- retire()
- const detachLive = attachLifecycleAgent(ctx, lifecycle.session)
- expect((await nextMetrics(iterator)).contextWindow).toBeUndefined()
- await vi.waitFor(() => { expect(deferred.pending).toHaveLength(2) })
- deferred.resolve(0, 64_000)
- await settleCapacityCompletion()
- deferred.resolve(1, 128_000)
- await settleCapacityCompletion()
- expect((await nextMetrics(iterator)).contextWindow).toBe(128_000)
- controller.abort()
- await iterator.return?.()
- detachLive()
- lifecycle.detach()
- await ctx.fiber.dispose()
- })
- it('drops capacity completion from a replaced session while its old agent remains live', async () => {
- const ctx = await hostContext()
- const deferred = new DeferredCatalogAdapter()
- ctx.llm.registerAdapter(['deferred'], deferred)
- const sessionId = SessionId('capacity-session-lifecycle')
- const retiredSession = attachLifecycleSession(ctx, sessionId)
- const retireAgent = attachLifecycleAgent(ctx, retiredSession.session)
- const api = createApiProxy(ctx, { provider: 'deepseek', model: 'deepseek-chat', cwd: '/tmp', workspaceRoot: '/tmp' })
- const controller = new AbortController()
- const iterator = api.events.mux(request({}), controller.signal)[Symbol.asyncIterator]()
- expect((await nextMetrics(iterator)).contextWindow).toBeUndefined()
- await vi.waitFor(() => { expect(deferred.pending).toHaveLength(1) })
- retiredSession.detach()
- const liveSession = attachLifecycleSession(ctx, sessionId, true)
- expect((await nextMetrics(iterator)).contextWindow).toBeUndefined()
- deferred.resolve(0, 64_000)
- await settleCapacityCompletion()
- retireAgent()
- const detachLiveAgent = attachLifecycleAgent(ctx, liveSession.session)
- const scheduled = await nextMetrics(iterator)
- expect(scheduled.logRevision).toBe(2)
- expect(scheduled.contextWindow).toBeUndefined()
- await vi.waitFor(() => { expect(deferred.pending).toHaveLength(2) })
- deferred.resolve(1, 128_000)
- await settleCapacityCompletion()
- expect((await nextMetrics(iterator)).contextWindow).toBe(128_000)
- controller.abort()
- await iterator.return?.()
- detachLiveAgent()
- liveSession.detach()
- await ctx.fiber.dispose()
- })
- it('does not project retired agent capacity into replacement session snapshots', async () => {
- const ctx = await hostContext()
- const deferred = new DeferredCatalogAdapter()
- ctx.llm.registerAdapter(['deferred'], deferred)
- const sessionId = SessionId('capacity-snapshot-lifecycle')
- const retiredSession = attachLifecycleSession(ctx, sessionId)
- const retireAgent = attachLifecycleAgent(ctx, retiredSession.session)
- const api = createApiProxy(ctx, { provider: 'deepseek', model: 'deepseek-chat', cwd: '/tmp', workspaceRoot: '/tmp' })
- const primaryController = new AbortController()
- const primary = api.events.mux(request({}), primaryController.signal)[Symbol.asyncIterator]()
- expect((await nextMetrics(primary)).contextWindow).toBeUndefined()
- await vi.waitFor(() => { expect(deferred.pending).toHaveLength(1) })
- deferred.resolve(0, 64_000)
- await settleCapacityCompletion()
- expect((await nextMetrics(primary)).contextWindow).toBe(64_000)
- retiredSession.detach()
- const replacement = attachLifecycleSession(ctx, sessionId)
- const createdBaseline = await nextMetrics(primary)
- replacement.session.append('user/message', {
- content: [{ type: 'text', text: 'replacement marker' }],
- source: { kind: 'plugin', plugin: 'test' },
- }, { surfaceOp: 'append' })
- const scheduledFlush = await nextMetrics(primary)
- const reconnectController = new AbortController()
- const reconnect = api.events.mux(request({}), reconnectController.signal)[Symbol.asyncIterator]()
- const reconnectBaseline = await nextMetrics(reconnect)
- expect(createdBaseline.logRevision).toBe(1)
- for (const metrics of [scheduledFlush, reconnectBaseline]) {
- expect(metrics.logRevision).toBe(2)
- }
- expect({
- created: createdBaseline.contextWindow,
- scheduled: scheduledFlush.contextWindow,
- reconnect: reconnectBaseline.contextWindow,
- }).toEqual({ created: undefined, scheduled: undefined, reconnect: undefined })
- retireAgent()
- const detachReplacementAgent = attachLifecycleAgent(ctx, replacement.session)
- expect((await nextMetrics(primary)).contextWindow).toBeUndefined()
- await vi.waitFor(() => { expect(deferred.pending).toHaveLength(2) })
- deferred.resolve(1, 128_000)
- await settleCapacityCompletion()
- expect((await nextMetrics(primary)).contextWindow).toBe(128_000)
- primaryController.abort()
- reconnectController.abort()
- await primary.return?.()
- await reconnect.return?.()
- detachReplacementAgent()
- replacement.detach()
- await ctx.fiber.dispose()
- })
- it('refreshes same-route capacity after adapter owner replacement', async () => {
- const ctx = await hostContext()
- const retiredAdapter = new DeferredCatalogAdapter()
- const retiredFiber = await installDeferredAdapter(ctx, retiredAdapter)
- const lifecycle = attachLifecycleSession(ctx, SessionId('capacity-adapter-lifecycle'))
- const detachAgent = attachLifecycleAgent(ctx, lifecycle.session)
- const api = createApiProxy(ctx, { provider: 'deepseek', model: 'deepseek-chat', cwd: '/tmp', workspaceRoot: '/tmp' })
- const controller = new AbortController()
- const iterator = api.events.mux(request({}), controller.signal)[Symbol.asyncIterator]()
- expect((await nextMetrics(iterator)).contextWindow).toBeUndefined()
- await vi.waitFor(() => { expect(retiredAdapter.pending).toHaveLength(1) })
- retiredAdapter.resolve(0, 64_000)
- await settleCapacityCompletion()
- expect((await nextMetrics(iterator)).contextWindow).toBe(64_000)
- await retiredFiber.dispose()
- expect((await nextMetrics(iterator)).contextWindow).toBeUndefined()
- const replacementAdapter = new DeferredCatalogAdapter()
- const replacementFiber = await installDeferredAdapter(ctx, replacementAdapter)
- expect((await nextMetrics(iterator)).contextWindow).toBeUndefined()
- await vi.waitFor(() => { expect(replacementAdapter.pending).toHaveLength(1) })
- replacementAdapter.resolve(0, 128_000)
- await settleCapacityCompletion()
- expect((await nextMetrics(iterator)).contextWindow).toBe(128_000)
- expect(lifecycle.session.requestHeader()?.config).toMatchObject({
- provider: 'deferred',
- model: 'lifecycle-model',
- })
- controller.abort()
- await iterator.return?.()
- detachAgent()
- lifecycle.detach()
- await replacementFiber.dispose()
- await ctx.fiber.dispose()
- })
- it('does not read the sessions service after its disposal status', async () => {
- let sessionsFiber: Fiber | undefined
- const ctx = await hostContext((fiber) => { sessionsFiber = fiber })
- if (sessionsFiber === undefined) throw new Error('sessions fiber missing')
- createApiProxy(ctx, { provider: 'deepseek', model: 'deepseek-chat', cwd: '/tmp', workspaceRoot: '/tmp' })
- const sessions = ctx.get('sessions')
- if (sessions === undefined) throw new Error('sessions service missing')
- const list = vi.spyOn(sessions, 'list').mockImplementation(() => {
- throw new Error('disposed sessions service read')
- })
- await expect(sessionsFiber.dispose()).resolves.toBeUndefined()
- expect(list).not.toHaveBeenCalled()
- await ctx.fiber.dispose()
- })
- it('invalidates pending capacity before SessionStore teardown can reach its callback', async () => {
- let sessionsFiber: Fiber | undefined
- const ctx = await hostContext((fiber) => { sessionsFiber = fiber })
- if (sessionsFiber === undefined) throw new Error('sessions fiber missing')
- const deferred = new DeferredCatalogAdapter()
- ctx.llm.registerAdapter(['deferred'], deferred)
- const lifecycle = attachLifecycleSession(ctx, SessionId('capacity-session-store-teardown'))
- const detachAgent = attachLifecycleAgent(ctx, lifecycle.session)
- const api = createApiProxy(ctx, { provider: 'deepseek', model: 'deepseek-chat', cwd: '/tmp', workspaceRoot: '/tmp' })
- const controller = new AbortController()
- const iterator = api.events.mux(request({}), controller.signal)[Symbol.asyncIterator]()
- expect((await nextMetrics(iterator)).contextWindow).toBeUndefined()
- await vi.waitFor(() => { expect(deferred.pending).toHaveLength(1) })
- const agents = ctx.get('agents')
- if (agents === undefined) throw new Error('agent registry missing')
- expect(agents.get(lifecycle.session.id)).toBeDefined()
- const getAgent = vi.spyOn(agents, 'get')
- await sessionsFiber.dispose()
- getAgent.mockClear()
- deferred.resolve(0, 64_000)
- await settleCapacityCompletion()
- expect(getAgent).not.toHaveBeenCalled()
- const pendingFrame = iterator.next()
- const outcome = await Promise.race([
- pendingFrame.then(() => 'frame' as const),
- new Promise<'idle'>((resolve) => { setImmediate(() => { resolve('idle') }) }),
- ])
- expect(outcome).toBe('idle')
- controller.abort()
- await expect(pendingFrame).resolves.toMatchObject({ done: true })
- await iterator.return?.()
- detachAgent()
- lifecycle.detach()
- await ctx.fiber.dispose()
- })
- it('publishes unknown metrics after AgentRegistry terminal disposal with mux active', async () => {
- let agentsFiber: Fiber | undefined
- const ctx = await hostContext(undefined, (fiber) => { agentsFiber = fiber })
- if (agentsFiber === undefined) throw new Error('agent registry fiber missing')
- const deferred = new DeferredCatalogAdapter()
- ctx.llm.registerAdapter(['deferred'], deferred)
- const lifecycle = attachLifecycleSession(ctx, SessionId('capacity-agent-registry-disposed'))
- const detachAgent = attachLifecycleAgent(ctx, lifecycle.session)
- const api = createApiProxy(ctx, { provider: 'deepseek', model: 'deepseek-chat', cwd: '/tmp', workspaceRoot: '/tmp' })
- const controller = new AbortController()
- const iterator = api.events.mux(request({}), controller.signal)[Symbol.asyncIterator]()
- expect((await nextMetrics(iterator)).contextWindow).toBeUndefined()
- await vi.waitFor(() => { expect(deferred.pending).toHaveLength(1) })
- deferred.resolve(0, 64_000)
- await settleCapacityCompletion()
- expect((await nextMetrics(iterator)).contextWindow).toBe(64_000)
- await agentsFiber.dispose()
- const refresh = nextMetrics(iterator).then(
- metrics => ({ kind: 'metrics' as const, metrics }),
- () => ({ kind: 'error' as const }),
- )
- const outcome = await Promise.race([
- refresh,
- new Promise<{ kind: 'idle' }>((resolve) => {
- setImmediate(() => { resolve({ kind: 'idle' }) })
- }),
- ])
- controller.abort()
- await refresh
- await iterator.return?.()
- detachAgent()
- lifecycle.detach()
- await ctx.fiber.dispose()
- expect(outcome.kind).toBe('metrics')
- if (outcome.kind === 'metrics') {
- expect(outcome.metrics.contextWindow).toBeUndefined()
- expect(outcome.metrics.logRevision).toBe(1)
- }
- })
- it.each(['agents', 'sessions'] as const)(
- 'drops capacity completion while %s is unavailable during unload',
- async (serviceName) => {
- let sessionsFiber: Fiber | undefined
- let agentsFiber: Fiber | undefined
- const ctx = await hostContext(
- (fiber) => { sessionsFiber = fiber },
- (fiber) => { agentsFiber = fiber },
- )
- const heldFiber = serviceName === 'sessions' ? sessionsFiber : agentsFiber
- if (heldFiber === undefined) throw new Error(`${serviceName} fiber missing`)
- const unloadStarted = Promise.withResolvers<undefined>()
- const releaseUnload = Promise.withResolvers<undefined>()
- heldFiber.ctx.effect(() => () => {
- unloadStarted.resolve(undefined)
- return releaseUnload.promise
- }, `test: hold ${serviceName} unload`)
- const deferred = new DeferredCatalogAdapter()
- ctx.llm.registerAdapter(['deferred'], deferred)
- const lifecycle = attachLifecycleSession(ctx, SessionId(`capacity-${serviceName}-unloading`))
- const detachAgent = attachLifecycleAgent(ctx, lifecycle.session)
- const api = createApiProxy(ctx, { provider: 'deepseek', model: 'deepseek-chat', cwd: '/tmp', workspaceRoot: '/tmp' })
- const controller = new AbortController()
- const iterator = api.events.mux(request({}), controller.signal)[Symbol.asyncIterator]()
- expect((await nextMetrics(iterator)).contextWindow).toBeUndefined()
- await vi.waitFor(() => { expect(deferred.pending).toHaveLength(1) })
- const agents = ctx.get('agents')
- if (agents === undefined) throw new Error('agent registry missing')
- const sessions = ctx.get('sessions')
- if (sessions === undefined) throw new Error('sessions service missing')
- const getAgent = vi.spyOn(agents, 'get')
- const getSession = vi.spyOn(sessions, 'get')
- const disposing = heldFiber.dispose()
- await unloadStarted.promise
- await vi.waitFor(() => { expect(ctx.get(serviceName)).toBeUndefined() })
- expect(heldFiber.state).toBe(FiberState.UNLOADING)
- getAgent.mockClear()
- getSession.mockClear()
- deferred.resolve(0, 64_000)
- await settleCapacityCompletion()
- const agentReads = getAgent.mock.calls.length
- const sessionReads = getSession.mock.calls.length
- const pendingFrame = iterator.next()
- const outcome = await Promise.race([
- pendingFrame.then(() => 'frame' as const),
- new Promise<'idle'>((resolve) => { setImmediate(() => { resolve('idle') }) }),
- ])
- controller.abort()
- await expect(pendingFrame).resolves.toMatchObject({ done: true })
- await iterator.return?.()
- releaseUnload.resolve(undefined)
- await disposing
- detachAgent()
- lifecycle.detach()
- await ctx.fiber.dispose()
- expect(agentReads).toBe(serviceName === 'sessions' ? 1 : 0)
- expect(sessionReads).toBe(0)
- expect(outcome).toBe('idle')
- },
- )
- })
|