dispatch.ts 6.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134
  1. /**
  2. * Fused scope-carrier dispatch for agent-subject operations, plus the assembly
  3. * context builder. The sanctioned ordinary spelling is
  4. * `agentEvents(ctx, agent).waterfall('agent/request', …)`: it builds the scope
  5. * carrier ({@link scopeTarget} keyed by the agent) AND injects the subject as
  6. * the first argument in one move, so a site cannot name a different subject.
  7. * The registry lifecycle pair is the deliberate exception: `enter()` captures
  8. * one stable carrier before commit and `announce()`/detach dispatch through it
  9. * directly, so both lifecycle edges use the same routing identity. The dev
  10. * scoped-dispatch invariant checks both shapes.
  11. *
  12. * @module @deepseek-ai/dsh-agent/dispatch
  13. */
  14. import type { Context, Events } from 'cordis'
  15. import { scopeTarget } from '@deepseek-ai/dsh-scope'
  16. import type { Scoped } from '@deepseek-ai/dsh-scope'
  17. import type { AssembleContext } from '@deepseek-ai/dsh-system-prompt'
  18. import type { Agent } from './types.ts'
  19. /** Extract the parameter tuple from an event handler type (its `this` is not part of the tuple). */
  20. type Params<F> = F extends (...args: infer P) => unknown ? P : never
  21. /** Extract the return type from an event handler type. */
  22. type Return<F> = F extends (...args: never[]) => infer R ? R : never
  23. /**
  24. * The event names whose subject is an agent: handler parameters start with an
  25. * `Agent` AND the handler declares a `Scoped<Agent>` `this` (the scope-carrier
  26. * contract). The `this` check keeps accidental first-parameter-happens-to-be-
  27. * an-Agent events (or zero-arg events, whose parameter tuple would satisfy a
  28. * bare rest-tuple check via callability) out of the fused-dispatch surface.
  29. */
  30. export type AgentSubjectEvent = {
  31. [K in keyof Events]: Events[K] extends (this: Scoped<Agent>, ...args: infer P) => unknown
  32. ? P extends [Agent, ...unknown[]] ? K : never
  33. : never
  34. }[keyof Events]
  35. /** The event arguments AFTER the injected agent subject. */
  36. type Tail<K extends AgentSubjectEvent> = Params<Events[K]> extends [Agent, ...infer R] ? R : never
  37. /**
  38. * The fused dispatcher {@link agentEvents} returns: each method dispatches the
  39. * named agent-subject event with the agent's scope carrier as `thisArg` and
  40. * the agent itself injected as the first event argument.
  41. */
  42. export interface AgentEventDispatch {
  43. /**
  44. * Fire-and-forget notification in the agent's scope. Every listener is
  45. * invoked; synchronous throws and returned-promise rejections are logged and
  46. * contained per listener, so a notification cannot veto lifecycle progress
  47. * or starve a later observer.
  48. * @param name - the agent-subject event to emit.
  49. * @param rest - the event's arguments after the injected agent.
  50. */
  51. emit<K extends AgentSubjectEvent>(name: K, ...rest: Tail<K>): void
  52. /**
  53. * Awaited in-order dispatch (Cordis `serial`) in the agent's scope.
  54. * @param name - the agent-subject event to dispatch.
  55. * @param rest - the event's arguments after the injected agent.
  56. * @returns the serial chain's result (the first bail value, if any).
  57. */
  58. serial<K extends AgentSubjectEvent>(name: K, ...rest: Tail<K>): Promise<Awaited<Return<Events[K]>>>
  59. /**
  60. * Around-middleware dispatch (Cordis `waterfall`) in the agent's scope. The
  61. * declared event parameters already end with the `next` callback, so `rest`
  62. * is exactly the event's arguments after the injected agent — the final
  63. * element being the innermost `next` (the default the listener chain wraps).
  64. * @param name - the agent-subject event to dispatch.
  65. * @param rest - the event's arguments after the injected agent.
  66. * @returns the waterfall's composed result.
  67. */
  68. waterfall<K extends AgentSubjectEvent>(name: K, ...rest: Tail<K>): Return<Events[K]>
  69. }
  70. /**
  71. * Build the fused dispatcher for `agent`'s events (see the module doc). Cheap
  72. * (one carrier + one small object) — dispatch sites create it per run/turn
  73. * rather than caching it on the agent.
  74. * @param ctx - the context to dispatch through (any context of the app).
  75. * @param agent - the subject agent; also the scope-carrier key.
  76. * @returns the fused dispatcher.
  77. */
  78. export function agentEvents(ctx: Context, agent: Agent): AgentEventDispatch {
  79. const carrier: Scoped<Agent> = scopeTarget(agent, agent)
  80. // The ordinary dispatch methods forward through Cordis' variadic mixins. The
  81. // fused (carrier, name, agent, ...rest) tuple is provably a valid argument
  82. // list for the matching thisArg overload, but TypeScript cannot relate the
  83. // generic Tail<K> spread back to that overload's conditional parameter
  84. // tuple — hence one contained, shape-preserving cast per method.
  85. return {
  86. emit(name, ...rest) {
  87. // Cordis emit invokes callbacks through Array.map: one synchronous throw
  88. // starves later listeners, and returned promises are discarded. Agent
  89. // notifications are non-vetoing, so resolve the same filtered callback
  90. // set ourselves and contain both failure modes independently.
  91. const args: unknown[] = [carrier, name, agent, ...rest]
  92. const callbacks = ctx.events.dispatch('emit', args)
  93. for (const callback of callbacks) {
  94. try {
  95. const returned: unknown = callback(...args)
  96. void Promise.resolve(returned).catch((error: unknown) => {
  97. ctx.logger.warn(`agent event "${name}" listener rejected: ${String(error)}`)
  98. })
  99. } catch (error: unknown) {
  100. ctx.logger.warn(`agent event "${name}" listener threw: ${String(error)}`)
  101. }
  102. }
  103. },
  104. async serial(name, ...rest) {
  105. // eslint-disable-next-line @typescript-eslint/unbound-method -- the events mixin accessor returns a pre-bound function
  106. const serial = ctx.serial as (thisArg: Scoped<Agent>, name: string, ...args: unknown[]) => Promise<never>
  107. return await serial(carrier, name, agent, ...rest)
  108. },
  109. waterfall(name, ...rest) {
  110. // eslint-disable-next-line @typescript-eslint/unbound-method -- the events mixin accessor returns a pre-bound function
  111. const waterfall = ctx.waterfall as (thisArg: Scoped<Agent>, name: string, ...args: unknown[]) => never
  112. return waterfall(carrier, name, agent, ...rest)
  113. },
  114. }
  115. }
  116. /**
  117. * The assembly context for one agent's prompt: the typed `agent` DX field and
  118. * the `scope` layer selector, set together (setting `agent` without `scope`
  119. * silently drops the agent's scoped sections/tools from the assembly — the
  120. * dev invariants flag it). THE way the loop (and any custom driver) builds
  121. * its per-step `ctx.systemPrompt.assemble(…)` input.
  122. * @param agent - the agent the assembly is for.
  123. * @returns the context to pass to `assemble()`.
  124. */
  125. export function assembleContextFor(agent: Agent): AssembleContext {
  126. return { agent, scope: agent }
  127. }