index.ts 7.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178
  1. /** Fire-and-forget webhook rule registry and Workspace-backed Session runtime. */
  2. import { Context, Service } from '@deepseek-ai/cordis'
  3. import { deepFreeze, errorChain } from '@deepseek-ai/dsh-llm'
  4. import { snapshotJsonValue } from '@deepseek-ai/dsh-session'
  5. import type { WebhookRuleId } from './brand.ts'
  6. import { createWebhookSession } from './session.ts'
  7. import type { VerifiedWebhookDelivery, WebhookRule, WebhookSessionRequest } from './types.ts'
  8. export * from './brand.ts'
  9. export type * from './types.ts'
  10. declare module '@deepseek-ai/cordis' {
  11. interface Context {
  12. webhookRuntime: WebhookRuntime
  13. }
  14. }
  15. /** Internal type erasure after public generic registration validates the provider kind. */
  16. interface AnyWebhookRule {
  17. readonly id: WebhookRuleId
  18. readonly kind: string
  19. run(
  20. delivery: Readonly<VerifiedWebhookDelivery>,
  21. signal: AbortSignal,
  22. ): WebhookSessionRequest | null | Promise<WebhookSessionRequest | null>
  23. }
  24. /** One effect-owned rule registration and the invocations that currently use it. */
  25. interface RuleRegistration {
  26. readonly rule: AnyWebhookRule
  27. readonly controller: AbortController
  28. readonly active: Set<Promise<void>>
  29. closing: boolean
  30. disposal?: Promise<void>
  31. }
  32. /** Validate and detach one delivery before sharing it across arbitrary rules. */
  33. function snapshotDelivery(delivery: VerifiedWebhookDelivery): VerifiedWebhookDelivery {
  34. if (typeof delivery.kind !== 'string' || delivery.kind.trim() === '') {
  35. throw new TypeError('webhook delivery kind must be a non-empty string')
  36. }
  37. if (typeof delivery.source !== 'string' || delivery.source.trim() === '') {
  38. throw new TypeError('webhook delivery source must be a non-empty string')
  39. }
  40. if (typeof delivery.deliveryId !== 'string' || delivery.deliveryId.trim() === '') {
  41. throw new TypeError('webhook delivery id must be a non-empty string')
  42. }
  43. if (!Number.isSafeInteger(delivery.receivedAt) || delivery.receivedAt < 0) {
  44. throw new TypeError('webhook delivery receivedAt must be a non-negative safe integer')
  45. }
  46. const snapshot = snapshotJsonValue(delivery)
  47. if (snapshot === undefined) throw new TypeError('webhook delivery must be lossless JSON')
  48. return deepFreeze(snapshot)
  49. }
  50. /** Fire-and-forget rule runtime. Session creation is the only built-in action. */
  51. export class WebhookRuntime extends Service {
  52. static inject = [
  53. 'agents',
  54. 'agentDefaultModel',
  55. 'agentPresets',
  56. 'permissionPresets',
  57. 'sessionTitle',
  58. 'workspaceRegistry',
  59. ]
  60. private readonly rules = new Map<WebhookRuleId, RuleRegistration>()
  61. private readonly selfCtx: Context
  62. private closing = false
  63. constructor(ctx: Context) {
  64. super(ctx, 'webhookRuntime')
  65. this.selfCtx = ctx
  66. ctx.effect(() => async () => {
  67. this.closing = true
  68. /* v8 ignore next -- caller-owned registration effects normally dispose first; this covers provider-first unload. */
  69. await Promise.all(
  70. [...this.rules.values()].map(rule => this.disposeRegistration(rule)),
  71. )
  72. }, 'webhookRuntime.lifecycle()')
  73. }
  74. /**
  75. * Register one trusted programmatic rule.
  76. * @param rule - unique id, provider kind, and arbitrary callback.
  77. * @returns awaitable effect disposer that aborts and drains this rule's active callbacks.
  78. */
  79. register<K extends string>(rule: WebhookRule<K>): () => Promise<void> {
  80. if (this.closing) throw new Error('webhook runtime is closing')
  81. if (typeof rule.id !== 'string' || rule.id.trim() === '') {
  82. throw new TypeError('webhook rule id must be a non-empty string')
  83. }
  84. if (typeof rule.kind !== 'string' || rule.kind.trim() === '') {
  85. throw new TypeError(`webhook rule "${String(rule.id)}" kind must be a non-empty string`)
  86. }
  87. if (typeof rule.run !== 'function') {
  88. throw new TypeError(`webhook rule "${String(rule.id)}" requires run()`)
  89. }
  90. // The public generic preserves adapter-specific authoring types. The runtime
  91. // stores one erased callback after validating the shared provider tag.
  92. const erased = rule as unknown as AnyWebhookRule
  93. let registration!: RuleRegistration
  94. const disposeEffect = this.ctx.effect(() => {
  95. /* v8 ignore next -- no await separates the public liveness check from this initializer. */
  96. if (this.closing) throw new Error('webhook runtime is closing')
  97. if (this.rules.has(rule.id)) throw new Error(`webhook rule "${rule.id}" is already registered`)
  98. registration = {
  99. rule: erased,
  100. controller: new AbortController(),
  101. active: new Set(),
  102. closing: false,
  103. }
  104. this.rules.set(rule.id, registration)
  105. return () => this.disposeRegistration(registration)
  106. }, `webhookRuntime.register(${rule.id})`)
  107. return async () => { await disposeEffect() }
  108. }
  109. /**
  110. * Start every currently matching rule and return before any callback settles.
  111. * @param delivery - authenticated provider data; snapshotted before dispatch.
  112. * @throws synchronously when the runtime is closing or the delivery is malformed.
  113. */
  114. dispatch<K extends string>(delivery: VerifiedWebhookDelivery<K>): void {
  115. if (this.closing) throw new Error('webhook runtime is closing')
  116. const snapshot = snapshotDelivery(delivery)
  117. for (const registration of [...this.rules.values()]) {
  118. if (registration.closing || registration.rule.kind !== snapshot.kind) continue
  119. this.startInvocation(registration, snapshot)
  120. }
  121. }
  122. /** Start one contained invocation and attach it to registration teardown. */
  123. private startInvocation(registration: RuleRegistration, delivery: VerifiedWebhookDelivery): void {
  124. const tracked = Promise.resolve().then(async () => {
  125. registration.controller.signal.throwIfAborted()
  126. const request = await registration.rule.run(delivery, registration.controller.signal)
  127. registration.controller.signal.throwIfAborted()
  128. if (request !== null) {
  129. await createWebhookSession(
  130. this.selfCtx,
  131. delivery,
  132. registration.rule.id,
  133. request,
  134. registration.controller.signal,
  135. )
  136. }
  137. }).catch((error: unknown) => {
  138. const invocation = `webhook: provider=${JSON.stringify(delivery.kind)} source=${JSON.stringify(delivery.source)} `
  139. + `delivery=${JSON.stringify(delivery.deliveryId)} rule=${JSON.stringify(registration.rule.id)}`
  140. if (registration.controller.signal.aborted) {
  141. this.selfCtx.logger.debug(`${invocation} stopped after disposal: ${errorChain(error)}`)
  142. } else {
  143. this.selfCtx.logger.warn(`${invocation} failed: ${errorChain(error)}`)
  144. }
  145. }).finally(() => {
  146. registration.active.delete(tracked)
  147. })
  148. registration.active.add(tracked)
  149. }
  150. /** Memoized registration teardown: hide, abort, then drain. */
  151. private disposeRegistration(registration: RuleRegistration): Promise<void> {
  152. registration.disposal ??= (async () => {
  153. registration.closing = true
  154. this.rules.delete(registration.rule.id)
  155. registration.controller.abort(new Error(`webhook rule "${registration.rule.id}" was disposed`))
  156. while (registration.active.size > 0) {
  157. await Promise.allSettled([...registration.active])
  158. }
  159. })()
  160. return registration.disposal
  161. }
  162. }
  163. export default WebhookRuntime