1
import { defineProperty } from '@deepseek-ai/cosmokit'2
import type { Promisify } from '@deepseek-ai/cosmokit'3
import { Context } from './context.ts'4
import { Fiber, FiberState } from './fiber.ts'5
import { DisposableList, symbols } from './utils.ts'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
*/13
export function isBailed(value: any) {14
return value !== null && value !== false && value !== undefined15
}17
/** Extract the parameter tuple from a function type. */18
export type Parameters<F> = F extends (...args: infer P) => any ? P : never19
/** Extract the return type from a function type. */20
export type ReturnType<F> = F extends (...args: any) => infer R ? R : never21
/** Extract the explicit `this` type from a function type. */22
export type ThisType<F> = F extends (this: infer T, ...args: any) => any ? T : never24
/**25
* Event dispatch strategy used by the event service.26
*27
* `emit` runs synchronous listeners without awaiting them, `parallel` awaits28
* all listeners together, `serial` awaits them in order until one bails,29
* `bail` stops on the first synchronous bail value, and `waterfall` composes30
* listeners around a final `next` callback.31
*/32
export type DispatchMode = 'emit' | 'parallel' | 'serial' | 'bail' | 'waterfall'34
declare 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]>): void54
/** 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]>): void56
/**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 the80
* 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): () => boolean98
/**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): () => boolean107
/* eslint-enable max-len */108
}109
}111
/** Options accepted by `ctx.on()` and `ctx.once()`. */112
export interface EventOptions {113
/** Add the listener before existing listeners for the same event. */114
prepend?: boolean115
/** Receive the event regardless of context filter checks. */116
global?: boolean117
}119
/** Registered listener record stored by the event service. */120
export interface Hook extends EventOptions {121
ctx: Context122
callback: (...args: any[]) => any123
}125
/**126
* Event bus installed as `ctx.events` and mixed into every context.127
*128
* The service supports concurrent, synchronous, serial, bail, and waterfall129
* dispatch and automatically disposes listeners with their owning fiber.130
*/131
export class EventsService {132
_hooks: Record<keyof any, Hook[]> = {}134
constructor(private ctx: Context) {135
defineProperty(this, symbols.tracker, {136
property: 'ctx',137
noShadow: true,138
})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
})148
this.on('internal/update', function (config, noSave, next) {149
const cbs = [...this._hooks['internal/update'] || []]150
const _next = () => {151
const cb = cbs.shift() ?? next152
return cb.call(this, config, noSave, _next)153
}154
return _next()155
}, { global: true, prepend: true })156
}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() : null167
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
}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
}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
}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 result208
}209
}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 result221
}222
}224
/**225
* Compose listeners around the final `next` callback.226
*227
* The last dispatch argument is treated as the innermost `next`. Listeners228
* run outermost-first; a listener that does not call `next()` vetoes the229
* 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() ?? inner239
return cb(...args)240
}241
args.push(next)242
return next()243
}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
}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 true274
}275
}277
/**278
* Register an event listener owned by the current fiber.279
*280
* The listener is removed automatically when the fiber unloads. Throws281
* `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
}293
// handle special events294
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 result299
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
}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 dispose318
}319
}321
/**322
* Built-in framework events used by core services and extension points.323
*324
* Plugin and status events track fiber lifecycle, service events observe325
* dependency registration, update/get/set/listener events allow core services326
* to intercept runtime operations, and `internal/dispatch` exposes event-bus327
* diagnostics before public events are delivered.328
*/329
export interface Events {330
/** A plugin fiber was created or its uid was cleared on disposal. */331
'internal/plugin'(fiber: Fiber): void332
/** A fiber changed lifecycle state; receives the fiber and its previous state. */333
'internal/status'(fiber: Fiber, oldValue: FiberState): void334
/**335
* Resolve raw plugin config after the fiber's injections become active.336
* @param config - the raw config for this activation.337
* @mode waterfall338
*/339
'internal/config'(this: Fiber, config: any, next: () => any): any340
/** Interception hook for a service binding (no core producer). */341
'internal/service'(this: Context, name: string, value: any): void342
/** Waterfall: a fiber config update is being applied; skip `next()` to veto. */343
'internal/update'(this: Fiber, config: any, noSave: boolean, next: () => void): void344
/** Waterfall: a service is being read through the context proxy. */345
'internal/get'(ctx: Context, name: string, error: Error, next: () => any): any346
/** Waterfall: a service is being written through the context proxy. */347
'internal/set'(ctx: Context, name: string, value: any, error: Error, next: () => boolean): boolean348
/** Bail: a listener is being registered; a non-null result replaces registration. */349
'internal/listener'(this: Context, name: string, listener: any, prepend: boolean): void350
/** An event is being dispatched to listeners (fired for non-internal events only). */351
'internal/dispatch'(mode: DispatchMode, name: string, args: any[], thisArg: any): void352
}