/** * `HarnessSdkServer`: the JSON-RPC method surface the `dsh-jsonrpc` plugin * serves to out-of-process SDK clients (e.g. the Python `deepseek_harness` * package). Requests: `initialize` → `session/prompt`* → `shutdown`. * Notifications pushed to the host: `session.event` (every durable session * event, verbatim), `session.finished` (per prompt turn settle), * `subagent.started` / `subagent.finished` (child-session lineage and run * outcomes). The server owns only the SDK-facing session map — the harness * itself is the context the plugin mounts in; plugins, persistence, and * the LLM adapter set all come from the external `cordis.yml`. * * @module @deepseek-ai/dsh-jsonrpc/server */ import type { Context } from 'cordis' import { resolve } from 'node:path' import type { ContentBlock } from '@deepseek-ai/dsh-llm' import type { Agent, AgentHandle } from '@deepseek-ai/dsh-agent' import { carrierKeyOf, type Scoped } from '@deepseek-ai/dsh-scope' import { SessionId, type TurnEndReason } from '@deepseek-ai/dsh-session' import type SubagentService from '@deepseek-ai/dsh-subagent' import type { SubagentRunEndInfo, SubagentRunInfo } from '@deepseek-ai/dsh-subagent' import * as LlmDeepSeek from '@deepseek-ai/dsh-llm-deepseek' import type { JsonRpcTransportPeer } from './transport.ts' /** Parameters of the `initialize` request (once per process, before any prompt). */ export interface InitializeParams { /** Working directory recorded on every SDK-created session's header. */ cwd: string /** Model name every SDK-created agent runs on (see {@link HarnessSdkServer.initialize} for adapter fallback). */ model: string } /** Result of the `initialize` request: the server's identity for the SDK handshake. */ export interface InitializeResult { /** Wire-stable server identity (`deepseek-harness-sdk-runtime`) and version. */ serverInfo: { name: string; version: string } } /** * Parameters of a `session/prompt` request: one user turn on one SDK session, * with at most one in flight per session. */ export interface SessionPromptParams { /** The SDK-side session id; an unknown id lazily creates the agent+session pair. */ sessionId: string /** The prompt content blocks, sent verbatim as the user message. */ contentBlocks: ContentBlock[] } /** Result of a `session/prompt` request: the prompt ran to turn settle (outcome rides on `session.finished`). */ export interface SessionPromptResult { /** Always `true`; the turn outcome is the paired `session.finished` notification. */ accepted: true } interface SessionRecord { handle: AgentHandle lastTurnEnd: TurnEndReason | undefined activePrompt: boolean } /** Recover the delegating parent carried by every service-owned subagent lifecycle event. */ function subagentParentOf(carrier: Scoped): Agent { // SubagentService emits this lifecycle pair only through scopeTarget(this, parent). return carrierKeyOf(carrier) as Agent } /** Whether the live id names a local child related to this exact delegating parent. */ function isLocalChild(ctx: Context, id: SessionId, parent: Agent): boolean { const child = ctx.agents.get(id) return child !== undefined && ( ctx.agents.isOwnedBy(id, parent) || child.session.header.parentSession === parent.session.id ) } /** * The SDK server over a booted harness context. Constructing it subscribes to * session and subagent lifecycle events, forwarding durable session * events and SDK-facing completion notifications while retaining local-run * identity across child disposal. The subscriptions live until * {@link shutdown}. One instance serves one transport peer for the process * lifetime — there is no re-`initialize`. */ export class HarnessSdkServer { private cwd = process.cwd() private model = 'deepseek' private llmFiber: { dispose(): Promise } | undefined private readonly sessions = new Map() private readonly sessionCreations = new Map>() private readonly localRuns = new Map>>() private readonly disposers: (() => void)[] = [] private shutdownTask: Promise> | undefined private shuttingDown = false constructor( private readonly ctx: Context, private readonly transport: JsonRpcTransportPeer, ) { this.disposers.push(ctx.on('session/event', (session, event) => { if (event.type === 'turn/end') { const rec = this.sessions.get(String(session.id)) if (rec) rec.lastTurnEnd = event.data.reason } this.transport.notify('session.event', { sessionId: String(session.id), event }) })) this.disposers.push(ctx.on('session/created', (session) => { const parentSession = session.header.parentSession if (parentSession === undefined) return this.transport.notify('subagent.started', { parentSessionId: String(parentSession), childSessionId: String(session.id), }) })) // In-process providers publish the child before start. Count starts related // by exact runtime ownership or durable parent lineage so provider-owned // roots remain local, completions survive child disposal, and reused ids // need no settlement-order assumption. const localRuns = this.localRuns this.disposers.push(ctx.on('subagent/start', function (this: Scoped, info: SubagentRunInfo) { const parent = subagentParentOf(this) if (!isLocalChild(ctx, info.id, parent)) return const providerRuns = localRuns.get(info.provider) ?? new Map>() const parentRuns = providerRuns.get(info.id) ?? new Map() parentRuns.set(parent, (parentRuns.get(parent) ?? 0) + 1) providerRuns.set(info.id, parentRuns) localRuns.set(info.provider, providerRuns) })) this.disposers.push(ctx.on('subagent/end', function (this: Scoped, info: SubagentRunEndInfo) { const parent = subagentParentOf(this) const providerRuns = localRuns.get(info.provider) const parentRuns = providerRuns?.get(info.id) const pendingCount = parentRuns?.get(parent) if (pendingCount !== undefined) { if (pendingCount === 1) parentRuns?.delete(parent) else parentRuns?.set(parent, pendingCount - 1) if (parentRuns?.size === 0) providerRuns?.delete(info.id) if (providerRuns?.size === 0) localRuns.delete(info.provider) } // This protocol reports LOCAL child sessions. A lineage-bearing child // has the session/created-driven start notification above. A remote run // has neither a cached local start nor a live child related to this // parent; an unrelated local agent with the same id never makes it local. if (pendingCount === undefined && !isLocalChild(ctx, info.id, parent)) return transport.notify('subagent.finished', { provider: info.provider, agentId: String(info.id), parentSessionId: String(parent.session.id), childSessionId: String(info.id), status: info.stopReason === 'completed' ? 'ok' : 'error', stopReason: info.stopReason, ...(info.lastAssistantMessage === undefined ? {} : { lastAssistantMessage: info.lastAssistantMessage }), }) })) } /** * Handle `initialize`: record the SDK deployment facts (cwd, model) and, when * no registered adapter serves `params.model`, mount the DeepSeek adapter for * it (credentials from `$DEEPSEEK_API_KEY`/`$DEEPSEEK_BASE_URL`) — a config * that already registered an adapter for the model wins. * @param params - the SDK handshake parameters. * @returns the server identity for the handshake. */ async initialize(params: InitializeParams): Promise { this.cwd = resolve(params.cwd) this.model = params.model if (!this.llmFiber && !this.hasAdapterFor(this.model)) { this.llmFiber = await this.ctx.plugin(LlmDeepSeek, { models: [this.model] }) } return { serverInfo: { name: 'deepseek-harness-sdk-runtime', version: '0.0.1' } } } /** * Handle `session/prompt`: get-or-create the session's agent, send the * content as the user message, await turn settle (quiescence), then notify * `session.finished` with the settled turn's outcome. A session accepts at * most one prompt at a time; an overlapping request fails immediately while * other sessions remain independent. * @param params - the target session id and prompt content. * @returns `{ accepted: true }` after the turn settled. */ async prompt(params: SessionPromptParams): Promise { const rec = await this.getOrCreateSession(params.sessionId) if (rec.activePrompt) throw new Error(`session already has an active prompt: ${params.sessionId}`) rec.activePrompt = true try { rec.lastTurnEnd = undefined rec.handle.agent.send(params.contentBlocks) await rec.handle.agent.whenIdle() const status = this.finishedStatus(rec.lastTurnEnd) this.transport.notify('session.finished', { sessionId: params.sessionId, status, reason: rec.lastTurnEnd, }) return { accepted: true } } finally { rec.activePrompt = false } } /** * Handle `shutdown`: dispose every SDK-created agent handle (awaiting loop * quiescence), unmount the adapter fiber this server mounted (if any), and * detach the event subscriptions. The CONTEXT stays up — the bin disposes it * as part of process exit. * @returns an empty object (the JSON-RPC result). */ shutdown(): Promise> { this.shutdownTask ??= this.performShutdown() return this.shutdownTask } private async performShutdown(): Promise> { this.shuttingDown = true const pendingCreations = [...this.sessionCreations.values()] await Promise.allSettled(pendingCreations) this.sessionCreations.clear() const records = [...this.sessions.values()] this.sessions.clear() this.localRuns.clear() const failures: unknown[] = [] while (this.disposers.length > 0) { try { this.disposers.pop()?.() } catch (error) { failures.push(error) } } const teardownResults = await Promise.allSettled([ ...records.map(rec => Promise.resolve().then(() => rec.handle.dispose())), ...(this.llmFiber === undefined ? [] : [Promise.resolve().then(() => this.llmFiber?.dispose())]), ]) this.llmFiber = undefined failures.push(...teardownResults .filter((result): result is PromiseRejectedResult => result.status === 'rejected') .map(result => result.reason as unknown)) if (failures.length === 1) throw failures[0] if (failures.length > 1) throw new AggregateError(failures, 'SDK server teardown failed') return {} } /** * Dispatch one incoming JSON-RPC request to its typed handler. Throws (→ a * JSON-RPC error response) on an unknown method. * @param method - the JSON-RPC method name. * @param params - the raw params object from the wire. * @returns the handler's result, to be serialized as the response. */ async handleRequest(method: string, params: Record | undefined): Promise { switch (method) { case 'initialize': return this.initialize(params as unknown as InitializeParams) case 'session/prompt': return this.prompt(params as unknown as SessionPromptParams) case 'shutdown': return this.shutdown() default: throw new Error(`unknown DeepSeek Harness SDK runtime method: ${method}`) } } private async getOrCreateSession(sessionId: string): Promise { if (this.shuttingDown) throw new Error('SDK server is shutting down') const existing = this.sessions.get(sessionId) if (existing) return existing const pending = this.sessionCreations.get(sessionId) if (pending) return pending const creation = this.createSession(sessionId) this.sessionCreations.set(sessionId, creation) void creation.then( () => { this.sessionCreations.delete(sessionId) }, () => { this.sessionCreations.delete(sessionId) }, ) return creation } private async createSession(sessionId: string): Promise { const handle = await this.ctx.agents.create({ sessionId: SessionId(sessionId), meta: { cwd: this.cwd }, agentOptions: { model: this.model }, }) const rec: SessionRecord = { handle, lastTurnEnd: undefined, activePrompt: false } this.sessions.set(sessionId, rec) return rec } private finishedStatus(reason: TurnEndReason | undefined): 'ok' | 'error' { if (!reason) return 'error' return reason.kind === 'completed' ? 'ok' : 'error' } private hasAdapterFor(model: string): boolean { return this.ctx.get('llm')?.models().includes(model) ?? false } }