api-proxy-models.spec.ts 29 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728
  1. /**
  2. * Web session model-directory and selection behavior: dynamic provider grouping,
  3. * provider-local catalog failures, logged-target restoration, advisory unlisted
  4. * models, and the prompt-assembly boundary for a running selection change.
  5. */
  6. import { describe, expect, it, vi } from 'vitest'
  7. import { Context, FiberState } from 'cordis'
  8. import type { Fiber } from 'cordis'
  9. import AgentRegistry, { agentEvents, installAgentLlmTarget } from '@deepseek-ai/dsh-agent'
  10. import type { Agent, AgentLlmTargetRef } from '@deepseek-ai/dsh-agent'
  11. import LlmService, { LlmAdapter, ReasoningEffortId } from '@deepseek-ai/dsh-llm'
  12. import type {
  13. GenerateOptions, LlmCallConfig, LlmModelInfo, LlmModelReasoningInfo, LlmProviderInfo,
  14. LlmResolvedModelInfo, StreamChunk,
  15. } from '@deepseek-ai/dsh-llm'
  16. import SessionStore, { SessionId } from '@deepseek-ai/dsh-session'
  17. import type { Session } from '@deepseek-ai/dsh-session'
  18. import SystemPrompt from '@deepseek-ai/dsh-system-prompt'
  19. import UserInteractionService from '@deepseek-ai/dsh-user-interaction'
  20. import type { RpcRequest } from '@deepseek-ai/dsh-host-apiproxy/api/rpc'
  21. import { RpcId } from '@deepseek-ai/dsh-host-apiproxy/api/rpc'
  22. import type { MuxFrame } from '@deepseek-ai/dsh-host-apiproxy/api'
  23. import { createApiProxy } from '../src/api-proxy.ts'
  24. let nextRpc = 1
  25. function request<P>(payload: P): RpcRequest<P> {
  26. return { rpcId: RpcId(`models-${String(nextRpc++)}`), payload }
  27. }
  28. class CatalogAdapter extends LlmAdapter {
  29. constructor(
  30. private readonly name: string,
  31. private readonly models: readonly LlmModelInfo[] | Error,
  32. private readonly reasoning?: LlmModelReasoningInfo,
  33. private readonly exactError?: Error,
  34. ) {
  35. super()
  36. }
  37. override providerInfo(provider: string): LlmProviderInfo {
  38. return { id: provider, name: this.name }
  39. }
  40. override listModels(): Promise<readonly LlmModelInfo[]> {
  41. return this.models instanceof Error
  42. ? Promise.reject(this.models)
  43. : Promise.resolve(this.models)
  44. }
  45. override resolveModel(provider: string, model: string): Promise<LlmResolvedModelInfo> {
  46. if (this.exactError !== undefined) return Promise.reject(this.exactError)
  47. return Promise.resolve({
  48. provider,
  49. id: model,
  50. name: model,
  51. context: { contextWindow: model === 'private-preview' ? 128_000 : 64_000 },
  52. ...this.reasoning === undefined ? {} : { reasoning: this.reasoning },
  53. })
  54. }
  55. override async *stream(_options: GenerateOptions): AsyncIterable<StreamChunk> {
  56. // Catalog tests never enter provider streaming.
  57. }
  58. }
  59. class DeferredCatalogAdapter extends CatalogAdapter {
  60. readonly pending: PromiseWithResolvers<LlmResolvedModelInfo>[] = []
  61. constructor() {
  62. super('Deferred', [
  63. { provider: 'deferred', id: 'lifecycle-model', name: 'Lifecycle model' },
  64. ])
  65. }
  66. override resolveModel(_provider: string, _model: string): Promise<LlmResolvedModelInfo> {
  67. const result = Promise.withResolvers<LlmResolvedModelInfo>()
  68. this.pending.push(result)
  69. return result.promise
  70. }
  71. resolve(index: number, contextWindow: number): void {
  72. const pending = this.pending[index]
  73. if (pending === undefined) throw new Error(`no pending resolution at index ${String(index)}`)
  74. pending.resolve({
  75. provider: 'deferred',
  76. id: 'lifecycle-model',
  77. name: 'Lifecycle model',
  78. context: { contextWindow },
  79. })
  80. }
  81. }
  82. const REASONING: LlmModelReasoningInfo = {
  83. efforts: [
  84. { id: ReasoningEffortId('off'), name: 'Off' },
  85. { id: ReasoningEffortId('high'), name: 'High' },
  86. { id: ReasoningEffortId('max'), name: 'Max' },
  87. ],
  88. defaultEffort: ReasoningEffortId('high'),
  89. }
  90. async function hostContext(
  91. onSessions?: (fiber: Fiber) => void,
  92. onAgents?: (fiber: Fiber) => void,
  93. ): Promise<Context> {
  94. const ctx = new Context()
  95. const sessionsFiber = await ctx.plugin(SessionStore)
  96. onSessions?.(sessionsFiber)
  97. await ctx.plugin(SystemPrompt, { persona: '' })
  98. await ctx.plugin(LlmService)
  99. await ctx.plugin(UserInteractionService)
  100. const agentsFiber = await ctx.plugin(AgentRegistry)
  101. onAgents?.(agentsFiber)
  102. ctx.llm.registerAdapter(['deepseek'], new CatalogAdapter('DeepSeek', [
  103. { provider: 'deepseek', id: 'deepseek-chat', name: 'DeepSeek Chat' },
  104. { provider: 'deepseek', id: 'deepseek-reasoner', name: 'DeepSeek Reasoner', description: 'Reasoning model' },
  105. ], REASONING))
  106. ctx.llm.registerAdapter(['broken'], new CatalogAdapter('Broken Provider', new Error('catalog offline')))
  107. ctx.llm.registerAdapter(['metadata-broken'], new CatalogAdapter('Metadata Broken', [
  108. { provider: 'metadata-broken', id: 'listed', name: 'Listed' },
  109. ], undefined, new Error('reasoning metadata offline')))
  110. ctx.llm.registerAdapter(['empty'], new CatalogAdapter('Empty Provider', []))
  111. ctx.llm.registerAdapter(['duplicate'], new CatalogAdapter('Duplicate Provider', [
  112. { provider: 'duplicate', id: 'same', name: 'Same' },
  113. { provider: 'duplicate', id: 'same', name: 'Same Again' },
  114. ]))
  115. return ctx
  116. }
  117. async function harness(logged?: {
  118. provider: string
  119. model: string
  120. reasoningEffort?: ReasoningEffortId
  121. }): Promise<{
  122. ctx: Context
  123. agent: Agent
  124. sessionId: SessionId
  125. }> {
  126. const ctx = await hostContext()
  127. const session = ctx.sessions.create()
  128. if (logged !== undefined) {
  129. session.append('request/header', { header: { config: logged }, reason: 'initial' })
  130. }
  131. const agent = {
  132. id: session.id,
  133. session,
  134. status: 'running',
  135. ctx,
  136. } as Agent
  137. ctx.agents.register(agent)
  138. return { ctx, agent, sessionId: session.id }
  139. }
  140. function expectValue<T>(response: { result: { ok: true; value: T } | { ok: false } }): T {
  141. if (!response.result.ok) throw new Error('expected successful response')
  142. return response.result.value
  143. }
  144. async function nextMetrics(
  145. iterator: AsyncIterator<RpcRequest<MuxFrame>>,
  146. ): Promise<Extract<MuxFrame, { type: 'session/metrics' }>['metrics']> {
  147. for (;;) {
  148. const next = await iterator.next()
  149. if (next.done) throw new Error('mux ended before a metrics frame')
  150. if (next.value.payload.type === 'session/metrics') return next.value.payload.metrics
  151. }
  152. }
  153. function attachLifecycleSession(
  154. ctx: Context,
  155. sessionId: SessionId,
  156. withMarker = false,
  157. ): { session: Session; detach: () => void } {
  158. const session = ctx.sessions.prepare(sessionId)
  159. session.append('request/header', {
  160. header: { config: { provider: 'deferred', model: 'lifecycle-model' } },
  161. reason: 'initial',
  162. })
  163. if (withMarker) {
  164. session.append('user/message', {
  165. content: [{ type: 'text', text: 'replacement marker' }],
  166. source: { kind: 'plugin', plugin: 'test' },
  167. }, { surfaceOp: 'append' })
  168. }
  169. const detach = ctx.sessions.enter(session)
  170. ctx.sessions.announce(session)
  171. return { session, detach }
  172. }
  173. function attachLifecycleAgent(
  174. ctx: Context,
  175. session: Session,
  176. ): () => void {
  177. const agent = {
  178. id: session.id,
  179. session,
  180. status: 'running',
  181. ctx,
  182. } as Agent
  183. const detach = ctx.agents.enter(agent, undefined)
  184. ctx.agents.announce(agent)
  185. return detach
  186. }
  187. function settleCapacityCompletion(): Promise<void> {
  188. return new Promise<void>((resolve) => { setImmediate(resolve) })
  189. }
  190. function installDeferredAdapter(
  191. ctx: Context,
  192. adapter: DeferredCatalogAdapter,
  193. ): Fiber & PromiseLike<Fiber> {
  194. return ctx.plugin(Object.assign((inner: Context) => {
  195. inner.llm.registerAdapter(['deferred'], adapter)
  196. }, { inject: ['llm'] }))
  197. }
  198. describe('Web session model selection', () => {
  199. it('groups successful providers, isolates failures, and preserves an unlisted current model', async () => {
  200. const { ctx, sessionId } = await harness({
  201. provider: 'deepseek',
  202. model: 'private-preview',
  203. reasoningEffort: ReasoningEffortId('max'),
  204. })
  205. const api = createApiProxy(ctx, { provider: 'deepseek', model: 'deepseek-chat', cwd: '/tmp', workspaceRoot: '/tmp' })
  206. const catalog = expectValue(await api.sessions.models(request({ sessionId })))
  207. expect(catalog.current).toEqual({
  208. provider: 'deepseek',
  209. model: 'private-preview',
  210. reasoningEffort: 'max',
  211. })
  212. expect(catalog.groups).toEqual([{
  213. id: 'deepseek',
  214. name: 'DeepSeek',
  215. models: [
  216. { id: 'deepseek-chat', name: 'DeepSeek Chat', reasoning: REASONING },
  217. {
  218. id: 'deepseek-reasoner',
  219. name: 'DeepSeek Reasoner',
  220. description: 'Reasoning model',
  221. reasoning: REASONING,
  222. },
  223. {
  224. id: 'private-preview',
  225. name: 'private-preview',
  226. unlisted: true,
  227. reasoning: REASONING,
  228. },
  229. ],
  230. }])
  231. expect(catalog.failures).toEqual([
  232. { id: 'broken', name: 'Broken Provider', message: 'catalog offline' },
  233. { id: 'metadata-broken', name: 'Metadata Broken', message: 'reasoning metadata offline' },
  234. {
  235. id: 'duplicate',
  236. name: 'Duplicate Provider',
  237. message: 'adapter returned invalid or duplicate model metadata for provider "duplicate"',
  238. },
  239. ])
  240. await ctx.fiber.dispose()
  241. })
  242. it('accepts an advisory-unlisted model, rejects an unavailable provider, and switches only after the next assembly', async () => {
  243. const { ctx, agent, sessionId } = await harness()
  244. const api = createApiProxy(ctx, { provider: 'deepseek', model: 'deepseek-chat', cwd: '/tmp', workspaceRoot: '/tmp' })
  245. const seed: LlmCallConfig = { provider: 'seed', model: 'seed', temperature: 0.2 }
  246. const signal = new AbortController().signal
  247. expect(expectValue(await api.sessions.models(request({ sessionId }))).current)
  248. .toEqual({ provider: 'deepseek', model: 'deepseek-chat' })
  249. expect((await ctx.systemPrompt.assemble()).variables)
  250. .toMatchObject({ provider: 'deepseek', model: 'deepseek-chat' })
  251. const selected = expectValue(await api.sessions.selectModel(request({
  252. sessionId,
  253. provider: 'deepseek',
  254. model: 'private-preview',
  255. reasoningEffort: 'max',
  256. })))
  257. expect(selected.selected).toEqual({
  258. provider: 'deepseek',
  259. model: 'private-preview',
  260. reasoningEffort: 'max',
  261. })
  262. await expect(agentEvents(ctx, agent).waterfall(
  263. 'agent/request', 1, 0, signal, () => Promise.resolve(seed),
  264. )).resolves.toMatchObject({ provider: 'deepseek', model: 'deepseek-chat' })
  265. expect((await ctx.systemPrompt.assemble()).variables)
  266. .toMatchObject({ provider: 'deepseek', model: 'private-preview' })
  267. await expect(agentEvents(ctx, agent).waterfall(
  268. 'agent/request', 1, 1, signal, () => Promise.resolve(seed),
  269. )).resolves.toMatchObject({
  270. provider: 'deepseek',
  271. model: 'private-preview',
  272. reasoningEffort: 'max',
  273. })
  274. const unsupported = await api.sessions.selectModel(request({
  275. sessionId,
  276. provider: 'deepseek',
  277. model: 'private-preview',
  278. reasoningEffort: 'medium',
  279. }))
  280. expect(unsupported.result).toMatchObject({
  281. ok: false,
  282. error: {
  283. code: 'model-unavailable',
  284. message: 'provider "deepseek" model "private-preview" does not support reasoning effort "medium"',
  285. },
  286. })
  287. const rejected = await api.sessions.selectModel(request({
  288. sessionId,
  289. provider: 'missing',
  290. model: 'model',
  291. }))
  292. expect(rejected.result).toEqual({
  293. ok: false,
  294. error: {
  295. code: 'model-unavailable',
  296. message: 'no adapter registered for provider "missing"',
  297. details: { provider: 'missing', model: 'model' },
  298. },
  299. })
  300. expect(expectValue(await api.sessions.models(request({ sessionId }))).current)
  301. .toEqual({ provider: 'deepseek', model: 'private-preview', reasoningEffort: 'max' })
  302. await ctx.fiber.dispose()
  303. })
  304. it('publishes unknown capacity immediately on selection, then the exact selected route capacity', async () => {
  305. const { ctx, sessionId } = await harness()
  306. const api = createApiProxy(ctx, { provider: 'deepseek', model: 'deepseek-chat', cwd: '/tmp', workspaceRoot: '/tmp' })
  307. expectValue(await api.sessions.models(request({ sessionId })))
  308. const controller = new AbortController()
  309. const iterator = api.events.mux(request({}), controller.signal)[Symbol.asyncIterator]()
  310. expect((await nextMetrics(iterator)).contextWindow).toBeUndefined()
  311. expect((await nextMetrics(iterator)).contextWindow).toBe(64_000)
  312. expectValue(await api.sessions.selectModel(request({
  313. sessionId,
  314. provider: 'deepseek',
  315. model: 'private-preview',
  316. })))
  317. expect((await nextMetrics(iterator)).contextWindow).toBeUndefined()
  318. expect((await nextMetrics(iterator)).contextWindow).toBe(128_000)
  319. controller.abort()
  320. await iterator.return?.()
  321. await ctx.fiber.dispose()
  322. })
  323. it('uses logged capacity without installing Web routing while scheduling foreign metrics', async () => {
  324. const ctx = await hostContext()
  325. const api = createApiProxy(ctx, {
  326. provider: 'deepseek',
  327. model: 'deepseek-chat',
  328. cwd: '/tmp',
  329. workspaceRoot: '/tmp',
  330. })
  331. const controller = new AbortController()
  332. const iterator = api.events.mux(request({}), controller.signal)[Symbol.asyncIterator]()
  333. const initialMetrics = nextMetrics(iterator)
  334. const session = ctx.sessions.create()
  335. expect((await initialMetrics).contextWindow).toBeUndefined()
  336. session.append('request/header', {
  337. header: { config: { provider: 'deepseek', model: 'private-preview' } },
  338. reason: 'change',
  339. })
  340. const foreign = {
  341. id: session.id,
  342. session,
  343. status: 'running',
  344. ctx,
  345. } as Agent
  346. const foreignTarget: AgentLlmTargetRef = {
  347. current: { provider: 'foreign', model: 'foreign-model' },
  348. assembled: undefined,
  349. }
  350. const disposeForeignTarget = installAgentLlmTarget(foreign.ctx, foreignTarget)
  351. const scheduledMetrics = nextMetrics(iterator)
  352. ctx.agents.register(foreign)
  353. expect((await scheduledMetrics).contextWindow).toBeUndefined()
  354. expect((await nextMetrics(iterator)).contextWindow).toBe(128_000)
  355. expect((await ctx.systemPrompt.assemble()).variables)
  356. .toMatchObject({ provider: 'foreign', model: 'foreign-model' })
  357. const seed: LlmCallConfig = { provider: 'seed', model: 'seed', temperature: 0.2 }
  358. const signal = new AbortController().signal
  359. await expect(agentEvents(ctx, foreign).waterfall(
  360. 'agent/request', 1, 0, signal, () => Promise.resolve(seed),
  361. )).resolves.toMatchObject({ provider: 'foreign', model: 'foreign-model' })
  362. disposeForeignTarget()
  363. expect((await ctx.systemPrompt.assemble()).variables).not.toHaveProperty('provider')
  364. await expect(agentEvents(ctx, foreign).waterfall(
  365. 'agent/request', 1, 1, signal, () => Promise.resolve(seed),
  366. )).resolves.toBe(seed)
  367. controller.abort()
  368. await iterator.return?.()
  369. await ctx.fiber.dispose()
  370. })
  371. it('drops capacity completion from a replaced agent that retains the exact session', async () => {
  372. const ctx = await hostContext()
  373. const deferred = new DeferredCatalogAdapter()
  374. ctx.llm.registerAdapter(['deferred'], deferred)
  375. const lifecycle = attachLifecycleSession(ctx, SessionId('capacity-agent-lifecycle'))
  376. const retire = attachLifecycleAgent(ctx, lifecycle.session)
  377. const api = createApiProxy(ctx, { provider: 'deepseek', model: 'deepseek-chat', cwd: '/tmp', workspaceRoot: '/tmp' })
  378. const controller = new AbortController()
  379. const iterator = api.events.mux(request({}), controller.signal)[Symbol.asyncIterator]()
  380. expect((await nextMetrics(iterator)).contextWindow).toBeUndefined()
  381. await vi.waitFor(() => { expect(deferred.pending).toHaveLength(1) })
  382. retire()
  383. const detachLive = attachLifecycleAgent(ctx, lifecycle.session)
  384. expect((await nextMetrics(iterator)).contextWindow).toBeUndefined()
  385. await vi.waitFor(() => { expect(deferred.pending).toHaveLength(2) })
  386. deferred.resolve(0, 64_000)
  387. await settleCapacityCompletion()
  388. deferred.resolve(1, 128_000)
  389. await settleCapacityCompletion()
  390. expect((await nextMetrics(iterator)).contextWindow).toBe(128_000)
  391. controller.abort()
  392. await iterator.return?.()
  393. detachLive()
  394. lifecycle.detach()
  395. await ctx.fiber.dispose()
  396. })
  397. it('drops capacity completion from a replaced session while its old agent remains live', async () => {
  398. const ctx = await hostContext()
  399. const deferred = new DeferredCatalogAdapter()
  400. ctx.llm.registerAdapter(['deferred'], deferred)
  401. const sessionId = SessionId('capacity-session-lifecycle')
  402. const retiredSession = attachLifecycleSession(ctx, sessionId)
  403. const retireAgent = attachLifecycleAgent(ctx, retiredSession.session)
  404. const api = createApiProxy(ctx, { provider: 'deepseek', model: 'deepseek-chat', cwd: '/tmp', workspaceRoot: '/tmp' })
  405. const controller = new AbortController()
  406. const iterator = api.events.mux(request({}), controller.signal)[Symbol.asyncIterator]()
  407. expect((await nextMetrics(iterator)).contextWindow).toBeUndefined()
  408. await vi.waitFor(() => { expect(deferred.pending).toHaveLength(1) })
  409. retiredSession.detach()
  410. const liveSession = attachLifecycleSession(ctx, sessionId, true)
  411. expect((await nextMetrics(iterator)).contextWindow).toBeUndefined()
  412. deferred.resolve(0, 64_000)
  413. await settleCapacityCompletion()
  414. retireAgent()
  415. const detachLiveAgent = attachLifecycleAgent(ctx, liveSession.session)
  416. const scheduled = await nextMetrics(iterator)
  417. expect(scheduled.logRevision).toBe(2)
  418. expect(scheduled.contextWindow).toBeUndefined()
  419. await vi.waitFor(() => { expect(deferred.pending).toHaveLength(2) })
  420. deferred.resolve(1, 128_000)
  421. await settleCapacityCompletion()
  422. expect((await nextMetrics(iterator)).contextWindow).toBe(128_000)
  423. controller.abort()
  424. await iterator.return?.()
  425. detachLiveAgent()
  426. liveSession.detach()
  427. await ctx.fiber.dispose()
  428. })
  429. it('does not project retired agent capacity into replacement session snapshots', async () => {
  430. const ctx = await hostContext()
  431. const deferred = new DeferredCatalogAdapter()
  432. ctx.llm.registerAdapter(['deferred'], deferred)
  433. const sessionId = SessionId('capacity-snapshot-lifecycle')
  434. const retiredSession = attachLifecycleSession(ctx, sessionId)
  435. const retireAgent = attachLifecycleAgent(ctx, retiredSession.session)
  436. const api = createApiProxy(ctx, { provider: 'deepseek', model: 'deepseek-chat', cwd: '/tmp', workspaceRoot: '/tmp' })
  437. const primaryController = new AbortController()
  438. const primary = api.events.mux(request({}), primaryController.signal)[Symbol.asyncIterator]()
  439. expect((await nextMetrics(primary)).contextWindow).toBeUndefined()
  440. await vi.waitFor(() => { expect(deferred.pending).toHaveLength(1) })
  441. deferred.resolve(0, 64_000)
  442. await settleCapacityCompletion()
  443. expect((await nextMetrics(primary)).contextWindow).toBe(64_000)
  444. retiredSession.detach()
  445. const replacement = attachLifecycleSession(ctx, sessionId)
  446. const createdBaseline = await nextMetrics(primary)
  447. replacement.session.append('user/message', {
  448. content: [{ type: 'text', text: 'replacement marker' }],
  449. source: { kind: 'plugin', plugin: 'test' },
  450. }, { surfaceOp: 'append' })
  451. const scheduledFlush = await nextMetrics(primary)
  452. const reconnectController = new AbortController()
  453. const reconnect = api.events.mux(request({}), reconnectController.signal)[Symbol.asyncIterator]()
  454. const reconnectBaseline = await nextMetrics(reconnect)
  455. expect(createdBaseline.logRevision).toBe(1)
  456. for (const metrics of [scheduledFlush, reconnectBaseline]) {
  457. expect(metrics.logRevision).toBe(2)
  458. }
  459. expect({
  460. created: createdBaseline.contextWindow,
  461. scheduled: scheduledFlush.contextWindow,
  462. reconnect: reconnectBaseline.contextWindow,
  463. }).toEqual({ created: undefined, scheduled: undefined, reconnect: undefined })
  464. retireAgent()
  465. const detachReplacementAgent = attachLifecycleAgent(ctx, replacement.session)
  466. expect((await nextMetrics(primary)).contextWindow).toBeUndefined()
  467. await vi.waitFor(() => { expect(deferred.pending).toHaveLength(2) })
  468. deferred.resolve(1, 128_000)
  469. await settleCapacityCompletion()
  470. expect((await nextMetrics(primary)).contextWindow).toBe(128_000)
  471. primaryController.abort()
  472. reconnectController.abort()
  473. await primary.return?.()
  474. await reconnect.return?.()
  475. detachReplacementAgent()
  476. replacement.detach()
  477. await ctx.fiber.dispose()
  478. })
  479. it('refreshes same-route capacity after adapter owner replacement', async () => {
  480. const ctx = await hostContext()
  481. const retiredAdapter = new DeferredCatalogAdapter()
  482. const retiredFiber = await installDeferredAdapter(ctx, retiredAdapter)
  483. const lifecycle = attachLifecycleSession(ctx, SessionId('capacity-adapter-lifecycle'))
  484. const detachAgent = attachLifecycleAgent(ctx, lifecycle.session)
  485. const api = createApiProxy(ctx, { provider: 'deepseek', model: 'deepseek-chat', cwd: '/tmp', workspaceRoot: '/tmp' })
  486. const controller = new AbortController()
  487. const iterator = api.events.mux(request({}), controller.signal)[Symbol.asyncIterator]()
  488. expect((await nextMetrics(iterator)).contextWindow).toBeUndefined()
  489. await vi.waitFor(() => { expect(retiredAdapter.pending).toHaveLength(1) })
  490. retiredAdapter.resolve(0, 64_000)
  491. await settleCapacityCompletion()
  492. expect((await nextMetrics(iterator)).contextWindow).toBe(64_000)
  493. await retiredFiber.dispose()
  494. expect((await nextMetrics(iterator)).contextWindow).toBeUndefined()
  495. const replacementAdapter = new DeferredCatalogAdapter()
  496. const replacementFiber = await installDeferredAdapter(ctx, replacementAdapter)
  497. expect((await nextMetrics(iterator)).contextWindow).toBeUndefined()
  498. await vi.waitFor(() => { expect(replacementAdapter.pending).toHaveLength(1) })
  499. replacementAdapter.resolve(0, 128_000)
  500. await settleCapacityCompletion()
  501. expect((await nextMetrics(iterator)).contextWindow).toBe(128_000)
  502. expect(lifecycle.session.requestHeader()?.config).toMatchObject({
  503. provider: 'deferred',
  504. model: 'lifecycle-model',
  505. })
  506. controller.abort()
  507. await iterator.return?.()
  508. detachAgent()
  509. lifecycle.detach()
  510. await replacementFiber.dispose()
  511. await ctx.fiber.dispose()
  512. })
  513. it('does not read the sessions service after its disposal status', async () => {
  514. let sessionsFiber: Fiber | undefined
  515. const ctx = await hostContext((fiber) => { sessionsFiber = fiber })
  516. if (sessionsFiber === undefined) throw new Error('sessions fiber missing')
  517. createApiProxy(ctx, { provider: 'deepseek', model: 'deepseek-chat', cwd: '/tmp', workspaceRoot: '/tmp' })
  518. const sessions = ctx.get('sessions')
  519. if (sessions === undefined) throw new Error('sessions service missing')
  520. const list = vi.spyOn(sessions, 'list').mockImplementation(() => {
  521. throw new Error('disposed sessions service read')
  522. })
  523. await expect(sessionsFiber.dispose()).resolves.toBeUndefined()
  524. expect(list).not.toHaveBeenCalled()
  525. await ctx.fiber.dispose()
  526. })
  527. it('invalidates pending capacity before SessionStore teardown can reach its callback', async () => {
  528. let sessionsFiber: Fiber | undefined
  529. const ctx = await hostContext((fiber) => { sessionsFiber = fiber })
  530. if (sessionsFiber === undefined) throw new Error('sessions fiber missing')
  531. const deferred = new DeferredCatalogAdapter()
  532. ctx.llm.registerAdapter(['deferred'], deferred)
  533. const lifecycle = attachLifecycleSession(ctx, SessionId('capacity-session-store-teardown'))
  534. const detachAgent = attachLifecycleAgent(ctx, lifecycle.session)
  535. const api = createApiProxy(ctx, { provider: 'deepseek', model: 'deepseek-chat', cwd: '/tmp', workspaceRoot: '/tmp' })
  536. const controller = new AbortController()
  537. const iterator = api.events.mux(request({}), controller.signal)[Symbol.asyncIterator]()
  538. expect((await nextMetrics(iterator)).contextWindow).toBeUndefined()
  539. await vi.waitFor(() => { expect(deferred.pending).toHaveLength(1) })
  540. const agents = ctx.get('agents')
  541. if (agents === undefined) throw new Error('agent registry missing')
  542. expect(agents.get(lifecycle.session.id)).toBeDefined()
  543. const getAgent = vi.spyOn(agents, 'get')
  544. await sessionsFiber.dispose()
  545. getAgent.mockClear()
  546. deferred.resolve(0, 64_000)
  547. await settleCapacityCompletion()
  548. expect(getAgent).not.toHaveBeenCalled()
  549. const pendingFrame = iterator.next()
  550. const outcome = await Promise.race([
  551. pendingFrame.then(() => 'frame' as const),
  552. new Promise<'idle'>((resolve) => { setImmediate(() => { resolve('idle') }) }),
  553. ])
  554. expect(outcome).toBe('idle')
  555. controller.abort()
  556. await expect(pendingFrame).resolves.toMatchObject({ done: true })
  557. await iterator.return?.()
  558. detachAgent()
  559. lifecycle.detach()
  560. await ctx.fiber.dispose()
  561. })
  562. it('publishes unknown metrics after AgentRegistry terminal disposal with mux active', async () => {
  563. let agentsFiber: Fiber | undefined
  564. const ctx = await hostContext(undefined, (fiber) => { agentsFiber = fiber })
  565. if (agentsFiber === undefined) throw new Error('agent registry fiber missing')
  566. const deferred = new DeferredCatalogAdapter()
  567. ctx.llm.registerAdapter(['deferred'], deferred)
  568. const lifecycle = attachLifecycleSession(ctx, SessionId('capacity-agent-registry-disposed'))
  569. const detachAgent = attachLifecycleAgent(ctx, lifecycle.session)
  570. const api = createApiProxy(ctx, { provider: 'deepseek', model: 'deepseek-chat', cwd: '/tmp', workspaceRoot: '/tmp' })
  571. const controller = new AbortController()
  572. const iterator = api.events.mux(request({}), controller.signal)[Symbol.asyncIterator]()
  573. expect((await nextMetrics(iterator)).contextWindow).toBeUndefined()
  574. await vi.waitFor(() => { expect(deferred.pending).toHaveLength(1) })
  575. deferred.resolve(0, 64_000)
  576. await settleCapacityCompletion()
  577. expect((await nextMetrics(iterator)).contextWindow).toBe(64_000)
  578. await agentsFiber.dispose()
  579. const refresh = nextMetrics(iterator).then(
  580. metrics => ({ kind: 'metrics' as const, metrics }),
  581. () => ({ kind: 'error' as const }),
  582. )
  583. const outcome = await Promise.race([
  584. refresh,
  585. new Promise<{ kind: 'idle' }>((resolve) => {
  586. setImmediate(() => { resolve({ kind: 'idle' }) })
  587. }),
  588. ])
  589. controller.abort()
  590. await refresh
  591. await iterator.return?.()
  592. detachAgent()
  593. lifecycle.detach()
  594. await ctx.fiber.dispose()
  595. expect(outcome.kind).toBe('metrics')
  596. if (outcome.kind === 'metrics') {
  597. expect(outcome.metrics.contextWindow).toBeUndefined()
  598. expect(outcome.metrics.logRevision).toBe(1)
  599. }
  600. })
  601. it.each(['agents', 'sessions'] as const)(
  602. 'drops capacity completion while %s is unavailable during unload',
  603. async (serviceName) => {
  604. let sessionsFiber: Fiber | undefined
  605. let agentsFiber: Fiber | undefined
  606. const ctx = await hostContext(
  607. (fiber) => { sessionsFiber = fiber },
  608. (fiber) => { agentsFiber = fiber },
  609. )
  610. const heldFiber = serviceName === 'sessions' ? sessionsFiber : agentsFiber
  611. if (heldFiber === undefined) throw new Error(`${serviceName} fiber missing`)
  612. const unloadStarted = Promise.withResolvers<undefined>()
  613. const releaseUnload = Promise.withResolvers<undefined>()
  614. heldFiber.ctx.effect(() => () => {
  615. unloadStarted.resolve(undefined)
  616. return releaseUnload.promise
  617. }, `test: hold ${serviceName} unload`)
  618. const deferred = new DeferredCatalogAdapter()
  619. ctx.llm.registerAdapter(['deferred'], deferred)
  620. const lifecycle = attachLifecycleSession(ctx, SessionId(`capacity-${serviceName}-unloading`))
  621. const detachAgent = attachLifecycleAgent(ctx, lifecycle.session)
  622. const api = createApiProxy(ctx, { provider: 'deepseek', model: 'deepseek-chat', cwd: '/tmp', workspaceRoot: '/tmp' })
  623. const controller = new AbortController()
  624. const iterator = api.events.mux(request({}), controller.signal)[Symbol.asyncIterator]()
  625. expect((await nextMetrics(iterator)).contextWindow).toBeUndefined()
  626. await vi.waitFor(() => { expect(deferred.pending).toHaveLength(1) })
  627. const agents = ctx.get('agents')
  628. if (agents === undefined) throw new Error('agent registry missing')
  629. const sessions = ctx.get('sessions')
  630. if (sessions === undefined) throw new Error('sessions service missing')
  631. const getAgent = vi.spyOn(agents, 'get')
  632. const getSession = vi.spyOn(sessions, 'get')
  633. const disposing = heldFiber.dispose()
  634. await unloadStarted.promise
  635. await vi.waitFor(() => { expect(ctx.get(serviceName)).toBeUndefined() })
  636. expect(heldFiber.state).toBe(FiberState.UNLOADING)
  637. getAgent.mockClear()
  638. getSession.mockClear()
  639. deferred.resolve(0, 64_000)
  640. await settleCapacityCompletion()
  641. const agentReads = getAgent.mock.calls.length
  642. const sessionReads = getSession.mock.calls.length
  643. const pendingFrame = iterator.next()
  644. const outcome = await Promise.race([
  645. pendingFrame.then(() => 'frame' as const),
  646. new Promise<'idle'>((resolve) => { setImmediate(() => { resolve('idle') }) }),
  647. ])
  648. controller.abort()
  649. await expect(pendingFrame).resolves.toMatchObject({ done: true })
  650. await iterator.return?.()
  651. releaseUnload.resolve(undefined)
  652. await disposing
  653. detachAgent()
  654. lifecycle.detach()
  655. await ctx.fiber.dispose()
  656. expect(agentReads).toBe(serviceName === 'sessions' ? 1 : 0)
  657. expect(sessionReads).toBe(0)
  658. expect(outcome).toBe('idle')
  659. },
  660. )
  661. })