import { KeyStoreParams, KeyType } from '../../modules/key'; import { ScoutType } from '../../types/hotmesh'; import { ILogger } from '../logger'; import { SerializerService as Serializer } from '../serializer'; import { Consumes } from '../../types/activity'; import { AppVID } from '../../types/app'; import { HookRule, HookSignal } from '../../types/hook'; import { HotMeshApps, HotMeshSettings } from '../../types/hotmesh'; import { ProviderClient, ProviderTransaction } from '../../types/provider'; import { ThrottleOptions } from '../../types/quorum'; import { StringAnyType, Symbols, StringStringType, SymbolSets } from '../../types/serializer'; import { IdsData, JobStatsRange, StatsType } from '../../types/stats'; import { WorkListTaskType } from '../../types/task'; import { Cache } from './cache'; declare abstract class StoreService { storeClient: Provider; namespace: string; appId: string; logger: ILogger; cache: Cache; serializer: Serializer; constructor(client: Provider); abstract transact(): TransactionProvider; abstract init(namespace: string, appId: string, logger: ILogger, guid?: string, role?: string): Promise; abstract mintKey(type: KeyType, params: KeyStoreParams): string; abstract getSettings(bCreate?: boolean): Promise; abstract setSettings(manifest: HotMeshSettings): Promise; abstract getApp(id: string, refresh?: boolean): Promise; abstract setApp(id: string, version: string): Promise; abstract activateAppVersion(id: string, version: string): Promise; abstract reserveScoutRole(scoutType: ScoutType, delay?: number): Promise; abstract releaseScoutRole(scoutType: ScoutType): Promise; abstract reserveSymbolRange(target: string, size: number, type: 'JOB' | 'ACTIVITY', tryCount?: number): Promise<[number, number, Symbols]>; abstract getSymbols(activityId: string): Promise; abstract addSymbols(activityId: string, symbols: Symbols): Promise; abstract getSymbolValues(): Promise; abstract addSymbolValues(symvals: Symbols): Promise; abstract getSymbolKeys(symbolNames: string[]): Promise; abstract setStats(jobKey: string, jobId: string, dateTime: string, stats: StatsType, appVersion: AppVID, transaction?: TransactionProvider): Promise; abstract getJobStats(jobKeys: string[]): Promise; abstract getJobIds(indexKeys: string[], idRange: [number, number]): Promise; abstract setStatus(collationKeyStatus: number, jobId: string, appId: string, transaction?: TransactionProvider): Promise; abstract setStatusAndCollateGuid(statusDelta: number, // typically (N - 1) threshold: number, // typically 0 (but supports 0,1,12,...) jobId: string, appId: string, guidField: string, // guidWeight: number, transaction?: ProviderTransaction): Promise; abstract getStatus(jobId: string, appId: string): Promise; abstract setStateNX(jobId: string, appId: string, status?: number, entity?: string, transaction?: ProviderTransaction, originId?: string, parentId?: string): Promise; abstract setState(state: StringAnyType, status: number | null, jobId: string, symbolNames: string[], dIds: StringStringType, transaction?: TransactionProvider): Promise; abstract getQueryState(jobId: string, fields: string[]): Promise; abstract getState(jobId: string, consumes: Consumes, dIds: StringStringType): Promise<[StringAnyType, number] | undefined>; abstract getRaw(jobId: string): Promise; abstract collate(jobId: string, activityId: string, amount: number, dIds: StringStringType, transaction?: TransactionProvider): Promise; abstract collateLeg2Entry(jobId: string, activityId: string, guid: string, dIds: StringStringType, transaction?: TransactionProvider): Promise<[number, number]>; abstract collateSynthetic(jobId: string, guid: string, amount: number, transaction?: TransactionProvider): Promise; abstract getSchema(activityId: string, appVersion: AppVID): Promise; abstract getSchemas(appVersion: AppVID): Promise>; abstract setSchemas(schemas: Record, appVersion: AppVID): Promise; abstract setSubscriptions(subscriptions: Record, appVersion: AppVID): Promise; abstract getSubscriptions(appVersion: AppVID): Promise>; abstract getSubscription(topic: string, appVersion: AppVID): Promise; abstract setTransitions(transitions: Record, appVersion: AppVID): Promise; abstract getTransitions(appVersion: AppVID): Promise; abstract setHookRules(hookRules: Record): Promise; abstract getHookRules(): Promise>; abstract getAllSymbols(): Promise; /** * Leg1: Attempts to set the hook signal. If a pending signal occupies * the key (race condition), overwrites it and returns the pending data. * When called with a transaction, queues the setnxex (no pending detection). * * When `redelivery` is provided (the webhook routing for this hook), * a consumed pending signal is republished as an engine stream message * in the SAME transaction that overwrites the marker — the wake * survives a crash at any instant. `pendingData` is then not returned, * since the store already owns the redelivery. */ abstract setHookSignal(hook: HookSignal, transaction?: TransactionProvider, redelivery?: { aid: string; topic: string; }): Promise<{ success: boolean; pendingData?: string; }>; /** * Leg2: Atomically gets the hook signal OR inserts a pending signal * if no hook is registered yet (early signal). Returns the hook * signal value, or undefined if we stored a pending signal instead. */ abstract getHookSignal(topic: string, resolved: string, pendingData?: string, pendingExpire?: number): Promise; abstract deleteHookSignal(topic: string, resolved: string): Promise; abstract addTaskQueues(keys: string[]): Promise; abstract getActiveTaskQueue(): Promise; abstract deleteProcessedTaskQueue(workItemKey: string, key: string, processedKey: string, scrub?: boolean): Promise; abstract processTaskQueue(sourceKey: string, destinationKey: string): Promise; abstract expireJob(jobId: string, inSeconds: number, txProvider?: TransactionProvider): Promise; abstract getDependencies(jobId: string): Promise; abstract delistSignalKey(key: string, target: string): Promise; abstract registerTimeHook(jobId: string, gId: string, activityId: string, type: WorkListTaskType, deletionTime: number, dad: string, transaction?: TransactionProvider): Promise; abstract getNextTask(listKey?: string): Promise<[ listKey: string, jobId: string, gId: string, activityId: string, type: WorkListTaskType ] | boolean>; abstract interrupt(topic: string, jobId: string, options: { [key: string]: any; }): Promise; abstract scrub(jobId: string): Promise; abstract findJobs(queryString?: string, limit?: number, batchSize?: number, cursor?: string): Promise<[string, string[]]>; abstract findJobFields(jobId: string, fieldMatchPattern?: string, limit?: number, batchSize?: number, cursor?: string): Promise<[string, StringStringType]>; abstract setCancel(jobId: string, appId: string): Promise; abstract setThrottleRate(options: ThrottleOptions): Promise; abstract getThrottleRates(): Promise; abstract getThrottleRate(topic: string): Promise; /** * Fetch activity inputs for a workflow. Used by the exporter to enrich * timeline events with activity arguments. * * @param workflowId - The workflow ID * @param symbolField - The compressed symbol field for activity arguments * @returns Map of job_id -> parsed input arguments and activityName:index -> parsed inputs */ getActivityInputs?(workflowId: string, symbolField: string): Promise<{ byJobId: Map; byNameIndex: Map; }>; /** * Fetch child workflow inputs in batch. Used by the exporter to enrich * child workflow events with their arguments. * * @param childJobKeys - Array of child job keys to fetch * @param symbolField - The compressed symbol field for workflow arguments * @returns Map of child_workflow_id -> parsed input arguments */ getChildWorkflowInputs?(childJobKeys: string[], symbolField: string): Promise>; /** * Fetch and decode a single compressed-symbol field from a job's HASH. * Reads one field with `hmget` rather than the whole hash, so it is the * cheap path for narrow getters like the workflow input arguments. * * @param jobId - The job ID * @param symbolField - The compressed symbol field (symbol + dimension) * @returns The decoded value, or undefined if the field is absent */ getJobArguments?(jobId: string, symbolField: string): Promise; /** * Fetch stream message history for a job from worker_streams. * Returns raw activity input/output data from soft-deleted messages. * * @param jobId - The job ID (metadata.jid in stream messages) * @param options - Optional filters for activity or message types * @returns Array of stream history entries ordered by creation time */ getStreamHistory?(jobId: string, options?: { activity?: string; types?: string[]; }): Promise; /** * Fetch job record and attributes by key. Used by the exporter to * reconstruct execution history for expired jobs. * * @param jobKey - The job key (e.g., "hmsh:durable:j:workflowId") * @returns Job row and all attributes */ getJobByKeyDirect?(jobKey: string): Promise<{ job: { id: string; key: string; status: number; created_at: Date; updated_at: Date; expired_at?: Date; is_live: boolean; }; attributes: Record; }>; /** * Fetch just the lineage columns (`parent_id`, `origin_id`) for a job via a * single indexed lookup on `key`. Powers the exporter's opt-in * `include_lineage` pointer — `parent_id` is the real spawning workflow (never * the synthetic collator `$C` job), `origin_id` is the root ancestor. * * @param jobKey - The job key (e.g., "hmsh:durable:j:workflowId") * @returns parent/origin ids (null when absent — a null `parent_id` marks a root) */ getJobLineage?(jobKey: string): Promise<{ parent_id: string | null; origin_id: string | null; } | null>; } export { StoreService };