import { HotMesh } from '../hotmesh'; import { EventsConfig } from '../../types/system_events'; import { Connection } from '../../types/durable'; import { EscalationEntry, ClaimEscalationResult, ClaimByMetadataResult, ReleaseEscalationResult, ResolveEscalationResult, CancelEscalationResult, ListEscalationsParams, StatsEscalationsParams, EscalationStats, CreateEscalationParams, UpdateEscalationParams, AppendMilestonesParams, ClaimEscalationParams, ClaimByMetadataParams, ReleaseEscalationParams, ResolveEscalationParams, ResolveByMetadataParams, EscalateToRoleParams, MigrateEscalationParams, ClaimManyParams, EscalateManyToRoleParams, UpdateManyPriorityParams, ResolveManyParams, ResolveAllOrNoneParams, ResolveAllOrNoneResult, PruneEscalationsParams, PruneEscalationsResult, ClaimManyByQueryParams, ResolveBatchItemParams, ResolveBatchItemByMetadataParams, ResolveBatchItemResult, AccumulateItemParams, AccumulateItemByMetadataParams, AccumulateItemResult, RemoveAccumulatedItemParams, RemoveAccumulatedItemByMetadataParams, RemoveAccumulatedItemResult } from '../../types/hmsh_escalations'; /** * Reserved signal-payload key carrying resolution provenance * (`EscalationResolution`) to the waiting workflow. The `$` prefix marks the * control-key namespace — consumer payload fields never collide with it. */ export declare const ESCALATION_RESOLUTION_KEY = "$resolution"; export type GetHotMeshFn = (topic: string | null, namespace?: string) => Promise; export interface EscalationClientConfig { /** Postgres connection options — used when creating a standalone EscalationClient. */ connection?: Connection; /** * Inject a pre-existing `getHotMeshClient` function (e.g. from Durable.Client). * When provided, the client reuses the caller's engine pool — no extra connections. */ getHotMeshClient?: GetHotMeshFn; /** * Optional system-event sink. When set, this client calls `events.publish` * post-commit for every escalation lifecycle transition it performs. * Fire-and-forget; a publish error never fails the committed operation. */ events?: EventsConfig; } /** * Standalone client for the `public.hmsh_escalations` signal-pause surface. * * Requires NO dependency on `services/durable/`. Any HotMesh consumer — AI * agent, YAML DAG worker, REST API — can interact with the escalation queue * directly with just a Postgres connection. * * Signal delivery (for `resolve()` / `resolveByMetadata()`) uses HotMesh's * `engine.signal()` internally. The engine is initialised lazily on first use * and cached for the lifetime of the process. * * @example * ```typescript * import { Escalations } from '@hotmeshio/hotmesh'; * import { Client as Postgres } from 'pg'; * * const client = new Escalations.Client({ * connection: { * class: Postgres, * options: { connectionString: 'postgresql://usr:pwd@localhost:5432/db' }, * }, * }); * * // Claim the next available approval for the 'manager' role * const result = await client.claimByMetadata({ * key: 'orderId', value: 'order-123', * assignee: 'alice@company.com', * roles: ['manager'], * }); * * if (result.ok) { * await client.resolve({ id: result.entry.id, resolverPayload: { approved: true } }); * } * ``` */ export declare class EscalationClientService { private readonly _engine; private readonly _events?; static instances: Map>; constructor(config?: EscalationClientConfig); /** * Fires the configured `events.publish` hook post-commit. * Detached (not awaited); a publish error never fails the caller. * @private */ private _emit; /** * Fires per-row events for bulk operations. Skipped rows are not emitted. * @private */ private _emitMany; private _makeEngineFactory; private _hashConnection; /** * Composes the signal `data` delivered to the waiting workflow. When the * resolving caller supplies `resolvedBy`, resolution provenance (escalation * id + resolver identity) rides under the reserved `$resolution` key — * making the waiter's payload contract deterministic per call site. Without * `resolvedBy` the payload passes through byte-identical, so callers that * exact-match on payloads are unaffected. The enrichment rides the SIGNAL * only — the stored `resolver_payload` column receives the caller's payload * untouched. */ private _signalData; private _deliverEscalationSignal; /** * Builds the wake as a webhook message so the store can commit it * INSIDE the resolve/cancel transaction — the wake becomes durable * with the status change, closing the crash window between commit and * post-commit signal delivery. Mirrors `_deliverEscalationSignal`'s * topic fallback chain; returns null when no hook rule is deployed * for any candidate topic (the caller then keeps post-commit * delivery as the only path, preserving prior behavior). */ private _buildWakeCommand; /** * Returns all escalation rows matching the given filters. Each row includes * a computed `available` field (true = claimable). Supports `sortBy`, * `sortOrder`, `orderBy[]`, and multi-role `roles[]` filter. */ list(params?: ListEscalationsParams): Promise; /** * Returns the count of escalation rows matching the given filters. * Uses the same filter parameters as `list()`. */ count(params?: ListEscalationsParams): Promise; /** * Retention: deletes terminal escalation rows (`resolved`/`cancelled`/`expired`) * whose `updated_at` is older than the horizon (a Postgres interval string, * e.g. `'90 days'`). Terminal rows are inert — every engine state transition * guards on `status = 'pending'` — so pruning them is safe for live waiters, * claims, and signal delivery; this call is the engine-blessed way to age * out the audit backlog. Each call deletes at most `limit` rows (default * 10,000) in one atomic statement; loop until `deleted` is 0 to drain a * large backlog. `list()`/`get()`/`stats()` reads over windows older than * the pruning horizon reflect only the rows retained. */ prune(params: PruneEscalationsParams): Promise; /** Returns a single escalation row by UUID. Returns `null` if not found. */ get(id: string, namespace?: string): Promise; /** Looks up an escalation by `signal_key` — the value passed to `condition()`. */ getBySignalKey(signalKey: string, namespace?: string): Promise; /** * Creates a standalone escalation row with `signal_key = null`. * Useful for external task tracking that doesn't need to resume a workflow. */ create(params: CreateEscalationParams): Promise; /** * Patches an existing escalation row. `metadata` is merged, not replaced. * Signal routing fields can be enriched after creation. */ update(params: UpdateEscalationParams): Promise; /** Appends milestone entries to the escalation's audit trail. */ appendMilestones(params: AppendMilestonesParams): Promise; /** * Atomically claims an escalation by UUID. Implicit model: `status` stays * `'pending'`; claim is expressed via `assigned_to` + `assigned_until`. * Returns `isExtension: true` when the same assignee re-claims a row they already hold. */ claim(params: ClaimEscalationParams): Promise; /** * Atomically claims the highest-priority pending escalation whose `metadata` * contains the given key/value. Optionally merges `metadata` into the claimed row * in the same atomic UPDATE. Returns `isExtension: true` when the same assignee * re-claims a row they already hold (extends the expiry). */ claimByMetadata(params: ClaimByMetadataParams): Promise; /** Releases a claimed escalation, returning it to available status. */ release(params: ReleaseEscalationParams): Promise; /** * Reassigns the escalation to a different role, clearing any current claim * and resetting status to `'pending'`. */ escalateToRole(params: EscalateToRoleParams): Promise; /** * Cancels a pending escalation and delivers a cancellation signal to the * waiting workflow so that `condition()` returns `null`. Terminal rows * return `already-terminal`. The cancellation wake commits inside the * cancel transaction (same durability contract as `resolve()`); when no * hook rule is deployed for any candidate topic, delivery falls back to * a best-effort post-commit publish. */ cancel(id: string, namespace?: string): Promise; /** * Emits local `cancelled` events for a batch of already-cancelled escalation * entries. Called by `WorkflowHandleService.terminate()` after the single * atomic transaction that interrupts the workflow and cancels its escalations * has committed. Fire-and-forget via the configured `events.publish` sink * (e.g. NATS) — instance-local, never broadcast via Postgres LISTEN/NOTIFY. */ emitCancelledBatch(entries: EscalationEntry[]): void; /** * Resolves a pending escalation by UUID. Uses an explicit Postgres transaction * with FOR UPDATE + WHERE guard: only one concurrent caller can commit the * status change; the committed resolved row with its `signal_key` is the * durable proof. Signal delivery is best-effort post-commit — the resolved * row is the recovery record for any missed delivery. Returns the updated * row as `entry` on success. * * Pass `params.metadata` to merge the resolution outcome ("what actually * happened") into the row's GIN-indexed `metadata` in the same atomic UPDATE, * making it `@>`-queryable alongside the creation metadata ("what was * intended"). This is distinct from `resolverPayload`, which is delivered to * the waiting workflow as `condition()`'s return value and is not GIN-indexed. * * Pass `params.assertClaim` to additionally require — inside the same guarded * UPDATE — that no active claim lock stands against that assignee. Blocks * with `claim-expired` (own claim window lapsed) or `claimed-by-other` (a * live window held by someone else); unclaimed rows and durable * pre-assignments (no window) resolve normally. */ resolve(params: ResolveEscalationParams, namespace?: string): Promise; /** * Resolves the highest-priority matching escalation by metadata filter, * then delivers its signal. Same transaction + WHERE guard semantics as `resolve()`. * * The `key`/`value` are the *selector* used to find the row; pass * `params.metadata` to additionally merge a resolution patch into that row's * GIN-indexed `metadata` in the same atomic UPDATE. See {@link resolve}. */ resolveByMetadata(params: ResolveByMetadataParams, namespace?: string): Promise; /** * Pre-serializes the `$resolution` control object merged into a completing * batch item's delivered collection. Rides the signal only — the stored * `resolver_payload` receives the bare collection. */ private _batchResolutionJson; /** * Shared post-statement handling for both batch-item forms: post-commit * wake fallback (the committed `resolver_payload` — the full assembled * collection — is the recovery record) and lifecycle events * (`batch-item` on interim fills, `resolved` on completion). */ private _settleBatchItemResult; /** * Fills ONE declared item of a batch escalation (a wait created with * `condition(signalId, { batch: [...] })` or a standalone `create()` with * `batch`). One atomic statement: the payload lands in * `envelope.batch_items[itemKey]` only while `itemKey` is still pending, * the `batch_pending`/`batch_count` facets recompute, and — on the LAST * item — the row resolves with the assembled collection as its * `resolver_payload` and the waiting workflow's wake commits WITH the * fill (same durability contract as `resolve()`). The caller that lands * the last item learns it from `outcome: 'completed'`. * * Fills are claim-agnostic by default — a batch accumulates contributions * from multiple principals. Pass `assertClaim` to opt into the same * claim-lock assertion as `resolve()`. Pass `metadata` to merge an outcome * patch into the row's GIN-indexed metadata in the same UPDATE (reserved * batch keys cannot be overridden). Pass `resolvedBy` on submissions to * deliver `$resolution` provenance with the completing item's signal. * * A duplicate submission of an already-filled key returns * `outcome: 'duplicate-item'` without touching the row — safe under * webhook retries. */ resolveBatchItem(params: ResolveBatchItemParams, namespace?: string): Promise; /** * Batch-item fill selecting the row by metadata facet — the highest * priority pending row whose `metadata` contains the key/value, mirroring * `resolveByMetadata()`'s selector semantics. See {@link resolveBatchItem} * for the fill contract. */ resolveBatchItemByMetadata(params: ResolveBatchItemByMetadataParams, namespace?: string): Promise; private _previewAccumulateRow; /** * Pre-builds one side's wake with a placeholder payload; the store * rewrites the `{data,data}` slot from the committed collection inside * the add statement, so a completing side's wake commits WITH the add. */ private _accumulateWake; private _assertAccumulatePatch; /** * Shared post-statement handling for both add forms: post-commit wake * fallback per completed side (the committed `resolver_payload` is the * recovery record) and lifecycle events (`accumulated` on an interim * add, `resolved` on completion) for the container and the reciprocal. */ private _settleAccumulateResult; private _accumulateWakes; /** * Adds ONE item to an accumulator escalation (a wait created with * `condition(signalId, { accumulate: {...} })` or a standalone `create()` * with `accumulate`). One atomic statement: the entry lands under * `envelope.accumulate_items[itemKey]` with a database-clock `at`, the * `accumulate_keys` / `accumulate_count` facets recompute, and, when the * add reaches `max` with `resolveAtMax`, the row resolves with the * ordered collection (`{ $accumulated, $trigger: 'count' }`) as its * `resolver_payload` and the waiting workflow's wake commits WITH the * add. The caller that completes the row learns it from * `outcome: 'completed'`. * * Pass `reciprocal` to write a second accumulator row in the same * statement, both or neither: the reciprocal holds the container's id as * its item key, each entry carries the other row's id as `reciprocalId`, * and a member row declared `accumulate: { max: 1 }` completes and wakes * here too. A blocked reciprocal leaves the container untouched and * answers `reciprocal-*`. * * Adds are claim-agnostic by default; pass `assertClaim` for the same * claim-lock assertion as `resolve()` on the container row. Pass * `metadata` to merge an outcome patch into the container's GIN-indexed * metadata in the same UPDATE (reserved accumulate keys are rejected). * A repeated key answers `duplicate-item` without touching the row * unless the accumulator was declared `unique: false`, which replaces * the entry in place. */ accumulateItem(params: AccumulateItemParams, namespace?: string): Promise; /** * Add selecting the container by metadata facet: the highest priority * pending row whose `metadata` contains the key/value, mirroring * `resolveByMetadata()`'s selector. See {@link accumulateItem} for the * add contract. */ accumulateItemByMetadata(params: AccumulateItemByMetadataParams, namespace?: string): Promise; private _settleRemoveResult; /** * Removes ONE held item from a pending accumulator escalation in one * guarded statement: the entry leaves `accumulate_items`, the facets * recompute, the row stays `pending`, and the waiter is never woken. * Pass `reciprocal` to also remove the container's id from that row, * both or neither. A key that is not held answers `item-absent`. */ removeAccumulatedItem(params: RemoveAccumulatedItemParams, namespace?: string): Promise; /** Removal selecting the container by metadata facet. See {@link removeAccumulatedItem}. */ removeAccumulatedItemByMetadata(params: RemoveAccumulatedItemByMetadataParams, namespace?: string): Promise; /** * Full-fidelity migration: inserts an escalation row preserving the original * UUID and lifecycle state. Returns `null` on duplicate (idempotent). */ migrate(params: MigrateEscalationParams, namespace?: string): Promise; /** * No-op in the implicit claim model — availability is computed at query time * from `assigned_until`. Kept for API compatibility. */ releaseExpired(namespace?: string): Promise; /** * Bulk-claims up to `ids.length` pending escalations in one statement. * Returns `{ claimed, skipped }` — skipped rows are either already claimed * by another assignee or non-existent. Implicit-claim semantics apply. */ claimMany(params: ClaimManyParams): Promise<{ claimed: number; skipped: number; }>; /** * Atomic query-form bulk claim: one UPDATE selects and claims every pending, * claimable row matching the selector (role/type/priority equality, metadata * `@>` containment). Prefer this over `list()` + `claimMany({ids})` when the * population is describable by filter — a row that re-parks between a search * and an ids-claim is invisible to the ids form but claimed by this one. * Returns the claimed rows. */ claimManyByQuery(params: ClaimManyByQueryParams): Promise<{ claimed: number; entries: EscalationEntry[]; }>; /** * Bulk-reassigns pending escalations to a new role, clearing any current claim. * Returns the count of rows updated. */ escalateManyToRole(params: EscalateManyToRoleParams): Promise; /** * Bulk-updates priority for pending escalations. Returns the count of rows updated. */ updateManyPriority(params: UpdateManyPriorityParams): Promise; /** * Bulk-resolves pending escalations by id-set. No signal delivery — intended * for redirect-to-triage flows where no workflow is waiting. Returns the * resolved rows. * * Pass `params.metadata` to merge a resolution patch into every winning * (still-pending) row's GIN-indexed `metadata` in the single atomic UPDATE. * See {@link resolve}. */ /** * Bulk-resolves standalone escalations (rows with `signal_key = null`). * Rows backing a live `condition()` waiter are excluded — they stay * pending so a targeted `resolve()`/`cancel()` can deliver their wake; * bulk resolution carries no wake and would strand the workflow. */ resolveMany(params: ResolveManyParams): Promise; /** * All-or-none bulk resolve with per-row payloads. Every item's escalation * must be pending (and, when `assertAssignee` is set, currently assigned to * that assignee) or NOTHING resolves — one atomic SQL statement locks the * rows in deterministic order, applies each row's own `resolverPayload`, * and commits each waiter's wake WITH its resolve (same durability contract * as `resolve()`). Unlike `resolveMany()`, rows backing a live `condition()` * waiter are first-class here: each receives its own payload as * `condition()`'s return value. * * On `ok: false` no row was written and `failed` lists exactly the blocking * rows with reasons (`not-found`, `already-resolved`, `already-cancelled`, * `already-expired`, `assignee-mismatch`) — rows that were themselves * resolvable stay pending and are not listed. * * Pass `params.metadata` to merge one shared outcome patch into every row's * GIN-indexed `metadata` in the same statement. Ids must be unique; an empty * `items` array returns `ok: true` with no entries. */ resolveAllOrNone(params: ResolveAllOrNoneParams, namespace?: string): Promise; /** * Returns dashboard-ready escalation counts. `period` controls the window * used for `created` and `resolved` counts (default `'24h'`). When `roles` * is an empty array, all counts are zero (RBAC guard). */ stats(params?: StatsEscalationsParams): Promise; /** * Returns the sorted list of distinct `type` values in the escalations table. * Useful for populating filter dropdowns. */ listDistinctTypes(namespace?: string): Promise; static shutdown(): Promise; }