import type { Readable } from "node:stream" //////////////////////////////////////////////////////////////////////////////////////////////////// // Base Task // //////////////////////////////////////////////////////////////////////////////////////////////////// export interface TaskResolveContext { runnerReturn: RunnerReturn resolve: (value: void | PromiseLike) => void } export type TaskResolver = (context: TaskResolveContext) => void export const defaultTaskResolver: TaskResolver = (context) => { context.resolve() } export type Task = () => Promise export interface DefineTaskOptions { name: string runner: () => RunnerReturn resolver?: TaskResolver | undefined } export const defineTask = (options: DefineTaskOptions): Task => { return async () => { const { name, runner, resolver = defaultTaskResolver } = options console.log(`\n\n[Pipeline] start ${name}...\n\n`) const runnerReturn = runner() return await new Promise((resolve) => { // print status to console when task resolved const internalResolve = (value: void | PromiseLike): void => { console.log(`\n\n[Pipeline] end ${name}...\n\n`) resolve(value) } resolver({ runnerReturn, resolve: internalResolve, }) }) } } //////////////////////////////////////////////////////////////////////////////////////////////////// // Stream Task // //////////////////////////////////////////////////////////////////////////////////////////////////// export type StreamCondition = (data: string) => boolean export const include = (keyword: string): StreamCondition => { return (data) => data.includes(keyword) } export const occurTimes = (keyword: string, times: number): StreamCondition => { let count = 0 return (data) => { if (data.includes(keyword)) { count = count + 1 } return count === times } } export const sequencialOccur = (keywords: string[]): StreamCondition => { let index = 0 return (data) => { const keyword = keywords[index] if (keyword !== undefined && data.includes(keyword)) { index = index + 1 } return index === keywords.length } } export interface ReadableReturn { stdout: Readable } export interface StreamTaskOptions { name: string runner: () => ReadableReturn itSuccessWhen: StreamCondition[] } export const defineStreamTask = (options: StreamTaskOptions): Task => { const { name, runner, itSuccessWhen } = options const successMarks = Array.from({ length: itSuccessWhen.length }).fill(false) let resolved = false return defineTask({ name, runner, resolver: (context) => { context.runnerReturn.stdout.on("data", (data) => { console.log(String(data)) itSuccessWhen.forEach((itSuccess, index) => { if (itSuccess(String(data))) { successMarks[index] = true } }) if (!successMarks.includes(false) && !resolved) { context.resolve() resolved = true } }) }, }) } //////////////////////////////////////////////////////////////////////////////////////////////////// // Sleep Task // //////////////////////////////////////////////////////////////////////////////////////////////////// export interface SleepTaskOptions { name: string duration: number } export const defineSleepTask = (options: SleepTaskOptions): Task => { const { name, duration } = options return defineTask({ name, runner: () => { // }, resolver: (context) => { setTimeout(() => { context.resolve() }, duration) }, }) }