1
import { defineProperty, isNullable } from '@deepseek-ai/cosmokit'2
import type { Awaitable, Dict } from '@deepseek-ai/cosmokit'3
import { Context } from './context.ts'4
import type { Plugin } from './registry.ts'5
import { buildOuterStack, composeError, DisposableList, getTraceable, isConstructor, isObject, symbols } from './utils.ts'6
import type { Impl } from './reflect.ts'7
import type { StandardSchemaV1 } from '@standard-schema/spec'9
declare module './context.ts' {10
export interface Context extends Pick<Fiber, 'effect'> {11
/** The fiber (plugin runtime instance) that owns this context. */12
fiber: Fiber13
}14
}16
const kValidationError = Symbol.for('ValidationError')18
/** Error raised when plugin configuration fails standard-schema validation. */19
export class ValidationError extends TypeError {20
name = 'ValidationError'22
/**23
* Build the aggregated message from schema issues.24
*25
* @param issues — the standard-schema issues, one message line each.26
*/27
constructor(issues: readonly StandardSchemaV1.Issue[]) {28
super(`invalid config:\n` + issues.map(issue => {29
if (issue.path) {30
return ` - ${issue.message} (at ${issue.path.join('.')})`31
} else {32
return ` - ${issue.message}`33
}34
}).join('\n'))35
}36
}38
Object.defineProperty(ValidationError.prototype, kValidationError, {39
value: true,40
})42
/**43
* Validate and normalize config for a plugin runtime before it starts.44
*45
* @param runtime — the plugin runtime whose `Config` schema to apply.46
* @param config — the raw user config.47
* @returns the validated config, or `config` unchanged if the runtime has no schema.48
* @throws {ValidationError} when validation reports issues.49
*/50
export function resolveConfig(runtime: Plugin.Runtime, config: any) {51
if (!runtime.Config) return config52
// TODO: async validation53
const result = runtime.Config['~standard'].validate(config)54
if ('then' in result) {55
throw new TypeError('Async config validation is not supported')56
}57
if (result.issues) {58
throw new ValidationError(result.issues)59
} else {60
return result.value61
}62
}64
interface AsyncDisposable<T extends Awaitable<void> = Awaitable<void>> extends PromiseLike<() => T> {65
(): T66
}68
/**69
* Function returned by an effect to release resources during disposal.70
*71
* Disposers run in reverse registration order when the owning fiber unloads;72
* they may be async, in which case unloading awaits them.73
*/74
export type Disposable<T = any> = () => T76
/**77
* Effect body result accepted by `ctx.effect()` and plugin startup.78
*79
* Either a single disposer, a promise of one, or a (possibly async) iterable80
* yielding several — generator effects register each yielded disposer as it81
* is produced.82
*/83
export type Effect<T = any> =84
| SyncEffect<T>85
| AsyncEffect<T>87
type SyncEffect<T = any> =88
| Disposable<T>89
| Iterable<Disposable<T>, void, void>91
type AsyncEffect<T = any> =92
| Promise<Disposable<T>>93
| AsyncIterable<Disposable<T>, void, void>95
/** Tree node used to expose nested effect labels for diagnostics. */96
export interface EffectMeta {97
/** Human-readable effect label, e.g. `ctx.on("event")` or `ctx.provide("name")`. */98
label: string99
/** Metadata of nested effects registered while this effect ran. */100
children: EffectMeta[]101
}103
interface EffectRunner<T> {104
epoch: T105
execute: () => any106
collect: (dispose: Disposable) => void107
getOuterStack: () => string[]108
}110
// Public effect disposers remain single-shot, but structural owners and outer111
// effects must still be able to join a cleanup that another caller started.112
const effectInertia = new WeakMap<Disposable, () => void | Promise<void>>()114
function runDisposable(dispose: Disposable) {115
const result = dispose()116
return effectInertia.get(dispose)?.() ?? result117
}119
/** Notify plugin teardown without allowing one observer to break ownership cleanup. */120
function emitPluginDisposed(context: Context, fiber: Fiber) {121
const args: any[] = ['internal/plugin', fiber]122
let callbacks: Function[]123
try {124
callbacks = context.events.dispatch('emit', args)125
} catch (error) {126
context.logger.error(error)127
return128
}129
for (const callback of callbacks) {130
try {131
const returned = callback(...args)132
void Promise.resolve(returned).catch(error => context.logger.error(error))133
} catch (error) {134
context.logger.error(error)135
}136
}137
}139
/**140
* Lifecycle state for one plugin fiber.141
*142
* `PENDING` — waiting for required services; `LOADING` — the plugin callback143
* is running; `ACTIVE` — loaded and providing; `FAILED` — the callback or its144
* config threw; `UNLOADING` — disposers are running; `DISPOSED` — the fiber145
* was removed and cannot restart.146
*/147
export const enum FiberState {148
PENDING,149
LOADING,150
ACTIVE,151
FAILED,152
DISPOSED,153
UNLOADING,154
}156
/** Framework error with a stable machine-readable code. */157
export class CordisError extends Error {158
/**159
* @param code — the stable error code; also the default message.160
* @param message — optional human-readable override.161
*/162
constructor(public code: CordisError.Code, message?: string) {163
super(message ?? CordisError.Code[code])164
}165
}167
/** Cordis error code definitions. */168
export namespace CordisError {169
export type Code = keyof typeof Code171
export const Code = {172
INACTIVE_EFFECT: 'cannot create effect on inactive context',173
} as const174
}176
const INACTIVE = '__INACTIVE__'178
/**179
* Runtime instance of one plugin application.180
*181
* A fiber tracks dependency state, validated config, lifecycle effects, and182
* cleanup for the plugin context returned by `ctx.plugin()`.183
*/184
export class Fiber {185
/** Unique id within the registry; 0 for the root fiber, `null` once disposed. */186
public uid: number | null187
/** The context this fiber's plugin runs in (extends the parent context). */188
public readonly ctx: Context189
/** The validated plugin config (updated by `update()`). */190
public config: any191
/** The raw plugin config, re-resolved before each activation. */192
public _config: any193
/** Current lifecycle state; transitions emit `internal/status`. */194
public state = FiberState.PENDING195
/** Dispose this fiber: unload the plugin, then settle once cleanup finished. */196
public readonly dispose: () => Promise<void>197
/** Snapshot of required service implementations while loaded; `undefined` otherwise. */198
public store: Dict<Impl> | undefined199
/** The in-flight load/unload transition, if one is currently running. */200
public inertia: Promise<void> | undefined202
public readonly _hooks: Dict<DisposableList<Function>> = Object.create(null)203
public readonly _disposables = new DisposableList<Disposable>()205
// Same as `this.ctx`, but with a more specific type.206
protected context: Context208
private _error: any209
private _runner: EffectRunner<string>210
private _store: Dict<Impl> = Object.create(null)212
/**213
* Create a fiber. Plugin authors normally obtain fibers from `ctx.plugin()`214
* rather than constructing them directly.215
*216
* @param parent — the context the plugin was loaded from.217
* @param config — raw config, validated against the runtime's schema.218
* @param inject — resolved dependency map (service name → intercept config).219
* @param runtime — the shared plugin runtime, or `null` for the root fiber.220
* @param getOuterStack — captures the caller stack for effect diagnostics.221
*/222
constructor(223
public parent: Context,224
config: any,225
public inject: Dict<any>,226
public runtime: Plugin.Runtime | null,227
getOuterStack: () => string[],228
) {229
this._config = config230
const collect = (dispose: Disposable) => {231
this._disposables.push(dispose)232
}234
if (runtime) {235
this.uid = parent.registry.counter236
this.ctx = this.context = parent.extend({ fiber: this })238
const injectEntries = Object.entries(this.inject)239
if (injectEntries.length) {240
this.ctx[Context.intercept] = Object.create(parent[Context.intercept])241
for (const [name, config] of injectEntries) {242
if (isNullable(config)) continue243
this.ctx[Context.intercept][name] = config244
}245
}247
this._runner = {248
epoch: INACTIVE,249
getOuterStack,250
execute: function () {251
if (isConstructor(runtime.callback)) {252
// eslint-disable-next-line new-cap253
const instance = new runtime.callback(this.ctx, this.config)254
for (const hook of instance?.[symbols.initHooks] ?? []) {255
hook()256
}257
return instance?.[symbols.init]?.()258
} else {259
return runtime.callback(this.ctx, this.config)260
}261
},262
collect,263
}265
this.dispose = parent.fiber.effect(() => {266
const remove = runtime.fibers.push(this)267
return async () => {268
this.uid = null269
emitPluginDisposed(this.context, this)270
if (this.ctx.registry.has(runtime.callback)) {271
remove()272
if (!runtime.fibers.length) {273
this.ctx.registry.delete(runtime.callback)274
}275
}276
this._setEpoch(INACTIVE)277
// A PENDING fiber can already own effects registered by an278
// internal/plugin observer. Its epoch is still INACTIVE, so279
// _setEpoch() has no transition to drive; explicitly unload that280
// pre-activation work before reporting disposal complete.281
if (!this.inertia) {282
this._updateState(() => {283
this.inertia = this._unload()284
return FiberState.UNLOADING285
})286
}287
// `this.inertia` itself should never reject — both `_reload` and288
// `_unload` swallow their own work errors via `ctx.logger.error`.289
// If it *does* reject, the only remaining cause is the logger290
// itself failing, which we can't recover from in this exact spot291
// (calling the logger again is what just failed). Let the292
// rejection propagate; process-level crash is the honest outcome.293
while (this.inertia) {294
await this.inertia295
}296
}297
}, 'ctx.plugin()')299
try {300
// Publish only after the parent owns a fully assigned disposer. A301
// synchronous observer may dispose either this fiber or its parent.302
this.context.emit('internal/plugin', this)303
} catch (error) {304
// Publication failed synchronously. The disposer removes the child305
// from both the parent and runtime before control escapes.306
void Promise.resolve(this.dispose()).catch(reason => this.ctx.logger.error(reason))307
throw error308
}310
// Keep the initial notification's historical PENDING view. The loader311
// may also extend `inject` in that notification, so resolve dependencies312
// only after publication. A reentrant parent unload makes the child313
// disposer responsible for draining any PENDING effects instead.314
if (this.uid !== null && parent.fiber.state !== FiberState.UNLOADING) {315
for (const name of Object.keys(this.inject)) {316
this._checkImpl(name)317
}318
this._refresh()319
}320
} else {321
this.uid = 0322
this.ctx = this.context = parent323
this.state = FiberState.ACTIVE324
this.store = Object.create(null)325
this._runner = {326
epoch: '',327
getOuterStack,328
execute: () => {},329
collect,330
}331
this.dispose = () => this.restart()332
}333
}335
/** The plugin's display name, inherited from the nearest named ancestor, else `'root'`. */336
get name() {337
let fiber: Fiber = this338
do {339
if (fiber.runtime?.name) return fiber.runtime.name340
fiber = fiber.parent.fiber341
} while (fiber !== fiber.parent.fiber)342
return 'root'343
}345
/**346
* Throw if the fiber has already been disposed.347
*348
* @returns nothing when the fiber is still active.349
* @throws {CordisError} `INACTIVE_EFFECT` when the fiber's uid has been cleared.350
*/351
assertActive() {352
if (this.uid !== null) return353
throw new CordisError('INACTIVE_EFFECT')354
}356
private _execute<T>(runner: EffectRunner<T>) {357
const oldEpoch = runner.epoch358
return composeError((info) => {359
const safeCollect = (dispose: void | Disposable) => {360
if (typeof dispose === 'function') {361
runner.collect(dispose)362
} else if (!isNullable(dispose)) {363
throw new TypeError('Invalid effect')364
}365
}366
const effect: Effect = runner.execute.call(this)367
if (typeof effect === 'function') {368
return runner.collect(effect)369
} else if (isNullable(effect)) {370
// return371
} else if (!isObject(effect)) {372
throw new TypeError('Invalid effect')373
} else if ('then' in effect) {374
return effect.then(safeCollect)375
} else if (Symbol.iterator in effect) {376
info.error = new Error()377
const iter = effect[Symbol.iterator]()378
while (true) {379
const result = iter.next()380
safeCollect(result.value)381
if (result.done) return382
}383
} else if (Symbol.asyncIterator in effect) {384
const iter = effect[Symbol.asyncIterator]()385
return (async () => {386
// force async stack trace387
await Promise.resolve()388
info.error = new Error()389
while (true) {390
if (runner.epoch !== oldEpoch) return391
const result = await iter.next()392
safeCollect(result.value)393
if (result.done) return394
}395
})()396
} else {397
throw new TypeError('Invalid effect')398
}399
}, runner.getOuterStack)400
}402
/**403
* Register a cleanup-aware effect on this fiber.404
*405
* `execute` runs immediately; the disposers it produces are collected and406
* run (in reverse order) either when the returned disposer is called or407
* when the fiber unloads, whichever comes first. Calling the disposer twice408
* is a no-op. Throws `CordisError('INACTIVE_EFFECT')` if the fiber is409
* already disposed, and `TypeError` if `execute` returns an invalid shape.410
*411
* @param execute — the effect body; see {@link Effect} for accepted shapes.412
* @param label — effect label shown in `getEffects()` diagnostics.413
* @returns a disposer that tears the effect down and settles once done.414
*/415
effect(execute: () => SyncEffect, label?: string): Disposable<Promise<void>>416
/** Same as above for async effects; the disposer is also awaitable. */417
effect(execute: () => Effect, label?: string): AsyncDisposable<Promise<void>>418
effect(execute: () => Effect, label = 'anonymous'): any {419
this.assertActive()420
if (this.state === FiberState.UNLOADING) {421
throw new CordisError('INACTIVE_EFFECT')422
}424
const disposables: Disposable[] = []425
let disposing = false426
let disposalTask: void | Promise<void>427
const dispose = () => {428
if (disposing) return disposalTask429
disposing = true430
let task!: void | Promise<void>431
for (const disposable of disposables.splice(0).reverse()) {432
if (task) {433
task = task.then(() => runDisposable(disposable))434
} else {435
const result = runDisposable(disposable)436
if (isObject(result) && 'then' in result) {437
task = result as any438
}439
}440
}441
return disposalTask = task442
}444
const meta: EffectMeta = { label, children: [] }445
const runner: EffectRunner<boolean> = {446
execute,447
epoch: true,448
collect: (dispose) => {449
disposables.push(dispose)450
this._disposables.delete(dispose)451
if (dispose[symbols.effect]) {452
meta.children.push(dispose[symbols.effect])453
}454
},455
getOuterStack: buildOuterStack(),456
}458
let task: void | Promise<void>459
let executing = true460
let resolveSetup: (() => void) | undefined461
let rejectSetup: ((reason: unknown) => void) | undefined462
let setupBarrier: Promise<void> | undefined463
let setupFailed = false464
let inFlight: void | Promise<void>465
let removeWrapper = () => false467
const waitForSetup = () => {468
setupBarrier ??= new Promise<void>((resolve, reject) => {469
resolveSetup = resolve470
rejectSetup = reject471
})472
return setupBarrier473
}475
const disposeAfter = (setup: PromiseLike<void>) => {476
return Promise.resolve(setup).then(477
() => dispose(),478
async (reason) => {479
await dispose()480
throw reason481
},482
)483
}485
const finalizeDisposal = (callback: () => void | Promise<void>) => {486
let result: void | Promise<void>487
try {488
result = callback()489
} catch (error) {490
removeWrapper()491
throw error492
}493
if (isObject(result) && 'then' in result) {494
const pending = Promise.resolve(result).finally(() => {495
removeWrapper()496
if (inFlight === pending) inFlight = undefined497
})498
return inFlight = pending499
}500
removeWrapper()501
return result502
}504
const wrapper = defineProperty(() => {505
// A synchronous setup failure can race an owner unload that already506
// captured this wrapper but has not invoked it yet. The failed effect is507
// never returned publicly, so let that internal caller await rollback.508
if (!runner.epoch) return setupFailed ? inFlight : undefined509
runner.epoch = false510
return finalizeDisposal(() => {511
if (executing) return disposeAfter(waitForSetup())512
return task ? disposeAfter(task) : dispose()513
})514
}, symbols.effect, meta) as AsyncDisposable515
effectInertia.set(wrapper, () => inFlight)517
// Make the effect visible to a reentrant owner unload before execute()518
// runs any plugin code. Async teardown stays owner-visible until it519
// settles, allowing an outer effect to join cleanup another caller began.520
removeWrapper = this._disposables.push(wrapper)521
try {522
task = this._execute(runner)523
} catch (reason) {524
executing = false525
setupFailed = true526
runner.epoch = false527
let cleanup: void | Promise<void>528
try {529
cleanup = finalizeDisposal(dispose)530
} finally {531
rejectSetup?.(reason)532
}533
if (isObject(cleanup) && 'then' in cleanup) {534
cleanup.catch(error => this.ctx.logger.error(error))535
}536
throw reason537
}538
executing = false539
if (setupBarrier) {540
Promise.resolve(task).then(resolveSetup, rejectSetup)541
}543
// prevent unhandled rejection — both from `task` itself and from the544
// disposer chain if it fails to settle cleanly.545
task?.catch(() => {546
if (!runner.epoch) return dispose()547
return finalizeDisposal(dispose)548
}).catch((error) => this.ctx.logger.error(error))550
const disposeAsync = () => {551
if (!runner.epoch) return552
runner.epoch = false553
return finalizeDisposal(dispose)554
}555
wrapper.then = async (onFulfilled, onRejected) => {556
return Promise.resolve(task)557
.then(() => disposeAsync)558
.then(onFulfilled, onRejected)559
}560
return wrapper561
}563
/**564
* Return metadata for currently registered effects.565
*566
* @returns one {@link EffectMeta} tree per labeled live effect.567
*/568
getEffects() {569
return [...this._disposables]570
.map<EffectMeta>(dispose => dispose[symbols.effect])571
.filter(Boolean)572
}574
private _getState() {575
if (this.uid === null) return FiberState.DISPOSED576
if (this._error) return FiberState.FAILED577
if (this._runner.epoch !== INACTIVE) return FiberState.ACTIVE578
return FiberState.PENDING579
}581
private _updateState(callback: () => void | FiberState) {582
const oldState = this.state583
this.state = callback() ?? this._getState()584
if (oldState === this.state) return585
// FIXME internal/fiber-info586
this.context.emit('internal/status', this, oldState)588
// only notify changes between ACTIVE and NON-ACTIVE states589
if (oldState !== FiberState.ACTIVE && this.state !== FiberState.ACTIVE) return590
for (const key of Reflect.ownKeys(this.ctx.reflect.store)) {591
const impl = this.ctx.reflect.store[key as symbol]592
if (impl.fiber !== this) continue593
this.ctx.reflect.notify([impl.name])594
}595
}597
_checkImpl(name: string) {598
const impl = this.ctx.reflect._getImpl(name, true)599
if (!impl) return delete this._store[name]600
try {601
if (impl.check && !impl.check.call(getTraceable(this.ctx, impl.value))) {602
return delete this._store[name]603
}604
} catch (error) {605
impl.fiber.ctx.logger.error(error)606
return delete this._store[name]607
}608
this._store[name] = impl609
}611
_refresh() {612
let epoch: string | boolean = false613
epoch = ''614
for (const name of Object.keys(this.inject)) {615
const impl = this._store[name]616
if (!impl) {617
epoch = INACTIVE618
break619
}620
epoch += ':' + impl.fiber.uid621
}622
this._setEpoch(epoch)623
}625
private _setEpoch(epoch: string) {626
const oldEpoch = this._runner.epoch627
if (epoch === oldEpoch) return628
this._runner.epoch = epoch629
if (this.inertia) return630
this._updateState(() => {631
if (epoch !== INACTIVE && oldEpoch === INACTIVE) {632
this.inertia = this._reload()633
return FiberState.LOADING634
} else {635
this.inertia = this._unload()636
return FiberState.UNLOADING637
}638
})639
}641
private _resolveConfig(config: any) {642
config = this.context.waterfall(this, 'internal/config', config, () => config)643
return this.runtime ? resolveConfig(this.runtime, config) : config644
}646
private async _reload() {647
this.store = { ...this._store }648
const oldEpoch = this._runner.epoch649
try {650
await Promise.resolve()651
// A disposer queued before this checkpoint may already have invalidated652
// the load. Do not run plugin code for a stale epoch; the state update653
// below will drain any effects collected while the fiber was PENDING.654
if (this._runner.epoch === oldEpoch) {655
this.config = this._resolveConfig(this._config)656
await this._execute(this._runner)657
this._error = undefined658
}659
} catch (reason) {660
// impl guarantees that the error is non-null (?)661
this.ctx.logger.error(reason)662
this._error = reason663
this._runner.epoch = INACTIVE664
}665
this._updateState(() => {666
if (this._runner.epoch === oldEpoch) {667
this.inertia = undefined668
} else {669
this.inertia = this._unload()670
return FiberState.UNLOADING671
}672
})673
}675
private async _unload() {676
await Promise.all(this._disposables.clear().map(async (dispose) => {677
try {678
await composeError(async (info) => {679
await Promise.resolve()680
info.error = new Error()681
await runDisposable(dispose)682
}, this._runner.getOuterStack)683
} catch (reason) {684
this.ctx.logger.error(reason)685
}686
}))687
this.store = undefined688
this._updateState(() => {689
if (this._runner.epoch === INACTIVE) {690
this.inertia = undefined691
} else {692
this.inertia = this._reload()693
return FiberState.LOADING694
}695
})696
}698
/**699
* Wait for current lifecycle work and rethrow startup errors.700
*701
* @returns this fiber, once it has settled into a stable state.702
* @throws the config-validation or plugin-startup error, if any.703
*/704
async await() {705
while (this.inertia) {706
await this.inertia707
}708
if (this._error) throw this._error709
return this710
}712
/**713
* Dispose and immediately reload this plugin with its current config.714
*715
* @returns a promise resolving once the reload settled.716
* @throws {CordisError} `INACTIVE_EFFECT` when the fiber is already disposed.717
*/718
async restart() {719
this.assertActive()720
this._setEpoch(INACTIVE)721
this._refresh()722
await this.await()723
}725
/**726
* Validate and apply new config, then restart the plugin.727
*728
* Runs the `internal/update` waterfall first, so update hooks (and HMR)729
* can veto or replace the restart.730
*731
* @param config — the new raw config; validated before anything restarts.732
* @param noSave — hint for persistence hooks not to write the change back.733
* @returns nothing; the restart runs behind the `internal/update` waterfall.734
* @throws {ValidationError} when the new config fails validation.735
*/736
update(config: any, noSave = false) {737
this.assertActive()738
this._config = config739
if (this.state !== FiberState.ACTIVE) {740
// Config resolution may access injected services, so defer it until the741
// fiber can activate.742
this._error = undefined743
this._setEpoch(INACTIVE)744
this._refresh()745
return746
}747
config = this._resolveConfig(config)748
this.context.waterfall(this, 'internal/update', config, noSave, () => {749
this.config = config750
this._error = undefined751
return this.restart()752
})753
}754
}