1
/**2
* JSON-RPC methods and notifications for out-of-process harness SDKs.3
* The surrounding context owns plugins, persistence, and configured adapters.4
*5
* @module @deepseek-ai/dsh-sdk-jsonrpc-server/server6
*/8
import type { Context } from '@deepseek-ai/cordis'9
import { resolve } from 'node:path'10
import { brandString } from '@deepseek-ai/dsh-brand'11
import type { Agent, AgentHandle } from '@deepseek-ai/dsh-agent'12
import { admitEncodedImages, type EncodedImageAttachment, type ImageAttachmentRef } from '@deepseek-ai/dsh-attachment'13
import { createUserMessage, ReasoningEffortId, type ContentBlock, type LlmRuntime } from '@deepseek-ai/dsh-llm'14
import { carrierKeyOf, type Scoped } from '@deepseek-ai/dsh-scope'15
import type { SessionId } from '@deepseek-ai/dsh-session'16
import type SubagentRuntime from '@deepseek-ai/dsh-subagent'17
import type { SubagentRunEndInfo } from '@deepseek-ai/dsh-subagent'18
import * as LlmDeepSeek from '@deepseek-ai/dsh-llm-deepseek-api-key'19
import type {20
InitializeParams,21
InitializeResult,22
JsonRpcTransportPeer,23
SessionEventNotification,24
SessionPromptParams,25
SessionPromptResult,26
SdkEncodedImageBlock,27
SubagentFinishedNotification,28
SubagentStartedNotification,29
} from '@deepseek-ai/dsh-sdk-protocol'31
interface SessionRecord {32
handle: AgentHandle33
}35
function encodedImage(block: SessionPromptParams['contentBlocks'][number]): block is SdkEncodedImageBlock {36
return block.type === 'image' && 'data' in block37
}39
async function durablePromptContent(ctx: Context, blocks: SessionPromptParams['contentBlocks']): Promise<ContentBlock[]> {40
const images = blocks.filter(encodedImage)41
if (images.length === 0) return blocks as ContentBlock[]42
const attachments = ctx.get('attachments')43
if (attachments === undefined) throw new Error('SDK image prompt requires an attachment store')44
const refs = await admitEncodedImages(attachments, images.map((image): EncodedImageAttachment => ({45
data: image.data,46
mediaType: image.mimeType,47
})))48
let next = 049
return blocks.map(block => encodedImage(block)50
? { type: 'image', attachment: refs[next++] as ImageAttachmentRef }51
: block)52
}54
/** Recover the delegating parent from the service-owned scoped carrier. */55
function subagentParentOf(carrier: Scoped<SubagentRuntime>): Agent {56
return carrierKeyOf(carrier) as Agent57
}59
/** Deployment-specific status mapping for SDK turn and subagent outcomes. */60
export interface HarnessSdkJsonRpcServerOptions {61
/** Report max-token termination as an accepted result instead of an infrastructure error. */62
maxTokensAsSuccess?: boolean63
}65
function successStatus(reason: string, options: HarnessSdkJsonRpcServerOptions): 'ok' | 'error' {66
if (reason === 'completed') return 'ok'67
return reason === 'max-tokens' && options.maxTokensAsSuccess === true ? 'ok' : 'error'68
}70
/**71
* SDK server over one booted harness context and transport peer. Construction72
* subscribes to session, agent, and subagent lifecycle events until shutdown;73
* reinitialization is unsupported.74
*/75
export class HarnessSdkJsonRpcServer {76
private cwd = process.cwd()77
private provider = 'deepseek-official'78
private model = 'deepseek-official'79
private reasoningEffort: ReturnType<typeof ReasoningEffortId> | undefined80
private maxTokens: number | undefined81
private llmFiber: { dispose(): Promise<void> } | undefined82
private readonly sessions = new Map<string, SessionRecord>()83
private readonly sessionCreations = new Map<string, Promise<SessionRecord>>()84
private readonly disposers: (() => void)[] = []85
private shutdownTask: Promise<Record<string, never>> | undefined86
private shuttingDown = false87
private initialized = false89
constructor(90
private readonly ctx: Context,91
private readonly transport: JsonRpcTransportPeer,92
private readonly options: HarnessSdkJsonRpcServerOptions = {},93
) {94
const serverOptions = this.options95
this.disposers.push(ctx.on('session/event', (session, event) => {96
const payload: SessionEventNotification = { sessionId: String(session.id), event }97
this.transport.notify('session.event', payload)98
}))99
this.disposers.push(ctx.on('agent/status', ({ agent, status }) => {100
this.transport.notify('session.status', { sessionId: String(agent.session.id), status })101
}))102
this.disposers.push(ctx.on('session/created', (session) => {103
const parentSession = session.header.parentSession104
if (parentSession === undefined) return105
const payload: SubagentStartedNotification = {106
parentSessionId: String(parentSession),107
childSessionId: String(session.id),108
}109
this.transport.notify('subagent.started', payload)110
}))111
this.disposers.push(ctx.on('subagent/end', function (this: Scoped<SubagentRuntime>, info: SubagentRunEndInfo) {112
const parent = subagentParentOf(this)113
// This protocol reports only in-process child sessions. The service114
// snapshots the provider name and local flag through child disposal;115
// matching ids or parent lineage alone never establishes locality.116
if (!info.local) return117
const payload: SubagentFinishedNotification = {118
provider: info.provider,119
agentId: String(info.id),120
parentSessionId: String(parent.session.id),121
childSessionId: String(info.id),122
status: successStatus(info.stopReason, serverOptions),123
stopReason: info.stopReason,124
...(info.lastAssistantMessage === undefined125
? {}126
: { lastAssistantMessage: [...info.lastAssistantMessage] }),127
}128
transport.notify('subagent.finished', payload)129
}))130
}132
/**133
* Validate and configure the SDK route, mounting the DeepSeek fallback only when unowned.134
* @param params - SDK handshake parameters.135
* @returns server identity for the handshake.136
*/137
async initialize(params: InitializeParams): Promise<InitializeResult> {138
if (params.reasoningEffort !== undefined139
&& (typeof params.reasoningEffort !== 'string' || params.reasoningEffort.length === 0)) {140
throw new TypeError('initialize reasoningEffort must be a non-empty string')141
}142
if (params.maxTokens !== undefined143
&& (!Number.isSafeInteger(params.maxTokens) || params.maxTokens <= 0)) {144
throw new TypeError('initialize maxTokens must be a positive safe integer')145
}146
const cwd = resolve(params.cwd)147
const provider = params.provider148
const model = params.model149
const reasoningEffort = params.reasoningEffort === undefined150
? undefined151
: ReasoningEffortId(params.reasoningEffort)152
if (!this.hasAdapterFor(provider)) {153
if (provider !== 'deepseek-official') throw new Error(`no adapter registered for provider "${provider}"`)154
this.llmFiber = await this.ctx.plugin(LlmDeepSeek)155
}156
// Adapter presence was read from this service above; a successful fallback mount also requires it.157
const llm = this.ctx.get('llm') as LlmRuntime158
await llm.resolveCallConfig({159
provider,160
model,161
...reasoningEffort === undefined ? {} : { reasoningEffort },162
...params.maxTokens === undefined ? {} : { maxTokens: params.maxTokens },163
})164
this.cwd = cwd165
this.provider = provider166
this.model = model167
this.reasoningEffort = reasoningEffort168
this.maxTokens = params.maxTokens169
this.initialized = true170
return { serverInfo: { name: 'deepseek-harness-sdk-runtime', version: '0.0.1' } }171
}173
/**174
* Queue one identified prompt without assigning later activity to it.175
* @param params - target session and user content.176
* @returns the durable message identity.177
*/178
async prompt(params: SessionPromptParams): Promise<SessionPromptResult> {179
if (!this.initialized) throw new Error('SDK server is not initialized')180
const rec = await this.getOrCreateSession(params.sessionId)181
// An agent-loop-only reload disposes the loop's agents while this record182
// survives; a retained agent accepts followup() silently, so validate the183
// record against the live registry before delivery.184
this.assertLiveAgent(rec, params.sessionId)185
const content = await durablePromptContent(this.ctx, params.contentBlocks)186
// Attachment admission crosses an async boundary where shutdown or an187
// agent-loop reload may detach the retained handle.188
this.assertLiveAgent(rec, params.sessionId)189
const message = createUserMessage({190
content,191
source: { kind: 'user' },192
})193
rec.handle.agent.followup(message)194
return { messageId: message.id }195
}197
private assertLiveAgent(rec: SessionRecord, sessionId: string): void {198
if (this.ctx.agents.get(rec.handle.agent.id) !== rec.handle.agent) {199
throw new Error(`session agent was disposed outside the server: ${sessionId}`)200
}201
}203
/**204
* Dispose server-owned agents, adapter, and subscriptions to quiescence.205
* The surrounding context remains running.206
* @returns empty JSON-RPC result.207
*/208
shutdown(): Promise<Record<string, never>> {209
this.shutdownTask ??= this.performShutdown()210
return this.shutdownTask211
}213
private async performShutdown(): Promise<Record<string, never>> {214
this.shuttingDown = true215
const pendingCreations = [...this.sessionCreations.values()]216
await Promise.allSettled(pendingCreations)217
this.sessionCreations.clear()218
const records = [...this.sessions.values()]219
this.sessions.clear()220
const failures: unknown[] = []221
while (this.disposers.length > 0) {222
try {223
this.disposers.pop()?.()224
} catch (error) {225
failures.push(error)226
}227
}228
const teardownResults = await Promise.allSettled([229
...records.map(rec => Promise.resolve().then(() => rec.handle.dispose())),230
...(this.llmFiber === undefined ? [] : [Promise.resolve().then(() => this.llmFiber?.dispose())]),231
])232
this.llmFiber = undefined233
failures.push(...teardownResults234
.filter((result): result is PromiseRejectedResult => result.status === 'rejected')235
.map(result => result.reason as unknown))236
if (failures.length === 1) throw failures[0]237
if (failures.length > 1) throw new AggregateError(failures, 'SDK server teardown failed')238
return {}239
}241
/**242
* Dispatch one incoming JSON-RPC request to its typed handler. Throws (→ a243
* JSON-RPC error response) on an unknown method.244
* @param method - the JSON-RPC method name.245
* @param params - the raw params object from the wire.246
* @returns the handler's result, to be serialized as the response.247
*/248
async handleRequest(method: string, params: Record<string, unknown> | undefined): Promise<unknown> {249
switch (method) {250
case 'initialize':251
return this.initialize(params as unknown as InitializeParams)252
case 'session/prompt':253
return this.prompt(params as unknown as SessionPromptParams)254
case 'shutdown':255
return this.shutdown()256
default:257
throw new Error(`unknown DeepSeek Harness SDK runtime method: ${method}`)258
}259
}261
private async getOrCreateSession(sessionId: string): Promise<SessionRecord> {262
if (this.shuttingDown) throw new Error('SDK server is shutting down')263
const existing = this.sessions.get(sessionId)264
if (existing) return existing265
const pending = this.sessionCreations.get(sessionId)266
if (pending) return pending267
const creation = this.createSession(sessionId)268
this.sessionCreations.set(sessionId, creation)269
void creation.then(270
() => { this.sessionCreations.delete(sessionId) },271
() => { this.sessionCreations.delete(sessionId) },272
)273
return creation274
}276
private async createSession(sessionId: string): Promise<SessionRecord> {277
// No preset composition: this server's compositions keep the model-facing278
// rows in the host plane, so this agent reads them from the global layer. A279
// deployment that configures a roster has to join one here first280
// (@deepseek-ai/dsh-agent-preset-registry README, "Composing a child agent").281
const handle = await this.ctx.agents.create({282
sessionId: brandString<SessionId>(sessionId),283
meta: { cwd: this.cwd },284
agentOptions: {285
provider: this.provider,286
model: this.model,287
...this.reasoningEffort === undefined ? {} : { reasoningEffort: this.reasoningEffort },288
...this.maxTokens === undefined ? {} : { maxTokens: this.maxTokens },289
},290
})291
const rec: SessionRecord = { handle }292
this.sessions.set(sessionId, rec)293
return rec294
}296
private hasAdapterFor(provider: string): boolean {297
return this.ctx.get('llm')?.listProviders().some(entry => entry.id === provider) ?? false298
}299
}