dispatch.ts 6.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148
  1. /**
  2. * Agent-scoped dispatch and prompt assembly helpers. Ordinary events use the
  3. * fused dispatcher so subject and scope key cannot diverge; registry lifecycle
  4. * code instead captures one stable carrier for both edges.
  5. * @module @deepseek-ai/dsh-agent/dispatch
  6. */
  7. import type { Context, Events } from 'cordis'
  8. import { scopeTarget } from '@deepseek-ai/dsh-scope'
  9. import type { Scoped } from '@deepseek-ai/dsh-scope'
  10. import type { AssembleContext } from '@deepseek-ai/dsh-system-prompt'
  11. import type { Agent } from './types.ts'
  12. /** Extract the parameter tuple from an event handler type (its `this` is not part of the tuple). */
  13. type Params<F> = F extends (...args: infer P) => unknown ? P : never
  14. /** Extract the return type from an event handler type. */
  15. type Return<F> = F extends (...args: never[]) => infer R ? R : never
  16. /**
  17. * The event names whose subject is an agent: handler parameters start with an
  18. * `Agent` AND the handler declares a `Scoped<Agent>` `this` (the scope-carrier
  19. * contract). The `this` check keeps accidental first-parameter-happens-to-be-
  20. * an-Agent events (or zero-arg events, whose parameter tuple would satisfy a
  21. * bare rest-tuple check via callability) out of the fused-dispatch surface.
  22. */
  23. export type AgentSubjectEvent = {
  24. [K in keyof Events]: Events[K] extends (this: Scoped<Agent>, ...args: infer P) => unknown
  25. ? P extends [Agent, ...unknown[]] ? K : never
  26. : never
  27. }[keyof Events]
  28. /** The event arguments AFTER the injected agent subject. */
  29. type Tail<K extends AgentSubjectEvent> = Params<Events[K]> extends [Agent, ...infer R] ? R : never
  30. /**
  31. * The fused dispatcher {@link agentEvents} returns: each method dispatches the
  32. * named agent-subject event with the agent's scope carrier as `thisArg` and
  33. * the agent itself injected as the first event argument.
  34. */
  35. export interface AgentEventDispatch {
  36. /**
  37. * Fire-and-forget notification in the agent's scope. Every listener is
  38. * invoked; synchronous throws and returned-promise rejections are logged and
  39. * contained per listener, so a notification cannot veto lifecycle progress
  40. * or starve a later observer.
  41. * @param name - the agent-subject event to emit.
  42. * @param rest - the event's arguments after the injected agent.
  43. */
  44. emit<K extends AgentSubjectEvent>(name: K, ...rest: Tail<K>): void
  45. /**
  46. * Awaited in-order dispatch (Cordis `serial`) in the agent's scope.
  47. * @param name - the agent-subject event to dispatch.
  48. * @param rest - the event's arguments after the injected agent.
  49. * @returns the serial chain's result (the first bail value, if any).
  50. */
  51. serial<K extends AgentSubjectEvent>(name: K, ...rest: Tail<K>): Promise<Awaited<Return<Events[K]>>>
  52. /**
  53. * Around-middleware dispatch (Cordis `waterfall`) in the agent's scope. The
  54. * declared event parameters already end with the `next` callback, so `rest`
  55. * is exactly the event's arguments after the injected agent — the final
  56. * element being the innermost `next` (the default the listener chain wraps).
  57. * @param name - the agent-subject event to dispatch.
  58. * @param rest - the event's arguments after the injected agent.
  59. * @returns the waterfall's composed result.
  60. */
  61. waterfall<K extends AgentSubjectEvent>(name: K, ...rest: Tail<K>): Return<Events[K]>
  62. }
  63. /**
  64. * Return the fused scope carrier for one agent subject.
  65. * @param agent - the subject agent and scope key.
  66. * @returns the carrier passed as the event dispatcher `this` value.
  67. */
  68. export function agentCarrier(agent: Agent): Scoped<Agent> {
  69. return scopeTarget(agent, agent)
  70. }
  71. /**
  72. * Build a dispatcher that couples the agent subject to its scope carrier.
  73. * @param ctx - the context to dispatch through (any context of the app).
  74. * @param agent - the subject agent; also the scope-carrier key.
  75. * @returns the fused dispatcher.
  76. */
  77. export function agentEvents(ctx: Context, agent: Agent): AgentEventDispatch {
  78. const carrier = agentCarrier(agent)
  79. // The ordinary dispatch methods forward through Cordis' variadic mixins. The
  80. // fused (carrier, name, agent, ...rest) tuple is provably a valid argument
  81. // list for the matching thisArg overload, but TypeScript cannot relate the
  82. // generic Tail<K> spread back to that overload's conditional parameter
  83. // tuple — hence one contained, shape-preserving cast per method.
  84. return {
  85. emit(name, ...rest) {
  86. // Cordis emit invokes callbacks through Array.map: one synchronous throw
  87. // starves later listeners, and returned promises are discarded. Agent
  88. // notifications are non-vetoing, so resolve the same filtered callback
  89. // set ourselves and contain both failure modes independently.
  90. const args: unknown[] = [carrier, name, agent, ...rest]
  91. const callbacks = ctx.events.dispatch('emit', args)
  92. for (const callback of callbacks) {
  93. try {
  94. const returned: unknown = callback(...args)
  95. void Promise.resolve(returned).catch((error: unknown) => {
  96. ctx.logger.warn(`agent event "${name}" listener rejected: ${String(error)}`)
  97. })
  98. } catch (error: unknown) {
  99. ctx.logger.warn(`agent event "${name}" listener threw: ${String(error)}`)
  100. }
  101. }
  102. },
  103. async serial(name, ...rest) {
  104. // eslint-disable-next-line @typescript-eslint/unbound-method -- the events mixin accessor returns a pre-bound function
  105. const serial = ctx.serial as (thisArg: Scoped<Agent>, name: string, ...args: unknown[]) => Promise<never>
  106. return await serial(carrier, name, agent, ...rest)
  107. },
  108. waterfall(name, ...rest) {
  109. // eslint-disable-next-line @typescript-eslint/unbound-method -- the events mixin accessor returns a pre-bound function
  110. const waterfall = ctx.waterfall as (thisArg: Scoped<Agent>, name: string, ...args: unknown[]) => never
  111. return waterfall(carrier, name, agent, ...rest)
  112. },
  113. }
  114. }
  115. /**
  116. * Emit one contained agent notification without allocating a retained dispatcher.
  117. * @param ctx - the context to dispatch through.
  118. * @param agent - the subject agent and scope key.
  119. * @param name - the agent-subject event to emit.
  120. * @param rest - the event arguments after the injected agent.
  121. */
  122. export function emitAgentEvent<K extends AgentSubjectEvent>(
  123. ctx: Context,
  124. agent: Agent,
  125. name: K,
  126. ...rest: Tail<K>
  127. ): void {
  128. agentEvents(ctx, agent).emit(name, ...rest)
  129. }
  130. /**
  131. * Build the prompt assembly context with agent and scope set together, so
  132. * agent-scoped prompt and tool contributions cannot be silently omitted.
  133. * @param agent - the agent the assembly is for.
  134. * @param signal - the current turn's explicit control signal, when assembly belongs to a turn.
  135. * @returns the context to pass to `assemble()`.
  136. */
  137. export function assembleContextFor(agent: Agent, signal?: AbortSignal): AssembleContext {
  138. return { agent, scope: agent, ...signal === undefined ? {} : { signal } }
  139. }