/** * EngineService — the workflow execution engine. * * Consumes stream messages from the router and dispatches them * to the appropriate activity handler (Trigger, Worker, Hook, …). * * Each section delegates to a purpose-specific module inside `engine/`. * Open the module when you need implementation detail; read this file * when you need the big picture. * * Lifecycle (maps to modules): * 1. INIT → init.ts (channel setup, router, config) * 2. VERSION → version.ts (app version resolution, caching) * 3. SCHEMA → schema.ts (activity lookup, handler factory) * 4. COMPILE → compiler.ts (YAML plan & deploy) * 5. REPORT → reporting.ts (stats, IDs, query resolution) * 6. DISPATCH → dispatch.ts (stream message → activity handler) * 7. COMPLETION → completion.ts (parent notify, cleanup, expiry) * 8. SIGNAL → signal.ts (webhook/timehook delivery, fan-out) * 9. PUB/SUB → pubsub.ts (topic messaging, subscriptions) * 10. STATE → state.ts (job state retrieval, export) */ import { ExporterService } from '../exporter'; import { ILogger } from '../logger'; import { Router } from '../router'; import { SearchService } from '../search'; import { StoreService } from '../store'; import { StreamService } from '../stream'; import { SubService } from '../sub'; import { TaskService } from '../task'; import { AppVID } from '../../types/app'; import { ActivityType } from '../../types/activity'; import { CacheMode } from '../../types/cache'; import { ExportOptions, JobExport } from '../../types/exporter'; import { JobState, JobData, JobOutput, JobStatus, JobInterruptOptions, JobCompletionOptions, ExtensionType } from '../../types/job'; import { HotMeshApps, HotMeshConfig, HotMeshManifest, HotMeshSettings } from '../../types/hotmesh'; import { ProviderClient, ProviderTransaction } from '../../types/provider'; import { JobMessageCallback } from '../../types/quorum'; import { StringAnyType, StringStringType } from '../../types/serializer'; import { GetStatsOptions, IdsResponse, JobStatsInput, StatsResponse } from '../../types/stats'; import { StreamCode, StreamData, StreamDataResponse, StreamStatus } from '../../types/stream'; import { WorkListTaskType } from '../../types/task'; declare class EngineService { namespace: string; apps: HotMeshApps | null; appId: string; guid: string; inited: string; exporter: ExporterService | null; /** @hidden */ search: SearchService | null; /** @hidden */ store: StoreService | null; /** @hidden */ stream: StreamService | null; /** @hidden */ subscribe: SubService | null; /** @hidden */ router: Router | null; /** @hidden */ taskService: TaskService | null; logger: ILogger; cacheMode: CacheMode; untilVersion: string | null; jobCallbacks: Record; /** * @private */ constructor(); /** * @private */ static init(namespace: string, appId: string, guid: string, config: HotMeshConfig, logger: ILogger): Promise; /** * @private */ getSettings(): Promise; /** * @private */ getVID(vid?: AppVID): Promise; /** * @private */ setCacheMode(cacheMode: CacheMode, untilVersion: string): void; /** * @private */ initActivity(topic: string, data?: JobData, context?: JobState): Promise; /** * @private */ getSchema(topic: string): Promise<[activityId: string, schema: ActivityType]>; /** * @private */ isPrivate(topic: string): boolean; /** * @private */ plan(pathOrYAML: string): Promise; /** * @private */ deploy(pathOrYAML: string): Promise; /** * @private */ getStats(topic: string, query: JobStatsInput): Promise; /** * @private */ getIds(topic: string, query: JobStatsInput, queryFacets?: string[]): Promise; /** * @private */ resolveQuery(topic: string, query: JobStatsInput): Promise; /** * @private */ processStreamMessage(streamData: StreamDataResponse): Promise; /** * @private */ execAdjacentParent(context: JobState, jobOutput: JobOutput, emit?: boolean, transaction?: ProviderTransaction): Promise; /** * @private */ hasParentJob(context: JobState, checkSevered?: boolean): boolean; /** * @private */ interrupt(topic: string, jobId: string, options?: JobInterruptOptions): Promise; /** * @private */ scrub(jobId: string): Promise; /** * @private */ runJobCompletionTasks(context: JobState, options?: JobCompletionOptions, transaction?: ProviderTransaction): Promise; /** * @private */ signal(topic: string, data: JobData, status?: StreamStatus, code?: StreamCode, transaction?: ProviderTransaction): Promise; /** * @private */ hookTime(jobId: string, gId: string, topicOrActivity: string, type?: WorkListTaskType): Promise; /** * @private */ signalAll(hookTopic: string, data: JobData, keyResolver: JobStatsInput, queryFacets?: string[]): Promise; /** * @private */ routeToSubscribers(topic: string, message: JobOutput): Promise; /** * @private */ processWebHooks(): Promise; /** * @private */ processTimeHooks(): Promise; /** * @private */ throttle(delayInMillis: number): Promise; /** * Apply a remote duress signal from the quorum. * Delegates to the router's duress manager. * @private */ applyRemoteDuress(throttleMs: number, level: string): void; /** * @private */ pub(topic: string, data: JobData, context?: JobState, extended?: ExtensionType): Promise; /** * @private */ sub(topic: string, callback: JobMessageCallback): Promise; /** * @private */ unsub(topic: string): Promise; /** * @private */ psub(wild: string, callback: JobMessageCallback): Promise; /** * @private */ punsub(wild: string): Promise; /** * @private */ pubsub(topic: string, data: JobData, context?: JobState | null, timeout?: number): Promise; /** * @private */ pubOneTimeSubs(context: JobState, jobOutput: JobOutput, emit?: boolean, transaction?: ProviderTransaction): Promise; /** * @private */ pubPermSubs(context: JobState, jobOutput: JobOutput, emit?: boolean, transaction?: ProviderTransaction): Promise; /** * @private */ getPublishesTopic(context: JobState): Promise; /** * @private */ add(streamData: StreamData | StreamDataResponse): Promise; /** * @private */ registerJobCallback(jobId: string, jobCallback: JobMessageCallback): void; /** * @private */ delistJobCallback(jobId: string): void; /** * @private */ hasOneTimeSubscription(context: JobState): boolean; /** * @private */ export(jobId: string, options?: ExportOptions): Promise; /** * @private */ getRaw(jobId: string): Promise; /** * @private */ getStatus(jobId: string): Promise; /** * @private */ getState(topic: string, jobId: string): Promise; /** * @private */ getQueryState(jobId: string, fields: string[]): Promise; /** * @private * @deprecated */ compress(terms: string[]): Promise; } export { EngineService };