返回源码地图

vendor/cordis/src/events.ts

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

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

1import { defineProperty } from '@deepseek-ai/cosmokit'
2import type { Promisify } from '@deepseek-ai/cosmokit'
3import { Context } from './context.ts'
4import { Fiber, FiberState } from './fiber.ts'
5import { DisposableList, symbols } from './utils.ts'
6
7/**
8 * Return whether an event result should stop a bail-style dispatch.
9 *
10 * @param value — a listener's return value.
11 * @returns `true` unless `value` is `null`, `false`, or `undefined`.
12 */
13export function isBailed(value: any) {
14 return value !== null && value !== false && value !== undefined
15}
16
17/** Extract the parameter tuple from a function type. */
18export type Parameters<F> = F extends (...args: infer P) => any ? P : never
19/** Extract the return type from a function type. */
20export type ReturnType<F> = F extends (...args: any) => infer R ? R : never
21/** Extract the explicit `this` type from a function type. */
22export type ThisType<F> = F extends (this: infer T, ...args: any) => any ? T : never
23
24/**
25 * Event dispatch strategy used by the event service.
26 *
27 * `emit` runs synchronous listeners without awaiting them, `parallel` awaits
28 * all listeners together, `serial` awaits them in order until one bails,
29 * `bail` stops on the first synchronous bail value, and `waterfall` composes
30 * listeners around a final `next` callback.
31 */
32export type DispatchMode = 'emit' | 'parallel' | 'serial' | 'bail' | 'waterfall'
33
34declare module './context.ts' {
35 export interface Context {
36 /* eslint-disable max-len */
37 /**
38 * Dispatch an event, running all listeners concurrently.
39 *
40 * @param name — the event name.
41 * @param args — arguments passed to every listener.
42 * @returns a promise resolving once every listener has settled.
43 */
44 parallel<K extends keyof Events>(name: K, ...args: Parameters<Events[K]>): Promise<void>
45 /** Same as above, with an explicit `this` for listeners (also used for filtering). */
46 parallel<K extends keyof Events>(thisArg: NoInfer<ThisType<Events[K]>>, name: K, ...args: Parameters<Events[K]>): Promise<void>
47 /**
48 * Dispatch an event synchronously, ignoring listener return values.
49 *
50 * @param name — the event name.
51 * @param args — arguments passed to every listener.
52 */
53 emit<K extends keyof Events>(name: K, ...args: Parameters<Events[K]>): void
54 /** Same as above, with an explicit `this` for listeners (also used for filtering). */
55 emit<K extends keyof Events>(thisArg: NoInfer<ThisType<Events[K]>>, name: K, ...args: Parameters<Events[K]>): void
56 /**
57 * Dispatch an event, awaiting listeners in order until one bails.
58 *
59 * @param name — the event name.
60 * @param args — arguments passed to each listener.
61 * @returns the first bail value (non-null, non-false, non-undefined), if any.
62 */
63 serial<K extends keyof Events>(name: K, ...args: Parameters<Events[K]>): Promisify<ReturnType<Events[K]>>
64 /** Same as above, with an explicit `this` for listeners (also used for filtering). */
65 serial<K extends keyof Events>(thisArg: NoInfer<ThisType<Events[K]>>, name: K, ...args: Parameters<Events[K]>): Promisify<ReturnType<Events[K]>>
66 /**
67 * Dispatch an event, calling listeners in order until one bails.
68 *
69 * @param name — the event name.
70 * @param args — arguments passed to each listener.
71 * @returns the first bail value (non-null, non-false, non-undefined), if any.
72 */
73 bail<K extends keyof Events>(name: K, ...args: Parameters<Events[K]>): ReturnType<Events[K]>
74 /** Same as above, with an explicit `this` for listeners (also used for filtering). */
75 bail<K extends keyof Events>(thisArg: NoInfer<ThisType<Events[K]>>, name: K, ...args: Parameters<Events[K]>): ReturnType<Events[K]>
76 /**
77 * Dispatch an event whose last argument is a `next` continuation.
78 *
79 * Each listener wraps the rest of the chain: calling `next()` invokes the
80 * next listener (finally the built-in behavior); not calling it vetoes.
81 *
82 * @param name — the event name.
83 * @param args — listener arguments; the final one is the innermost `next`.
84 * @returns the outermost listener's return value.
85 */
86 waterfall<K extends keyof Events>(name: K, ...args: Parameters<Events[K]>): ReturnType<Events[K]>
87 /** Same as above, with an explicit `this` for listeners (also used for filtering). */
88 waterfall<K extends keyof Events>(thisArg: NoInfer<ThisType<Events[K]>>, name: K, ...args: Parameters<Events[K]>): ReturnType<Events[K]>
89 /**
90 * Register an event listener owned by the current fiber.
91 *
92 * @param name — the event name to listen for.
93 * @param listener — called with the dispatch arguments.
94 * @param options — listener options; a boolean is shorthand for `prepend`.
95 * @returns a disposer removing the listener; `true` if it was still registered.
96 */
97 on<K extends keyof Events>(name: K, listener: Events[K], options?: boolean | EventOptions): () => boolean
98 /**
99 * Same as `on()`, but the listener disposes itself after its first call.
100 *
101 * @param name — the event name to listen for.
102 * @param listener — called at most once with the dispatch arguments.
103 * @param options — listener options; a boolean is shorthand for `prepend`.
104 * @returns a disposer removing the listener; `true` if it was still registered.
105 */
106 once<K extends keyof Events>(name: K, listener: Events[K], options?: boolean | EventOptions): () => boolean
107 /* eslint-enable max-len */
108 }
109}
110
111/** Options accepted by `ctx.on()` and `ctx.once()`. */
112export interface EventOptions {
113 /** Add the listener before existing listeners for the same event. */
114 prepend?: boolean
115 /** Receive the event regardless of context filter checks. */
116 global?: boolean
117}
118
119/** Registered listener record stored by the event service. */
120export interface Hook extends EventOptions {
121 ctx: Context
122 callback: (...args: any[]) => any
123}
124
125/**
126 * Event bus installed as `ctx.events` and mixed into every context.
127 *
128 * The service supports concurrent, synchronous, serial, bail, and waterfall
129 * dispatch and automatically disposes listeners with their owning fiber.
130 */
131export class EventsService {
132 _hooks: Record<keyof any, Hook[]> = {}
133
134 constructor(private ctx: Context) {
135 defineProperty(this, symbols.tracker, {
136 property: 'ctx',
137 noShadow: true,
138 })
139
140 this.on('internal/listener', function (this: Context, name, listener, options: EventOptions) {
141 if (name === 'internal/update' && !options.global) {
142 const hooks = this.fiber._hooks['internal/update'] ??= new DisposableList()
143 const method = options.prepend ? 'unshift' : 'push'
144 return hooks[method](listener)
145 }
146 })
147
148 this.on('internal/update', function (config, noSave, next) {
149 const cbs = [...this._hooks['internal/update'] || []]
150 const _next = () => {
151 const cb = cbs.shift() ?? next
152 return cb.call(this, config, noSave, _next)
153 }
154 return _next()
155 }, { global: true, prepend: true })
156 }
157
158 /**
159 * Resolve listeners for one dispatch and apply context filtering.
160 *
161 * @param type — the dispatch mode, reported on `internal/dispatch`.
162 * @param args — the raw dispatch arguments; consumed up to the event name.
163 * @returns the matching listener callbacks, bound to the dispatch `this`.
164 */
165 dispatch(type: string, args: any[]) {
166 const thisArg = typeof args[0] === 'object' || typeof args[0] === 'function' ? args.shift() : null
167 const name: string = args.shift()
168 if (!name.startsWith('internal/')) {
169 this.emit('internal/dispatch', type, name, args, thisArg)
170 }
171 const filter = thisArg?.[Context.filter]
172 return (this._hooks[name] || [])
173 .filter(hook => hook.global || !filter || filter.call(thisArg, hook.ctx))
174 .map(hook => hook.callback.bind(thisArg))
175 }
176
177 /**
178 * Run listeners concurrently and wait for all of them.
179 *
180 * @param args — optional `this`, the event name, then listener arguments.
181 * @returns a promise resolving once every listener has settled.
182 */
183 async parallel(...args: any[]) {
184 const results = await Promise.allSettled(this.dispatch('emit', args).map(async cb => cb(...args)))
185 const errors = results.filter((result): result is PromiseRejectedResult => result.status === 'rejected')
186 if (errors.length) throw new AggregateError(errors.map(error => error.reason))
187 }
188
189 /**
190 * Run listeners synchronously without waiting for returned promises.
191 *
192 * @param args — optional `this`, the event name, then listener arguments.
193 */
194 emit(...args: any[]) {
195 this.dispatch('emit', args).map(cb => cb(...args))
196 }
197
198 /**
199 * Run listeners in order, awaiting each, until one returns a bail value.
200 *
201 * @param args — optional `this`, the event name, then listener arguments.
202 * @returns the first bail value (see {@link isBailed}), if any.
203 */
204 async serial(...args: any[]) {
205 for (const cb of this.dispatch('serial', args)) {
206 const result = await cb(...args)
207 if (isBailed(result)) return result
208 }
209 }
210
211 /**
212 * Run listeners synchronously until one returns a bail value.
213 *
214 * @param args — optional `this`, the event name, then listener arguments.
215 * @returns the first bail value (see {@link isBailed}), if any.
216 */
217 bail(...args: any[]) {
218 for (const cb of this.dispatch('bail', args)) {
219 const result = cb(...args)
220 if (isBailed(result)) return result
221 }
222 }
223
224 /**
225 * Compose listeners around the final `next` callback.
226 *
227 * The last dispatch argument is treated as the innermost `next`. Listeners
228 * run outermost-first; a listener that does not call `next()` vetoes the
229 * rest of the chain, including the built-in behavior.
230 *
231 * @param args — optional `this`, the event name, listener arguments, then `next`.
232 * @returns the outermost listener's return value.
233 */
234 waterfall(...args: any[]) {
235 const cbs = this.dispatch('waterfall', args)
236 const inner = args.pop()
237 const next = () => {
238 const cb = cbs.shift() ?? inner
239 return cb(...args)
240 }
241 args.push(next)
242 return next()
243 }
244
245 /**
246 * Store a listener record as an effect on the current fiber.
247 *
248 * @param label — effect label shown in fiber diagnostics.
249 * @param hooks — the listener list for one event.
250 * @param callback — the listener to store.
251 * @param options — placement and filtering options.
252 * @returns a disposer that unregisters the listener.
253 */
254 register(label: string, hooks: Hook[], callback: any, options: EventOptions): () => void {
255 const method = options.prepend ? 'unshift' : 'push'
256 return this.ctx.fiber.effect(() => {
257 hooks[method]({ ctx: this.ctx, callback, ...options })
258 return () => this.unregister(hooks, callback)
259 }, label)
260 }
261
262 /**
263 * Remove a stored listener record.
264 *
265 * @param hooks — the listener list for one event.
266 * @param callback — the listener to remove.
267 * @returns `true` if the listener was found and removed.
268 */
269 unregister(hooks: Hook[], callback: any) {
270 const index = hooks.findIndex(hook => hook.callback === callback)
271 if (index >= 0) {
272 hooks.splice(index, 1)
273 return true
274 }
275 }
276
277 /**
278 * Register an event listener owned by the current fiber.
279 *
280 * The listener is removed automatically when the fiber unloads. Throws
281 * `CordisError('INACTIVE_EFFECT')` if the fiber is already disposed.
282 *
283 * @param name — the event name to listen for.
284 * @param listener — called with the dispatch arguments.
285 * @param options — listener options; a boolean is shorthand for `prepend`.
286 * @returns a disposer removing the listener; `true` if it was still registered.
287 */
288 on(name: string | symbol, listener: (...args: any) => any, options?: boolean | EventOptions) {
289 if (typeof options !== 'object') {
290 options = { prepend: options }
291 }
292
293 // handle special events
294 this.ctx.fiber.assertActive()
295 listener = this.ctx.reflect.bind(listener)
296 const result = this.bail(this.ctx, 'internal/listener', name, listener, options)
297 if (result) return result
298
299 const hooks = this._hooks[name] ||= []
300 const label = `ctx.on(${typeof name === 'string' ? JSON.stringify(name) : name.toString()})`
301 return this.register(label, hooks, listener, options)
302 }
303
304 /**
305 * Register an event listener that disposes itself after the first call.
306 *
307 * @param name — the event name to listen for.
308 * @param listener — called at most once with the dispatch arguments.
309 * @param options — listener options; a boolean is shorthand for `prepend`.
310 * @returns a disposer removing the listener; `true` if it was still registered.
311 */
312 once(name: string, listener: (...args: any) => any, options?: boolean | EventOptions) {
313 const dispose = this.on(name, function (...args: any[]) {
314 dispose()
315 return listener.apply(this, args)
316 }, options)
317 return dispose
318 }
319}
320
321/**
322 * Built-in framework events used by core services and extension points.
323 *
324 * Plugin and status events track fiber lifecycle, service events observe
325 * dependency registration, update/get/set/listener events allow core services
326 * to intercept runtime operations, and `internal/dispatch` exposes event-bus
327 * diagnostics before public events are delivered.
328 */
329export interface Events {
330 /** A plugin fiber was created or its uid was cleared on disposal. */
331 'internal/plugin'(fiber: Fiber): void
332 /** A fiber changed lifecycle state; receives the fiber and its previous state. */
333 'internal/status'(fiber: Fiber, oldValue: FiberState): void
334 /**
335 * Resolve raw plugin config after the fiber's injections become active.
336 * @param config - the raw config for this activation.
337 * @mode waterfall
338 */
339 'internal/config'(this: Fiber, config: any, next: () => any): any
340 /** Interception hook for a service binding (no core producer). */
341 'internal/service'(this: Context, name: string, value: any): void
342 /** Waterfall: a fiber config update is being applied; skip `next()` to veto. */
343 'internal/update'(this: Fiber, config: any, noSave: boolean, next: () => void): void
344 /** Waterfall: a service is being read through the context proxy. */
345 'internal/get'(ctx: Context, name: string, error: Error, next: () => any): any
346 /** Waterfall: a service is being written through the context proxy. */
347 'internal/set'(ctx: Context, name: string, value: any, error: Error, next: () => boolean): boolean
348 /** Bail: a listener is being registered; a non-null result replaces registration. */
349 'internal/listener'(this: Context, name: string, listener: any, prepend: boolean): void
350 /** An event is being dispatched to listeners (fired for non-internal events only). */
351 'internal/dispatch'(mode: DispatchMode, name: string, args: any[], thisArg: any): void
352}