index.ts 3.7 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192
  1. /**
  2. * SDK-facing JSON-RPC plugin over stdio. An external `cordis.yml` decides
  3. * whether to load it; see the single-executable Agent Note and package README.
  4. * Stdout is reserved for protocol frames, so the tree must not load a stdout logger.
  5. * This plugin answers `shutdown`, disposes the complete root runtime, and exits 0; the app bin
  6. * owns EOF and signal exits. Keep named plugin exports with no default export so
  7. * Loader `unwrapExports` preserves `name`, `inject`, `Config`, and `apply`.
  8. *
  9. * @module @deepseek-ai/dsh-jsonrpc
  10. */
  11. import type { Context } from 'cordis'
  12. import type { Readable, Writable } from 'node:stream'
  13. import Schema from 'schemastery'
  14. import { JsonRpcLineTransport } from '@deepseek-ai/dsh-sdk-protocol'
  15. import { HarnessSdkServer } from './server.ts'
  16. export * from './server.ts'
  17. export const name = 'jsonrpc'
  18. // Only the agent factory is required; initialize reads the optional LLM seam with ctx.get().
  19. export const inject = ['agents']
  20. /** JSON-RPC deployment config plus runtime-only test seams. */
  21. export interface JsonRpcConfig {
  22. /** Report max-token turn/subagent termination as a successful SDK result. */
  23. maxTokensAsSuccess?: boolean
  24. /** Transport input override; production uses `process.stdin`. */
  25. input?: Readable
  26. /** Transport output override; production uses `process.stdout`. */
  27. output?: Writable
  28. /** Process-exit override; production uses `process.exit`. */
  29. exit?: (code: number) => void
  30. }
  31. export const Config: Schema<JsonRpcConfig> = Schema.object({
  32. maxTokensAsSuccess: Schema.boolean().default(false),
  33. })
  34. /**
  35. * Serve SDK requests over the configured streams. Effect disposal shuts down
  36. * SDK-created agents and closes the transport. A `shutdown` response is flushed
  37. * before the root runtime is disposed and the process exits 0; the app bin
  38. * owns root-context disposal for EOF and signals.
  39. */
  40. export function apply(ctx: Context, config: JsonRpcConfig): void {
  41. // Cordis applies the schema default before invoking the plugin.
  42. const resolvedConfig = config as JsonRpcConfig & { maxTokensAsSuccess: boolean }
  43. // Protocol shutdown owns the complete runtime process, so it must await the
  44. // root lifecycle (including persistence) before exiting.
  45. const rootFiber = ctx.root.fiber
  46. /* v8 ignore next -- production stdio wiring; tests always inject the runtime seams */
  47. const input = config.input ?? process.stdin
  48. /* v8 ignore next -- production stdio wiring; tests always inject the runtime seams */
  49. const output = config.output ?? process.stdout
  50. /* v8 ignore next -- production exit wiring; tests always inject the runtime seams */
  51. const exit = config.exit ?? ((code: number): void => { process.exit(code) })
  52. const transport = new JsonRpcLineTransport(input, output)
  53. const server = new HarnessSdkServer(ctx, transport, {
  54. maxTokensAsSuccess: resolvedConfig.maxTokensAsSuccess,
  55. })
  56. // Share one exit task so racing shutdown requests cannot dispose the root or
  57. // exit the process more than once.
  58. let exitTask: Promise<void> | undefined
  59. const disposeAndExit = (): Promise<void> => {
  60. exitTask ??= (async () => {
  61. await Promise.allSettled([Promise.resolve().then(() => transport.flush())])
  62. await Promise.allSettled([Promise.resolve().then(() => rootFiber.dispose())])
  63. exit(0)
  64. })()
  65. return exitTask
  66. }
  67. transport.onRequest(async (method, params) => {
  68. const result = await server.handleRequest(method, params)
  69. if (method === 'shutdown') {
  70. // Run after the handler result is written; the task then flushes, disposes, and exits.
  71. setImmediate(() => { void disposeAndExit() })
  72. }
  73. return result
  74. })
  75. ctx.effect(() => {
  76. transport.start()
  77. return async () => {
  78. await server.shutdown()
  79. transport.close()
  80. }
  81. }, 'jsonrpc.serve')
  82. }