import { MAX_SAFE_INTEGER } from './alias'; import { Inits } from './type'; import { List } from './list'; import { singleton } from './function'; import { causeAsyncException } from './exception'; class Node implements List.Node { constructor( public value: T, ) { } public next?: this = undefined; public prev?: this = undefined; } export interface Observer { monitor(namespace: Readonly>, listener: Monitor, options?: ObserverOptions): () => void; on(namespace: Readonly, listener: Subscriber, options?: ObserverOptions): () => void; off(namespace: Readonly, listener: Subscriber): void; off(namespace: Readonly>): void; once(namespace: Readonly, listener: Subscriber): () => void; } export interface ObserverOptions { once?: boolean; } export interface Publisher { emit(namespace: Readonly, data: D, tracker?: (data: D, results: R[]) => void): void; emit(this: Publisher, namespace: Readonly, data?: D, tracker?: (data: D, results: R[]) => void): void; reflect(namespace: Readonly>, data: D): R[]; reflect(this: Publisher, namespace: Readonly>, data?: D): R[]; } export type Monitor = (data: D, namespace: Readonly>) => void; export type Subscriber = (data: D, namespace: Readonly>) => R; class ListenerNode { constructor( public readonly name: N[number], public readonly parent?: ListenerNode, ) { } public mid = 0; public sid = 0; public readonly monitors = new List>>(); public readonly subscribers = new List>>(); public readonly index = new Map>(); public readonly children = new List>>(); public reset(listeners: List>>): void { switch (listeners) { case this.monitors: this.mid = 0; for (let node = listeners.head, i = listeners.length; node && i--; node = node.next) { node.value.id = ++this.mid; } return; case this.subscribers: this.sid = 0; for (let node = listeners.head, i = listeners.length; node && i--; node = node.next) { node.value.id = ++this.sid; } return; default: throw new Error('Unreachable'); } } public clear(disposable = false): boolean { const { monitors, subscribers, index, children } = this; const stack = []; for (let child = children.head, i = children.length; child && i--;) { if (child.value.clear(true)) { const next = child.next; disposable ? stack.push(child.value.name) : index.delete(child.value.name); children.delete(child); child = next; } else { child = child.next; } } if (children.length) while (stack.length) { index.delete(stack.pop()); } subscribers.clear(); return monitors.length === 0 && children.length === 0; } } export type ListenerItem = | MonitorItem | SubscriberItem; interface MonitorItem { id: number; readonly type: ListenerType.Monitor; readonly namespace: Readonly>; readonly listener: Monitor; readonly options: ObserverOptions; } interface SubscriberItem { id: number; readonly type: ListenerType.Subscriber; readonly namespace: Readonly; readonly listener: Subscriber; readonly options: ObserverOptions; } const enum ListenerType { Monitor, Subscriber, } const enum SeekMode { Extensible, Breakable, Closest, } export interface ObservationOptions { readonly limit?: number; } export class Observation implements Observer, Publisher { constructor(opts?: ObservationOptions) { this.limit = opts?.limit ?? 10; } private readonly node = new ListenerNode(undefined); private readonly limit: number; public monitor(namespace: Readonly>, monitor: Monitor, options: ObserverOptions = {}): () => void { if (typeof monitor !== 'function') throw new Error(`Spica: Observation: Invalid listener: ${monitor}`); const node = this.seek(namespace, SeekMode.Extensible); const monitors = node.monitors; if (monitors.length === this.limit) throw new Error(`Spica: Observation: Exceeded max listener limit`); node.mid === MAX_SAFE_INTEGER && node.reset(monitors); const inode = monitors.push(new Node({ id: ++node.mid, type: ListenerType.Monitor, namespace, listener: monitor, options, })); return singleton(() => void monitors.delete(inode)); } public on(namespace: Readonly, subscriber: Subscriber, options: ObserverOptions = {}): () => void { if (typeof subscriber !== 'function') throw new Error(`Spica: Observation: Invalid listener: ${subscriber}`); const node = this.seek(namespace, SeekMode.Extensible); const subscribers = node.subscribers; if (subscribers.length === this.limit) throw new Error(`Spica: Observation: Exceeded max listener limit`); node.sid === MAX_SAFE_INTEGER && node.reset(subscribers); const inode = subscribers.push(new Node({ id: ++node.sid, type: ListenerType.Subscriber, namespace, listener: subscriber, options, })); return singleton(() => void subscribers.delete(inode)); } public once(namespace: Readonly, subscriber: Subscriber): () => void { return this.on(namespace, subscriber, { once: true }); } public off(namespace: Readonly, subscriber: Subscriber): void; public off(namespace: Readonly>): void; public off(namespace: Readonly>, subscriber?: Subscriber): void { if (subscriber) { const list = this.seek(namespace, SeekMode.Breakable)?.subscribers; const node = list?.find(node => node.value.listener === subscriber); assert(node?.next || node?.prev || list?.head === node); node && list?.delete(node); } else { void this.seek(namespace, SeekMode.Breakable)?.clear(); } } public emit(namespace: Readonly, data: D, tracker?: (data: D, results: R[]) => void): void public emit(this: Publisher, namespace: Readonly, data?: D, tracker?: (data: D, results: R[]) => void): void public emit(namespace: Readonly, data: D, tracker?: (data: D, results: R[]) => void): void { this.drain(namespace, data, tracker); } public reflect(namespace: Readonly>, data: D): R[] public reflect(this: Publisher, namespace: Readonly>, data?: D): R[] public reflect(namespace: Readonly>, data: D): R[] { let results!: R[]; this.emit(namespace as N, data, (_, r) => results = r); assert(results); return results; } private relaies!: WeakSet>; public relay(source: Observer): () => void { this.relaies ??= new WeakSet(); assert(!this.relaies.has(source)); if (this.relaies.has(source)) throw new Error(`Spica: Observation: Relay source is already registered`); this.relaies.add(source); return source.monitor([] as Inits, (data, namespace) => void this.emit(namespace as N, data)); } public refs(namespace: Readonly>): ListenerItem[] { const node = this.seek(namespace, SeekMode.Breakable); if (node === undefined) return []; return this.listenersBelow(node) .reduce[]>((acc, listeners) => void listeners.foldl((_, node) => acc.push(node.value), 0) || acc , []); } private drain(namespace: Readonly, data: D, tracker?: (data: D, results: R[]) => void): void { let node = this.seek(namespace, SeekMode.Breakable); const results: R[] = []; for (let lists = node ? this.listenersBelow(node, ListenerType.Subscriber) : [], i = 0; i < lists.length; ++i) { const items = lists[i]; if (items.length === 0) continue; const recents: Node>[] = []; const max = items.last!.value.id; let min = 0; let prev: typeof recents[0] | undefined; for (let node = items.head; node && min < node.value.id && node.value.id <= max;) { min = node.value.id; const item = node.value; item.options.once && items.delete(node); try { const result = item.listener(data, namespace); tracker && results.push(result); } catch (reason) { causeAsyncException(reason); } (node.next !== undefined || node.prev !== undefined || items.head === node) && recents.push(node); node = node.next ?? prev?.next ?? rollback(recents, item => item.next) ?? items.head; prev = node?.prev; } } node ??= this.seek(namespace, SeekMode.Closest); for (let lists = this.listenersAbove(node, ListenerType.Monitor), i = 0; i < lists.length; ++i) { const items = lists[i]; if (items.length === 0) continue; const recents: Node>[] = []; const max = items.last!.value.id; let min = 0; let prev: typeof recents[0] | undefined; for (let node = items.head; node && min < node.value.id && node.value.id <= max;) { min = node.value.id; const item = node.value; item.options.once && items.delete(node); try { item.listener(data, namespace); } catch (reason) { causeAsyncException(reason); } (node.next !== undefined || node.prev !== undefined || items.head === node) && recents.push(node); node = node.next ?? prev?.next ?? rollback(recents, item => item.next) ?? items.head; prev = node?.prev; } } if (tracker) { try { tracker(data, results); } catch (reason) { causeAsyncException(reason); } } } private seek(namespace: Readonly>, mode: SeekMode.Extensible | SeekMode.Closest): ListenerNode; private seek(namespace: Readonly>, mode: SeekMode): ListenerNode | undefined; private seek(namespace: Readonly>, mode: SeekMode): ListenerNode | undefined { let node = this.node; for (let i = 0; i < namespace.length; ++i) { const name = namespace[i]; const { index, children } = node; let child = index.get(name); if (child === undefined) { switch (mode) { case SeekMode.Breakable: return; case SeekMode.Closest: return node; } child = new ListenerNode(name, node); index.set(name, child); children.push(new Node(child)); } node = child; } return node; } private listenersAbove({ parent, monitors }: ListenerNode, type: ListenerType.Monitor): List>>[]; private listenersAbove({ parent, monitors }: ListenerNode): List>>[] { const acc = [monitors]; while (parent) { acc.push(parent.monitors); parent = parent.parent; } return acc; } private listenersBelow(node: ListenerNode): List>>[]; private listenersBelow(node: ListenerNode, type: ListenerType.Subscriber): List>>[]; private listenersBelow(node: ListenerNode, type?: ListenerType.Subscriber): List>>[] { return this.listenersBelow$(node, type, [])[0]; } private listenersBelow$( { monitors, subscribers, index, children }: ListenerNode, type: ListenerType.Subscriber | undefined, acc: List>>[], ): readonly [List>>[], number] { switch (type) { case ListenerType.Subscriber: acc.push(subscribers); break; default: acc.push(monitors, subscribers); } let count = 0; for (let child = children.head, i = children.length; child && i--;) { const cnt = this.listenersBelow$(child.value, type, acc)[1]; count += cnt; if (cnt === 0) { const next = child.next; index.delete(child.value.name); children.delete(child); child = next; } else { child = child.next; } } return [acc, monitors.length + subscribers.length + count]; } } function rollback(array: T[], matcher: (value: T) => unknown): T | undefined { for (let i = array.length; i--;) { if (matcher(array[i])) return array[i]; array.pop(); } }