{"version":3,"sources":["../../src/errors.ts","../../src/utils/sleep.ts","../../src/storage/postgres.ts"],"names":[],"mappings":";;;AA+DO,IAAM,oBAAA,GAAN,cAAmC,KAAA,CAAM;AAAA,EAC9C,WAAA,CACkB,QAAA,EACA,MAAA,EAChB,MAAA,EACA;AACA,IAAA,KAAA;AAAA,MACE,MAAA,GACI,CAAA,2BAAA,EAA8B,QAAQ,CAAA,CAAA,EAAI,MAAM,CAAA,EAAA,EAAK,MAAM,CAAA,CAAA,GAC3D,CAAA,2BAAA,EAA8B,QAAQ,CAAA,CAAA,EAAI,MAAM,CAAA;AAAA,KACtD;AARgB,IAAA,IAAA,CAAA,QAAA,GAAA,QAAA;AACA,IAAA,IAAA,CAAA,MAAA,GAAA,MAAA;AAQhB,IAAA,IAAA,CAAK,IAAA,GAAO,sBAAA;AAAA,EACd;AACF,CAAA;;;ACtEO,SAAS,KAAA,CAAM,IAAY,MAAA,EAAqC;AACrE,EAAA,IAAI,MAAM,CAAA,EAAG;AAEX,IAAA,OAAO,QAAQ,OAAA,EAAQ;AAAA,EACzB;AAEA,EAAA,OAAO,IAAI,OAAA,CAAQ,CAAC,OAAA,EAAS,MAAA,KAAW;AAMtC,IAAc,WAAW,MAAM;AAE7B,MAAA,OAAA,EAAQ;AAAA,IACV,GAAG,EAAE;AAQoD,EAC3D,CAAC,CAAA;AACH;;;ACmBO,IAAM,kBAAN,MAAgD;AAAA,EAMrD,YAAY,IAAA,EAA8B;AACxC,IAAA,IAAI,CAAC,KAAK,IAAA,EAAM;AACd,MAAA,MAAM,IAAI,MAAM,oDAAoD,CAAA;AAAA,IACtE;AACA,IAAA,IAAA,CAAK,OAAO,IAAA,CAAK,IAAA;AACjB,IAAA,IAAA,CAAK,SAAA,GAAY,KAAK,SAAA,IAAa,iBAAA;AACnC,IAAA,IAAA,CAAK,aAAA,GAAgB,KAAK,aAAA,IAAiB,UAAA;AAC3C,IAAA,IAAA,CAAK,UAAA,GAAa,KAAK,UAAA,IAAc,EAAA;AAErC,IAAA,IAAI,CAAC,0BAAA,CAA2B,IAAA,CAAK,IAAA,CAAK,SAAS,CAAA,EAAG;AACpD,MAAA,MAAM,IAAI,KAAA,CAAM,CAAA,6BAAA,EAAgC,IAAA,CAAK,SAAS,CAAA,CAAA,CAAG,CAAA;AAAA,IACnE;AAAA,EACF;AAAA,EAEA,MAAM,IAAA,CAAK,QAAA,EAAkB,MAAA,EAA2C;AACtE,IAAA,MAAM,EAAE,IAAA,EAAK,GAAI,MAAM,KAAK,IAAA,CAAK,KAAA;AAAA,MAC/B,CAAA,kBAAA,EAAqB,KAAK,SAAS,CAAA,sCAAA,CAAA;AAAA,MACnC,CAAC,UAAU,MAAM;AAAA,KACnB;AACA,IAAA,OAAO,IAAA,CAAK,CAAC,CAAA,EAAG,KAAA,IAAS,IAAA;AAAA,EAC3B;AAAA,EAEA,MAAM,KAAK,KAAA,EAAiC;AAC1C,IAAA,MAAM,KAAK,IAAA,CAAK,KAAA;AAAA,MACd,CAAA,YAAA,EAAe,KAAK,SAAS,CAAA;AAAA;AAAA;AAAA,6EAAA,CAAA;AAAA,MAI7B,CAAC,MAAM,QAAA,EAAU,KAAA,CAAM,QAAQ,IAAA,CAAK,SAAA,CAAU,KAAK,CAAC;AAAA,KACtD;AAAA,EACF;AAAA,EAEA,MAAM,MAAA,CAAO,QAAA,EAAkB,MAAA,EAA+B;AAC5D,IAAA,MAAM,KAAK,IAAA,CAAK,KAAA;AAAA,MACd,CAAA,YAAA,EAAe,KAAK,SAAS,CAAA,sCAAA,CAAA;AAAA,MAC7B,CAAC,UAAU,MAAM;AAAA,KACnB;AAAA,EACF;AAAA,EAEA,MAAM,WAAA,CACJ,QAAA,EACA,MAAA,EACA,OAAA,EACe;AACf,IAAA,MAAM,EAAE,WAAU,GAAI,OAAA;AACtB,IAAA,MAAM,MAAA,GAAS,MAAM,IAAA,CAAK,IAAA,CAAK,OAAA,EAAQ;AAEvC,IAAA,IAAI;AAEF,MAAA,IAAI,MAAM,IAAA,CAAK,OAAA,CAAQ,MAAA,EAAQ,QAAA,EAAU,MAAM,CAAA,EAAG;AAChD,QAAA,OAAO,IAAA,CAAK,QAAA,CAAS,MAAA,EAAQ,QAAA,EAAU,MAAM,CAAA;AAAA,MAC/C;AAEA,MAAA,IAAI,aAAa,CAAA,EAAG;AAClB,QAAA,MAAM,IAAI,oBAAA,CAAqB,QAAA,EAAU,MAAA,EAAQ,cAAc,CAAA;AAAA,MACjE;AAGA,MAAA,MAAM,KAAA,GAAQ,KAAK,GAAA,EAAI;AACvB,MAAA,OAAO,IAAA,CAAK,GAAA,EAAI,GAAI,KAAA,GAAQ,SAAA,EAAW;AACrC,QAAA,MAAM,SAAA,GAAY,SAAA,IAAa,IAAA,CAAK,GAAA,EAAI,GAAI,KAAA,CAAA;AAC5C,QAAA,MAAM,MAAM,IAAA,CAAK,GAAA,CAAI,IAAA,CAAK,UAAA,EAAY,SAAS,CAAC,CAAA;AAChD,QAAA,IAAI,MAAM,IAAA,CAAK,OAAA,CAAQ,MAAA,EAAQ,QAAA,EAAU,MAAM,CAAA,EAAG;AAChD,UAAA,OAAO,IAAA,CAAK,QAAA,CAAS,MAAA,EAAQ,QAAA,EAAU,MAAM,CAAA;AAAA,QAC/C;AAAA,MACF;AAEA,MAAA,MAAM,IAAI,oBAAA;AAAA,QACR,QAAA;AAAA,QACA,MAAA;AAAA,QACA,sBAAsB,SAAS,CAAA,EAAA;AAAA,OACjC;AAAA,IACF,SAAS,GAAA,EAAK;AACZ,MAAA,MAAA,CAAO,OAAA,EAAQ;AACf,MAAA,MAAM,GAAA;AAAA,IACR;AAAA,EACF;AAAA,EAEA,MAAc,OAAA,CACZ,MAAA,EACA,QAAA,EACA,MAAA,EACkB;AAClB,IAAA,MAAM,EAAE,IAAA,EAAK,GAAI,MAAM,MAAA,CAAO,KAAA;AAAA,MAC5B,CAAA,6DAAA,CAAA;AAAA,MACA,CAAC,IAAA,CAAK,aAAA,EAAe,GAAG,QAAQ,CAAA,CAAA,EAAI,MAAM,CAAA,CAAE;AAAA,KAC9C;AACA,IAAA,OAAO,IAAA,CAAK,CAAC,CAAA,EAAG,EAAA,KAAO,IAAA;AAAA,EACzB;AAAA,EAEQ,QAAA,CAAS,MAAA,EAAsB,QAAA,EAAkB,MAAA,EAAsB;AAC7E,IAAA,MAAM,YAAY,IAAA,CAAK,aAAA;AACvB,IAAA,MAAM,GAAA,GAAM,CAAA,EAAG,QAAQ,CAAA,CAAA,EAAI,MAAM,CAAA,CAAA;AACjC,IAAA,IAAI,QAAA,GAAW,KAAA;AACf,IAAA,OAAO;AAAA,MACL,MAAM,OAAA,GAAU;AACd,QAAA,IAAI,QAAA,EAAU;AACd,QAAA,QAAA,GAAW,IAAA;AACX,QAAA,IAAI;AACF,UAAA,MAAM,MAAA,CAAO,KAAA;AAAA,YACX,CAAA,qDAAA,CAAA;AAAA,YACA,CAAC,WAAW,GAAG;AAAA,WACjB;AAAA,QACF,CAAA,SAAE;AACA,UAAA,MAAA,CAAO,OAAA,EAAQ;AAAA,QACjB;AAAA,MACF,CAAA;AAAA;AAAA;AAAA,MAGA,MAAM,OAAA,GAAU;AAAA,MAEhB;AAAA,KACF;AAAA,EACF;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA,EAOA,MAAM,YAAA,GAA8B;AAClC,IAAA,MAAM,KAAK,IAAA,CAAK,KAAA;AAAA,MACd,CAAA,2BAAA,EAA8B,KAAK,SAAS,CAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA,QAAA;AAAA,KAO9C;AACA,IAAA,MAAM,KAAK,IAAA,CAAK,KAAA;AAAA,MACd,CAAA,2BAAA,EAA8B,KAAK,SAAS,CAAA;AAAA,YAAA,EACpC,KAAK,SAAS,CAAA,aAAA;AAAA,KACxB;AAAA,EACF;AACF;AAEO,SAAS,sBAAsB,IAAA,EAA+C;AACnF,EAAA,OAAO,IAAI,gBAAgB,IAAI,CAAA;AACjC","file":"postgres.cjs","sourcesContent":["import type { SerializedError } from './types.js';\n\n/** Mark an error as permanently fatal — the executor must not retry it. */\nexport class PermanentError extends Error {\n  readonly permanent = true as const;\n  readonly code?: string;\n\n  constructor(message: string, options?: { cause?: unknown; code?: string }) {\n    super(message);\n    this.name = 'PermanentError';\n    if (options?.cause !== undefined) (this as { cause?: unknown }).cause = options.cause;\n    if (options?.code) this.code = options.code;\n  }\n}\n\n/** Mark an error as transient — explicitly eligible for retry. */\nexport class TransientError extends Error {\n  readonly transient = true as const;\n  readonly code?: string;\n\n  constructor(message: string, options?: { cause?: unknown; code?: string }) {\n    super(message);\n    this.name = 'TransientError';\n    if (options?.cause !== undefined) (this as { cause?: unknown }).cause = options.cause;\n    if (options?.code) this.code = options.code;\n  }\n}\n\n/** Thrown when a step exceeds its configured timeout. Retryable by default. */\nexport class StepTimeoutError extends Error {\n  readonly timeout = true as const;\n\n  constructor(\n    public readonly stepName: string,\n    public readonly timeoutMs: number,\n  ) {\n    super(`step \"${stepName}\" timed out after ${timeoutMs}ms`);\n    this.name = 'StepTimeoutError';\n  }\n}\n\n/** Thrown when execution is cancelled via AbortSignal. */\nexport class FlowAbortedError extends Error {\n  readonly aborted = true as const;\n\n  constructor(reason?: unknown) {\n    const msg =\n      reason instanceof Error\n        ? reason.message\n        : typeof reason === 'string'\n          ? reason\n          : 'flow aborted';\n    super(msg);\n    this.name = 'FlowAbortedError';\n    if (reason !== undefined) (this as { cause?: unknown }).cause = reason;\n  }\n}\n\n/**\n * Thrown when the storage adapter cannot acquire a lock within the configured\n * wait timeout — typically because another worker is currently executing the\n * same idempotency key.\n */\nexport class LockAcquisitionError extends Error {\n  constructor(\n    public readonly flowName: string,\n    public readonly flowId: string,\n    reason?: string,\n  ) {\n    super(\n      reason\n        ? `failed to acquire lock for ${flowName}/${flowId}: ${reason}`\n        : `failed to acquire lock for ${flowName}/${flowId}`,\n    );\n    this.name = 'LockAcquisitionError';\n  }\n}\n\n/** Thrown by the executor when a step fails and compensation runs. */\nexport class FlowError extends Error {\n  constructor(\n    message: string,\n    public readonly flowId: string,\n    public readonly flowName: string,\n    public readonly failedStep: string,\n    public readonly originalError: unknown,\n    public readonly compensationErrors: Array<{ step: string; error: unknown }> = [],\n  ) {\n    super(message);\n    this.name = 'FlowError';\n    (this as { cause?: unknown }).cause = originalError;\n  }\n}\n\n/** Heuristic: errors are retryable unless explicitly marked permanent. */\nexport function isPermanent(err: unknown): boolean {\n  return (\n    err instanceof PermanentError ||\n    (typeof err === 'object' && err !== null && (err as { permanent?: boolean }).permanent === true)\n  );\n}\n\nexport function isTransient(err: unknown): boolean {\n  return (\n    err instanceof TransientError ||\n    err instanceof StepTimeoutError ||\n    (typeof err === 'object' && err !== null && (err as { transient?: boolean }).transient === true)\n  );\n}\n\nexport function serializeError(err: unknown): SerializedError {\n  if (err instanceof Error) {\n    const out: SerializedError = { name: err.name, message: err.message };\n    if (err.stack) out.stack = err.stack;\n    const withCode = err as { code?: unknown };\n    if (typeof withCode.code === 'string') out.code = withCode.code;\n    if (isTransient(err)) out.transient = true;\n    return out;\n  }\n  return { name: 'UnknownError', message: String(err) };\n}\n","import { FlowAbortedError } from '../errors.js';\n\n/**\n * Promise-based sleep that respects an AbortSignal. Rejects with\n * FlowAbortedError if the signal fires mid-wait.\n */\nexport function sleep(ms: number, signal?: AbortSignal): Promise<void> {\n  if (ms <= 0) {\n    if (signal?.aborted) return Promise.reject(new FlowAbortedError(signal.reason));\n    return Promise.resolve();\n  }\n\n  return new Promise((resolve, reject) => {\n    if (signal?.aborted) {\n      reject(new FlowAbortedError(signal.reason));\n      return;\n    }\n\n    const timer = setTimeout(() => {\n      signal?.removeEventListener('abort', onAbort);\n      resolve();\n    }, ms);\n\n    const onAbort = () => {\n      clearTimeout(timer);\n      signal?.removeEventListener('abort', onAbort);\n      reject(new FlowAbortedError(signal?.reason));\n    };\n\n    signal?.addEventListener('abort', onAbort, { once: true });\n  });\n}\n","import type { AcquireLockOptions, FlowState, Lock, StorageAdapter } from '../types.js';\nimport { LockAcquisitionError } from '../errors.js';\nimport { sleep } from '../utils/sleep.js';\n\n/**\n * Structural subset of `pg.Pool` — we only need `query` and `connect`.\n * Import `Pool` from `pg` and pass an instance; avoids a hard dep on `pg`.\n */\nexport interface PgPoolLike {\n  query<R = unknown>(text: string, values?: readonly unknown[]): Promise<{ rows: R[] }>;\n  connect(): Promise<PgClientLike>;\n}\n\nexport interface PgClientLike {\n  query<R = unknown>(text: string, values?: readonly unknown[]): Promise<{ rows: R[] }>;\n  release(err?: unknown): void;\n}\n\nexport interface PostgresStorageOptions {\n  pool: PgPoolLike;\n  /** Table name for flow state. Default: `kompensa_states`. */\n  tableName?: string;\n  /**\n   * Identifier used to namespace Postgres advisory locks, avoiding collision\n   * with locks used by other subsystems. Default: `kompensa`.\n   */\n  lockNamespace?: string;\n  /** Polling interval while waiting for a contested lock, in ms. Default: 50. */\n  lockPollMs?: number;\n}\n\n/**\n * Durable Postgres-backed storage. Uses JSONB for state and session-level\n * advisory locks for multi-worker safety.\n *\n * Locks automatically release when the holding connection closes — so worker\n * crashes don't permanently wedge an idempotency key. Advisory locks do not\n * have a server-side TTL; the `ttlMs` option is advisory only and ignored.\n *\n * @example\n * import { Pool } from 'pg';\n * import { PostgresStorage } from 'kompensa/storage/postgres';\n *\n * const storage = new PostgresStorage({\n *   pool: new Pool({ connectionString: process.env.DATABASE_URL }),\n * });\n * await storage.ensureSchema();\n *\n * const flow = createFlow('checkout', { storage }).step(...)\n */\nexport class PostgresStorage implements StorageAdapter {\n  private readonly pool: PgPoolLike;\n  private readonly tableName: string;\n  private readonly lockNamespace: string;\n  private readonly lockPollMs: number;\n\n  constructor(opts: PostgresStorageOptions) {\n    if (!opts.pool) {\n      throw new Error('kompensa: PostgresStorage requires a `pool` option');\n    }\n    this.pool = opts.pool;\n    this.tableName = opts.tableName ?? 'kompensa_states';\n    this.lockNamespace = opts.lockNamespace ?? 'kompensa';\n    this.lockPollMs = opts.lockPollMs ?? 50;\n\n    if (!/^[A-Za-z_][A-Za-z0-9_]*$/.test(this.tableName)) {\n      throw new Error(`kompensa: invalid tableName \"${this.tableName}\"`);\n    }\n  }\n\n  async load(flowName: string, flowId: string): Promise<FlowState | null> {\n    const { rows } = await this.pool.query<{ state: FlowState }>(\n      `SELECT state FROM ${this.tableName} WHERE flow_name = $1 AND flow_id = $2`,\n      [flowName, flowId],\n    );\n    return rows[0]?.state ?? null;\n  }\n\n  async save(state: FlowState): Promise<void> {\n    await this.pool.query(\n      `INSERT INTO ${this.tableName} (flow_name, flow_id, state, updated_at)\n       VALUES ($1, $2, $3::jsonb, NOW())\n       ON CONFLICT (flow_name, flow_id)\n       DO UPDATE SET state = EXCLUDED.state, updated_at = EXCLUDED.updated_at`,\n      [state.flowName, state.flowId, JSON.stringify(state)],\n    );\n  }\n\n  async delete(flowName: string, flowId: string): Promise<void> {\n    await this.pool.query(\n      `DELETE FROM ${this.tableName} WHERE flow_name = $1 AND flow_id = $2`,\n      [flowName, flowId],\n    );\n  }\n\n  async acquireLock(\n    flowName: string,\n    flowId: string,\n    options: AcquireLockOptions,\n  ): Promise<Lock> {\n    const { timeoutMs } = options;\n    const client = await this.pool.connect();\n\n    try {\n      // Fast-path: try once without waiting.\n      if (await this.tryLock(client, flowName, flowId)) {\n        return this.makeLock(client, flowName, flowId);\n      }\n\n      if (timeoutMs <= 0) {\n        throw new LockAcquisitionError(flowName, flowId, 'lock is held');\n      }\n\n      // Poll until acquired or timeout elapses.\n      const start = Date.now();\n      while (Date.now() - start < timeoutMs) {\n        const remaining = timeoutMs - (Date.now() - start);\n        await sleep(Math.min(this.lockPollMs, remaining));\n        if (await this.tryLock(client, flowName, flowId)) {\n          return this.makeLock(client, flowName, flowId);\n        }\n      }\n\n      throw new LockAcquisitionError(\n        flowName,\n        flowId,\n        `wait timeout after ${timeoutMs}ms`,\n      );\n    } catch (err) {\n      client.release();\n      throw err;\n    }\n  }\n\n  private async tryLock(\n    client: PgClientLike,\n    flowName: string,\n    flowId: string,\n  ): Promise<boolean> {\n    const { rows } = await client.query<{ ok: boolean }>(\n      `SELECT pg_try_advisory_lock(hashtext($1), hashtext($2)) AS ok`,\n      [this.lockNamespace, `${flowName}:${flowId}`],\n    );\n    return rows[0]?.ok === true;\n  }\n\n  private makeLock(client: PgClientLike, flowName: string, flowId: string): Lock {\n    const namespace = this.lockNamespace;\n    const key = `${flowName}:${flowId}`;\n    let released = false;\n    return {\n      async release() {\n        if (released) return;\n        released = true;\n        try {\n          await client.query(\n            `SELECT pg_advisory_unlock(hashtext($1), hashtext($2))`,\n            [namespace, key],\n          );\n        } finally {\n          client.release();\n        }\n      },\n      // Advisory locks have no server-side TTL — refresh is a no-op. If the\n      // holder crashes, the connection dies and the lock auto-releases.\n      async refresh() {\n        /* no-op */\n      },\n    };\n  }\n\n  /**\n   * Create the state table if it doesn't exist. Safe to call repeatedly.\n   * Use this once at startup or migrate via your preferred migration tool\n   * using the SQL in `kompensa/storage/postgres/schema.sql`.\n   */\n  async ensureSchema(): Promise<void> {\n    await this.pool.query(\n      `CREATE TABLE IF NOT EXISTS ${this.tableName} (\n         flow_name  TEXT NOT NULL,\n         flow_id    TEXT NOT NULL,\n         state      JSONB NOT NULL,\n         updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),\n         PRIMARY KEY (flow_name, flow_id)\n       )`,\n    );\n    await this.pool.query(\n      `CREATE INDEX IF NOT EXISTS ${this.tableName}_updated_idx\n         ON ${this.tableName} (updated_at)`,\n    );\n  }\n}\n\nexport function createPostgresStorage(opts: PostgresStorageOptions): PostgresStorage {\n  return new PostgresStorage(opts);\n}\n"]}