import { scheduleMacroTask } from "#Source/basic/index.ts" import type { LoggerFriendly, LoggerFriendlyOptions } from "#Source/log/index.ts" import { Logger } from "#Source/log/index.ts" /** * @description 表示 Tube 中的订阅回调。 */ export type Subscriber = ((data: D) => void) | ((data: D) => Promise) /** * @description 表示用于移除订阅关系的函数。 */ export type Unsubscribe = () => void /** * @description `open` 事件的订阅配置。 */ export interface OpenEventSubscribeOptions { subscriber: Subscriber } interface OpenEventSubscribeStates { subscriber: OpenEventSubscribeOptions["subscriber"] unsubscribe: Unsubscribe } /** * @description `close` 事件的订阅配置。 */ export interface CloseEventSubscribeOptions { subscriber: Subscriber } interface CloseEventSubscribeStates { subscriber: CloseEventSubscribeOptions["subscriber"] unsubscribe: Unsubscribe } /** * @description `start` 事件的订阅配置。 */ export interface StartEventSubscribeOptions { subscriber: Subscriber } interface StartEventSubscribeStates { subscriber: StartEventSubscribeOptions["subscriber"] unsubscribe: Unsubscribe } /** * @description `end` 事件的订阅配置。 */ export interface EndEventSubscribeOptions { subscriber: Subscriber } interface EndEventSubscribeStates { subscriber: EndEventSubscribeOptions["subscriber"] unsubscribe: Unsubscribe } /** * @description `error` 事件的订阅配置。 */ export interface ErrorEventSubscribeOptions { subscriber: Subscriber } interface ErrorEventSubscribeStates { subscriber: ErrorEventSubscribeOptions["subscriber"] unsubscribe: Unsubscribe } /** * @description `wet` 事件的订阅配置。 */ export interface WetEventSubscribeOptions { subscriber: Subscriber } interface WetEventSubscribeStates { subscriber: WetEventSubscribeOptions["subscriber"] unsubscribe: Unsubscribe } /** * @description 数据事件的订阅配置。 */ export interface DataEventSubscribeOptions { subscriber: Subscriber /** * @description 是否重放历史记录。 * * @default DataOptions.replayHistory */ replayHistory?: boolean } interface DataEventSubscribeStates { subscriber: DataEventSubscribeOptions["subscriber"] replayHistory: boolean unsubscribe: Unsubscribe } interface TubeRunningStates { /** * @description True: 可以进数据。 * false: 不可以进数据。 */ isOpen: boolean /** * @description True: 数据进入启动。 * false: 数据进入停止。 */ isStart: boolean /** * @description True: 发生错误。 * false: 未发生错误。 */ isError: boolean /** * @description True: 数据为空。 * false: 数据不为空。 */ isEmpty: boolean /** * @description True: 已打开过。 * false: 未打开过。 */ hasOpened: boolean /** * @description True: 已开始过。 * false: 未开始过。 */ hasStarted: boolean /** * @description True: 已结束过。 * false: 未结束过。 */ hasEnded: boolean /** * @description True: 已关闭过。 * false: 未关闭过。 */ hasClosed: boolean } interface DataReplayStates { dataHistory: D[] } /** * @description `Tube` 的配置项。 */ export interface TubeOptions extends LoggerFriendlyOptions { /** * @default true */ autoStartOnOpen?: boolean /** * @default true */ autoOpenOnStart?: boolean /** * @default true */ autoEndOnClose?: boolean /** * @default true */ autoCloseOnEnd?: boolean /** * @default true */ autoStartOnData?: boolean /** * @default true */ autoOpenOnData?: boolean /** * @default true */ autoOpenOnError?: boolean /** * @default true */ autoEndOnError?: boolean /** * @default true */ autoCloseOnError?: boolean /** * @description 保留多少个历史记录。 * * @default 3 */ historyCount: number /** * @description 是否重放历史记录。 * * @default false */ replayHistory: boolean } /** * @description 表示一条带生命周期、错误语义和历史回放能力的数据通道。 * 一般的事件顺序是:Open - Data - Data - ... - Data - Data - Close。 * Error 可能出现在 Open 和 Close 之间的任何位置。 */ export class Tube implements LoggerFriendly { logger: Logger // 基础配置 protected autoStartOnOpen: boolean protected autoOpenOnStart: boolean protected autoEndOnClose: boolean protected autoCloseOnEnd: boolean protected autoStartOnData: boolean protected autoOpenOnData: boolean protected autoOpenOnError: boolean protected autoEndOnError: boolean protected autoCloseOnError: boolean protected historyCount: number protected replayHistory: boolean // 整体状态 protected runningStates: TubeRunningStates // 事件订阅 protected openEventSubscribeStatesMap: Map, OpenEventSubscribeStates> protected closeEventSubscribeStatesMap: Map, CloseEventSubscribeStates> protected startEventSubscribeStatesMap: Map, StartEventSubscribeStates> protected endEventSubscribeStatesMap: Map, EndEventSubscribeStates> protected errorEventSubscribeStatesMap: Map, ErrorEventSubscribeStates> protected wetEventSubscribeStatesMap: Map, WetEventSubscribeStates> protected dataEventSubscribeStatesMap: Map, DataEventSubscribeStates> // 其它 protected errorHistory: E[] protected dataHistory: D[] protected dataReplayStatesMap: Map, DataReplayStates> /** * @description 创建一个 Tube 实例。 */ constructor(options: TubeOptions) { this.logger = Logger.fromOptions(options).setDefaultName("Tube") this.autoStartOnOpen = options.autoStartOnOpen ?? true this.autoOpenOnStart = options.autoOpenOnStart ?? true this.autoEndOnClose = options.autoEndOnClose ?? true this.autoCloseOnEnd = options.autoCloseOnEnd ?? true this.autoStartOnData = options.autoStartOnData ?? true this.autoOpenOnData = options.autoOpenOnData ?? true this.autoOpenOnError = options.autoOpenOnError ?? true this.autoEndOnError = options.autoEndOnError ?? true this.autoCloseOnError = options.autoCloseOnError ?? true this.historyCount = options.historyCount ?? 3 this.replayHistory = options.replayHistory ?? false if (this.historyCount < 1) { throw new Error("historyCount must be greater than or equal to 1") } this.runningStates = { isOpen: false, isStart: false, isError: false, isEmpty: true, hasOpened: false, hasStarted: false, hasEnded: false, hasClosed: false, } this.openEventSubscribeStatesMap = new Map() this.closeEventSubscribeStatesMap = new Map() this.startEventSubscribeStatesMap = new Map() this.endEventSubscribeStatesMap = new Map() this.errorEventSubscribeStatesMap = new Map() this.wetEventSubscribeStatesMap = new Map() this.dataEventSubscribeStatesMap = new Map() this.errorHistory = [] this.dataHistory = [] this.dataReplayStatesMap = new Map() } /** * @description 安全地获取最近一条数据;若当前还没有任何数据则抛出错误。 * 当数据列表为空时,此方法会抛出错误,请妥善处理。 * 如果不希望错误发生,可以通过 {@link isWet} 方法判断数据列表是否为空。 */ getLatestDataOrThrow(): D { const latestData = this.dataHistory.at(-1) if (latestData === undefined) { throw new Error("latestData is undefined") } return latestData } /** * @description 安全地获取最近一次错误;若当前还没有错误记录则抛出错误。 * 当错误列表为空时,此方法会抛出错误,请妥善处理。 * 如果不希望错误发生,可以通过 {@link isError} 方法判断错误列表是否为空。 */ getLatestErrorOrThrow(): E { const latestError = this.errorHistory.at(-1) if (latestError === undefined) { throw new Error("latestError is undefined") } return latestError } /** * @description 判断 Tube 当前是否处于打开状态。 */ isOpen(): boolean { return this.runningStates.isOpen === true } /** * @description 判断 Tube 是否曾经被打开过。 */ hasOpened(): boolean { return this.runningStates.hasOpened === true } /** * @description 打开 Tube,并在需要时自动启动。 */ async open(): Promise { if (this.isOpen() === true) { return } if (this.hasClosed() === true) { return } this.runningStates.isOpen = true this.runningStates.hasOpened = true this.triggerOpenEventForAll() if (this.autoStartOnOpen === true) { await this.start() } } /** * @description 判断 Tube 当前是否处于关闭状态。 */ isClose(): boolean { return this.runningStates.isOpen === false } /** * @description 判断 Tube 是否曾经被关闭过。 */ hasClosed(): boolean { return this.runningStates.hasClosed === true } /** * @description 关闭 Tube,并在需要时自动结束。 */ async close(): Promise { if (this.isClose() === true) { return } if (this.autoEndOnClose === true) { await this.end() } this.runningStates.isOpen = false this.runningStates.hasClosed = true this.triggerCloseEventForAll() } /** * @description 判断 Tube 当前是否处于启动状态。 */ isStart(): boolean { return this.runningStates.isStart === true } /** * @description 判断 Tube 是否曾经进入过启动状态。 */ hasStarted(): boolean { return this.runningStates.hasStarted === true } /** * @description 启动 Tube,并在需要时自动打开。 */ async start(): Promise { if (this.isStart() === true) { return } if (this.autoOpenOnStart === true) { await this.open() } if (this.isOpen() === false) { return } this.runningStates.isStart = true this.runningStates.hasStarted = true this.triggerStartEventForAll() } /** * @description 判断 Tube 当前是否处于结束状态。 */ isEnd(): boolean { return this.runningStates.isStart === false } /** * @description 判断 Tube 是否曾经结束过。 */ hasEnded(): boolean { return this.runningStates.hasEnded === true } /** * @description 结束 Tube,并在需要时自动关闭。 */ async end(): Promise { if (this.isEnd() === true) { return } if (this.isOpen() === false) { return } this.runningStates.isStart = false this.runningStates.hasEnded = true this.triggerEndEventForAll() if (this.autoCloseOnEnd === true) { await this.close() } } /** * @description 判断 Tube 是否已经记录过错误。 */ isError(): boolean { return this.runningStates.isError === true } protected async error(): Promise { if (this.isOpen() === false) { if (this.autoOpenOnError === true) { await this.open() } } if (this.isOpen() === false) { return } // NOTE: 允许在错误状态下继续推送数据,且允许发生多次错误。 // if (this.isError() === true) { // return // } this.runningStates.isError = true this.triggerErrorEventForAll() if (this.autoEndOnError === true) { await this.end() } if (this.autoCloseOnError === true) { await this.close() } } /** * @description 判断 Tube 是否已经接收过至少一条数据。 */ isWet(): boolean { return this.runningStates.isEmpty === false } // oxlint-disable-next-line require-await protected async wet(): Promise { if (this.isWet() === true) { return } this.runningStates.isEmpty = false this.triggerWetEventForAll() } protected async data(): Promise { if (this.isOpen() === false) { if (this.autoOpenOnData === true) { await this.open() } else { return } } if (this.isOpen() === false) { return } if (this.isStart() === false) { if (this.autoStartOnData === true) { await this.start() } else { return } } if (this.isStart() === false) { return } } protected triggerOpenEventForOne(subscriber: OpenEventSubscribeOptions["subscriber"]): void { scheduleMacroTask({ task: async () => await subscriber() }) } protected triggerOpenEventForAll(): void { this.openEventSubscribeStatesMap.forEach((subscribeStates) => { scheduleMacroTask({ task: async () => await subscribeStates.subscriber() }) }) } /** * @description 订阅 Tube 的 `open` 事件。 */ subscribeOpenEvent(options: OpenEventSubscribeOptions): Unsubscribe { const states = this.openEventSubscribeStatesMap.get(options.subscriber) if (states !== undefined) { return states.unsubscribe } const subscriber = options.subscriber const unsubscribe = (): void => { this.openEventSubscribeStatesMap.delete(subscriber) } this.openEventSubscribeStatesMap.set(subscriber, { subscriber, unsubscribe, }) return unsubscribe } /** * @description 取消指定的 `open` 事件订阅。 */ unsubscribeOpenEvent(subscriber: OpenEventSubscribeOptions["subscriber"]): void { const subscribeStates = this.openEventSubscribeStatesMap.get(subscriber) if (subscribeStates === undefined) { return } subscribeStates.unsubscribe() } protected triggerCloseEventForOne(subscriber: CloseEventSubscribeOptions["subscriber"]): void { scheduleMacroTask({ task: async () => await subscriber() }) } protected triggerCloseEventForAll(): void { this.closeEventSubscribeStatesMap.forEach((subscribeStates) => { scheduleMacroTask({ task: async () => await subscribeStates.subscriber() }) }) } /** * @description 订阅 Tube 的 `close` 事件。 */ subscribeCloseEvent(options: CloseEventSubscribeOptions): Unsubscribe { const states = this.closeEventSubscribeStatesMap.get(options.subscriber) if (states !== undefined) { return states.unsubscribe } const subscriber = options.subscriber const unsubscribe = (): void => { this.closeEventSubscribeStatesMap.delete(subscriber) } this.closeEventSubscribeStatesMap.set(subscriber, { subscriber, unsubscribe, }) return unsubscribe } /** * @description 取消指定的 `close` 事件订阅。 */ unsubscribeCloseEvent(subscriber: CloseEventSubscribeOptions["subscriber"]): void { const subscribeStates = this.closeEventSubscribeStatesMap.get(subscriber) if (subscribeStates === undefined) { return } subscribeStates.unsubscribe() } protected triggerStartEventForOne(subscriber: StartEventSubscribeOptions["subscriber"]): void { scheduleMacroTask({ task: async () => await subscriber() }) } protected triggerStartEventForAll(): void { this.startEventSubscribeStatesMap.forEach((subscribeStates) => { scheduleMacroTask({ task: async () => await subscribeStates.subscriber() }) }) } /** * @description 订阅 Tube 的 `start` 事件。 */ subscribeStartEvent(options: StartEventSubscribeOptions): Unsubscribe { const states = this.startEventSubscribeStatesMap.get(options.subscriber) if (states !== undefined) { return states.unsubscribe } const subscriber = options.subscriber const unsubscribe = (): void => { this.startEventSubscribeStatesMap.delete(subscriber) } this.startEventSubscribeStatesMap.set(subscriber, { subscriber, unsubscribe, }) return unsubscribe } /** * @description 取消指定的 `start` 事件订阅。 */ unsubscribeStartEvent(subscriber: StartEventSubscribeOptions["subscriber"]): void { const subscribeStates = this.startEventSubscribeStatesMap.get(subscriber) if (subscribeStates === undefined) { return } subscribeStates.unsubscribe() } protected triggerEndEventForOne(subscriber: EndEventSubscribeOptions["subscriber"]): void { scheduleMacroTask({ task: async () => await subscriber() }) } protected triggerEndEventForAll(): void { this.endEventSubscribeStatesMap.forEach((subscribeStates) => { scheduleMacroTask({ task: async () => await subscribeStates.subscriber() }) }) } /** * @description 订阅 Tube 的 `end` 事件。 */ subscribeEndEvent(options: EndEventSubscribeOptions): Unsubscribe { const states = this.endEventSubscribeStatesMap.get(options.subscriber) if (states !== undefined) { return states.unsubscribe } const subscriber = options.subscriber const unsubscribe = (): void => { this.endEventSubscribeStatesMap.delete(subscriber) } this.endEventSubscribeStatesMap.set(subscriber, { subscriber, unsubscribe, }) return unsubscribe } /** * @description 取消指定的 `end` 事件订阅。 */ unsubscribeEndEvent(subscriber: EndEventSubscribeOptions["subscriber"]): void { const subscribeStates = this.endEventSubscribeStatesMap.get(subscriber) if (subscribeStates === undefined) { return } subscribeStates.unsubscribe() } protected triggerWetEventForOne(subscriber: WetEventSubscribeOptions["subscriber"]): void { scheduleMacroTask({ task: async () => await subscriber() }) } protected triggerWetEventForAll(): void { this.wetEventSubscribeStatesMap.forEach((subscribeStates) => { scheduleMacroTask({ task: async () => await subscribeStates.subscriber() }) }) } /** * @description 订阅 Tube 首次变为非空时触发的 `wet` 事件。 */ subscribeWetEvent(options: WetEventSubscribeOptions): Unsubscribe { const states = this.wetEventSubscribeStatesMap.get(options.subscriber) if (states !== undefined) { return states.unsubscribe } const subscriber = options.subscriber const unsubscribe = (): void => { this.wetEventSubscribeStatesMap.delete(subscriber) } this.wetEventSubscribeStatesMap.set(subscriber, { subscriber, unsubscribe, }) return unsubscribe } /** * @description 取消指定的 `wet` 事件订阅。 */ unsubscribeWetEvent(subscriber: WetEventSubscribeOptions["subscriber"]): void { const subscribeStates = this.wetEventSubscribeStatesMap.get(subscriber) if (subscribeStates === undefined) { return } subscribeStates.unsubscribe() } /** * @description 向 Tube 推送一个错误,并触发错误生命周期。 */ async pushError(error: E): Promise { this.errorHistory.push(error) await this.error() } protected triggerErrorEventForOne(subscriber: ErrorEventSubscribeOptions["subscriber"]): void { const latestErrorData = this.errorHistory.at(-1) if (latestErrorData === undefined) { throw new Error("latestErrorData is undefined") } scheduleMacroTask({ task: async () => await subscriber(latestErrorData) }) } protected triggerErrorEventForAll(): void { const latestErrorData = this.errorHistory.at(-1) if (latestErrorData === undefined) { throw new Error("latestErrorData is undefined") } this.errorEventSubscribeStatesMap.forEach((subscribeStates) => { scheduleMacroTask({ task: async () => await subscribeStates.subscriber(latestErrorData) }) }) } /** * @description 订阅 Tube 的 `error` 事件。 */ subscribeErrorEvent(options: ErrorEventSubscribeOptions): Unsubscribe { const states = this.errorEventSubscribeStatesMap.get(options.subscriber) if (states !== undefined) { return states.unsubscribe } const subscriber = options.subscriber const unsubscribe = (): void => { this.errorEventSubscribeStatesMap.delete(subscriber) } this.errorEventSubscribeStatesMap.set(subscriber, { subscriber, unsubscribe, }) return unsubscribe } /** * @description 取消指定的 `error` 事件订阅。 */ unsubscribeErrorEvent(subscriber: ErrorEventSubscribeOptions["subscriber"]): void { const subscribeStates = this.errorEventSubscribeStatesMap.get(subscriber) if (subscribeStates === undefined) { return } subscribeStates.unsubscribe() } /** * @description 向 Tube 推送一条数据,并在需要时自动打开、启动与变为非空。 */ async pushData(data: D): Promise { await this.data() this.dataHistory.push(data) if (this.dataHistory.length > this.historyCount) { this.dataHistory.shift() } const allDataEventSubscribeStatesList = Array.from(this.dataEventSubscribeStatesMap.values()) const replayingDataSubscriberList = Array.from(this.dataReplayStatesMap.keys()) const notReplayingDataEventSubscribeStatesList = allDataEventSubscribeStatesList.filter( (subscribeStates) => { return replayingDataSubscriberList.includes(subscribeStates.subscriber) === false }, ) // 对于不处于重放状态的监听器,直接触发 notReplayingDataEventSubscribeStatesList.forEach((subscribeStates) => { scheduleMacroTask({ task: async () => await subscribeStates.subscriber(data) }) }) // 对于处于重放状态的监听器,将数据推入重放队列 this.dataReplayStatesMap.forEach((replayStates) => { replayStates.dataHistory.push(data) }) if (this.isWet() === false) { await this.wet() } } protected startReplay(subscriber: Subscriber): void { const isReplaying = this.dataReplayStatesMap.has(subscriber) if (isReplaying === true) { return } const dataHistoryToReplay = [...this.dataHistory] this.dataReplayStatesMap.set(subscriber, { dataHistory: dataHistoryToReplay }) scheduleMacroTask({ task: async () => { while (dataHistoryToReplay.length !== 0) { const data = dataHistoryToReplay.shift()! await subscriber(data) } await this.stopReplay(subscriber) }, }) } // oxlint-disable-next-line require-await protected async stopReplay(subscriber: Subscriber): Promise { this.dataReplayStatesMap.delete(subscriber) } /** * @description 订阅 Tube 的数据事件,并可选地重放已有历史。 */ subscribeData(options: DataEventSubscribeOptions): Unsubscribe { const states = this.dataEventSubscribeStatesMap.get(options.subscriber) if (states !== undefined) { return states.unsubscribe } const subscriber = options.subscriber const replayHistory = options.replayHistory ?? this.replayHistory const unsubscribe = (): void => { this.dataEventSubscribeStatesMap.delete(subscriber) this.dataReplayStatesMap.delete(subscriber) } this.dataEventSubscribeStatesMap.set(subscriber, { subscriber, replayHistory, unsubscribe, }) if (replayHistory === true) { // NOTE: 开始 Replay 的方法必须同步调用,即第一时间将目标订阅者设为 Replay 状态,避免丢失数据。 this.startReplay(subscriber) } return unsubscribe } /** * @description 取消指定的数据订阅。 */ unsubscribeData(subscriber: Subscriber): void { const subscribeStates = this.dataEventSubscribeStatesMap.get(subscriber) if (subscribeStates === undefined) { return } subscribeStates.unsubscribe() } }