/** * Pub/sub messaging, subscriptions, and one-time callbacks. * * This module handles the full publish/subscribe lifecycle: * - `pub()` → fire-and-forget job trigger * - `pubsub()` → trigger + await result (with timeout) * - `sub()/unsub()` → permanent topic subscriptions * - `psub()/punsub()` → pattern-based subscriptions * - `add()` → raw stream message publish */ import { Router } from '../router'; import { StoreService } from '../store'; import { StreamService } from '../stream'; import { SubService } from '../sub'; import { ILogger } from '../logger'; import { AppVID } from '../../types/app'; import { ExtensionType, JobData, JobOutput, JobState } from '../../types/job'; import { ProviderClient, ProviderTransaction } from '../../types/provider'; import { JobMessageCallback } from '../../types/quorum'; import { StreamData, StreamDataResponse } from '../../types/stream'; interface PubSubContext { guid: string; appId: string; store: StoreService; stream: StreamService; subscribe: SubService; router: Router | null; logger: ILogger; jobCallbacks: Record; getVID(vid?: AppVID): Promise; initActivity(topic: string, data?: JobData, context?: JobState): Promise; getState(topic: string, jobId: string): Promise; getPublishesTopic(context: JobState): Promise; } export declare function pub(instance: PubSubContext, topic: string, data: JobData, context?: JobState, extended?: ExtensionType): Promise; export declare function sub(instance: PubSubContext, topic: string, callback: JobMessageCallback): Promise; export declare function unsub(instance: PubSubContext, topic: string): Promise; export declare function psub(instance: PubSubContext, wild: string, callback: JobMessageCallback): Promise; export declare function punsub(instance: PubSubContext, wild: string): Promise; export declare function pubsub(instance: PubSubContext, topic: string, data: JobData, context?: JobState | null, timeout?: number): Promise; export declare function add(instance: PubSubContext, streamData: StreamData | StreamDataResponse): Promise; export declare function registerJobCallback(instance: PubSubContext, jobId: string, jobCallback: JobMessageCallback): void; export declare function removeJobCallback(instance: PubSubContext, jobId: string): void; export declare function hasOneTimeSubscription(context: JobState): boolean; /** * Resolves the `publishes` topic for the activity that produced * this job's output — used to notify permanent subscribers. */ export declare function getPublishesTopic(instance: PubSubContext, context: JobState): Promise; export {};