|
|
@@ -6,6 +6,7 @@ import { tmpdir } from 'node:os'
|
|
|
import { PassThrough, Writable } from 'node:stream'
|
|
|
import { afterEach, describe, expect, it, vi } from 'vitest'
|
|
|
import { Context } from '@deepseek-ai/cordis'
|
|
|
+import Loader from '@deepseek-ai/cordis-plugin-loader'
|
|
|
import * as agentCore from '@deepseek-ai/dsh-agent-spine-demo'
|
|
|
import JsonlSessionPersistence from '@deepseek-ai/dsh-session-persistence-jsonl'
|
|
|
import * as jsonrpc from '../src/index.ts'
|
|
|
@@ -57,12 +58,17 @@ async function settle(): Promise<void> {
|
|
|
/** Mount the real plugin on a minimal harness with in-memory stdio and exit. */
|
|
|
async function mountPlugin(
|
|
|
storageDir: string,
|
|
|
- options: { writeDelayMs?: number; failFlush?: boolean } = {},
|
|
|
+ options: {
|
|
|
+ writeDelayMs?: number
|
|
|
+ failFlush?: boolean
|
|
|
+ beforeServer?: (ctx: Context) => Promise<void> | void
|
|
|
+ } = {},
|
|
|
): Promise<ApplyHarness> {
|
|
|
const ctx = new Context()
|
|
|
await ctx.plugin(agentCore, { workspaceContext: false })
|
|
|
await ctx.plugin(JsonlSessionPersistence, { root: storageDir })
|
|
|
await new Promise(resolve => setTimeout(resolve, 50))
|
|
|
+ await options.beforeServer?.(ctx)
|
|
|
|
|
|
const input = new PassThrough()
|
|
|
const events: WireEvent[] = []
|
|
|
@@ -170,6 +176,58 @@ describe('dsh-sdk-jsonrpc-server plugin apply', () => {
|
|
|
}
|
|
|
})
|
|
|
|
|
|
+ it('does not answer initialize until async sibling Loader entries settle', async () => {
|
|
|
+ const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-apply-readiness-'))
|
|
|
+ vi.stubEnv('DEEPSEEK_API_KEY', 'test-key')
|
|
|
+ let markStarted!: () => void
|
|
|
+ let release!: () => void
|
|
|
+ const started = new Promise<void>((resolve) => { markStarted = resolve })
|
|
|
+ const ready = new Promise<void>((resolve) => { release = resolve })
|
|
|
+ let delayedEntry: Promise<string> | undefined
|
|
|
+ const harness = await mountPlugin(storageDir, {
|
|
|
+ beforeServer: async (ctx) => {
|
|
|
+ await ctx.plugin(Loader)
|
|
|
+ ctx.loader.builtins['delayed-readiness'] = {
|
|
|
+ async apply() {
|
|
|
+ markStarted()
|
|
|
+ await ready
|
|
|
+ },
|
|
|
+ }
|
|
|
+ delayedEntry = ctx.loader.create({ name: 'cordis:delayed-readiness' })
|
|
|
+ await started
|
|
|
+ },
|
|
|
+ })
|
|
|
+ try {
|
|
|
+ const initialize = {
|
|
|
+ jsonrpc: '2.0',
|
|
|
+ id: 'init-delayed',
|
|
|
+ method: 'initialize',
|
|
|
+ params: { cwd: storageDir, provider: 'deepseek-official', model: 'apply-model' },
|
|
|
+ }
|
|
|
+ const probe = { jsonrpc: '2.0', id: 'probe-during-delay', method: 'nope/unknown' }
|
|
|
+ harness.sendRaw(`${JSON.stringify(initialize)}\n${JSON.stringify(probe)}\n`)
|
|
|
+
|
|
|
+ // The transport processes independent requests concurrently. Receiving
|
|
|
+ // this later probe proves the preceding initialize handler has reached
|
|
|
+ // its Loader wait, without relying on a scheduler delay.
|
|
|
+ await harness.waitForFrame(frame => frame.id === 'probe-during-delay', 'probe while initialize waits')
|
|
|
+ expect(harness.frames().some(frame => frame.id === 'init-delayed')).toBe(false)
|
|
|
+
|
|
|
+ release()
|
|
|
+ await delayedEntry
|
|
|
+ const response = await harness.waitForFrame(frame => frame.id === 'init-delayed', 'initialize response after Loader settlement')
|
|
|
+ expect(response).toMatchObject({
|
|
|
+ id: 'init-delayed',
|
|
|
+ result: { serverInfo: { name: 'deepseek-harness-sdk-runtime' } },
|
|
|
+ })
|
|
|
+ } finally {
|
|
|
+ release()
|
|
|
+ await Promise.allSettled(delayedEntry === undefined ? [] : [delayedEntry])
|
|
|
+ await harness.dispose()
|
|
|
+ await rm(storageDir, { recursive: true, force: true })
|
|
|
+ }
|
|
|
+ })
|
|
|
+
|
|
|
it('drives a session/prompt turn end-to-end and forwards session notifications as output frames', async () => {
|
|
|
const storageDir = await mkdtemp(join(tmpdir(), 'dsh-jsonrpc-apply-prompt-'))
|
|
|
const llmServer = await mockCompletionServer()
|