返回源码地图

packages/webhook/webhook/src/index.ts

main snapshot · da00f7f5358f · 正文引用章节 17;完整原文可核对,不声称全文件人工逐行审计

完整原文供逐行核对;页面收录不代表每行都经过人工语义审核。MIT 许可见 许可证。

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