1
/** Fire-and-forget webhook rule registry and Workspace-backed Session runtime. */3
import { Context, Service } from '@deepseek-ai/cordis'4
import { errorChain } from '@deepseek-ai/dsh-llm'5
import { deepFreeze, snapshotJsonValue } from '@deepseek-ai/dsh-util-values'6
import type { WebhookRuleId } from './brand.ts'7
import { createWebhookSession } from './session.ts'8
import type { VerifiedWebhookDelivery, WebhookRule, WebhookSessionRequest } from './types.ts'10
export * from './brand.ts'11
export type * from './types.ts'13
declare module '@deepseek-ai/cordis' {14
interface Context {15
webhookRuntime: WebhookRuntime16
}17
}19
/** Internal type erasure after public generic registration validates the provider kind. */20
interface AnyWebhookRule {21
readonly id: WebhookRuleId22
readonly kind: string23
run(24
delivery: Readonly<VerifiedWebhookDelivery>,25
signal: AbortSignal,26
): WebhookSessionRequest | null | Promise<WebhookSessionRequest | null>27
}29
/** One effect-owned rule registration and the invocations that currently use it. */30
interface RuleRegistration {31
readonly rule: AnyWebhookRule32
readonly controller: AbortController33
readonly active: Set<Promise<void>>34
closing: boolean35
disposal?: Promise<void>36
}38
/** Validate and detach one delivery before sharing it across arbitrary rules. */39
function 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
}57
/** Fire-and-forget rule runtime. Session creation is the only built-in action. */58
export class WebhookRuntime extends Service {59
static inject = [60
'agents',61
'agentDefaultModel',62
'agentPresets',63
'permissionPresets',64
'sessionTitle',65
'workspaceRegistry',66
]68
private readonly rules = new Map<WebhookRuleId, RuleRegistration>()69
private readonly selfCtx: Context70
private closing = false72
constructor(ctx: Context) {73
super(ctx, 'webhookRuntime')74
this.selfCtx = ctx75
ctx.effect(() => async () => {76
this.closing = true77
/* 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
}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
}101
// The public generic preserves adapter-specific authoring types. The runtime102
// stores one erased callback after validating the shared provider tag.103
const erased = rule as AnyWebhookRule104
let registration!: RuleRegistration105
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
}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) continue131
this.startInvocation(registration, snapshot)132
}133
}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
}164
/** Memoized registration teardown: hide, abort, then drain. */165
private disposeRegistration(registration: RuleRegistration): Promise<void> {166
registration.disposal ??= (async () => {167
registration.closing = true168
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.disposal175
}176
}178
export default WebhookRuntime