{"version":3,"sources":["../src/index.ts","../src/detect.ts","../src/ws.ts","../src/transport.ts","../src/embedded.ts","../src/nanoWorker.ts","../src/restPollDefault.ts"],"sourcesContent":["// @nanobpm/nano-sdk — a drop-in replacement for @camunda8/orchestration-cluster-api.\n//\n// It re-exports the entire upstream SDK surface, then overrides\n// `createCamundaClient` (and the default export) so that, when connected to a\n// Nano server, process-instance creation and job workers transparently upgrade\n// to Nano's Falcon protocol. Against stock Camunda 8 it is a no-op: the\n// upstream REST behaviour is used unchanged.\nimport {\n  createCamundaClient as createCamundaClientBase,\n} from \"@camunda8/orchestration-cluster-api\";\nimport { detectNano, normalizeBase, type NanoInfo } from \"./detect.js\";\nimport { FalconTransport } from \"./transport.js\";\nimport { EmbeddedTransport, type EmbeddedHost } from \"./embedded.js\";\nimport { NanoJobWorker } from \"./nanoWorker.js\";\nimport { withRestPollDefault, wrapRestPollDefault } from \"./restPollDefault.js\";\n\n// Re-export everything else so consumers can swap the import path and nothing\n// breaks.\nexport * from \"@camunda8/orchestration-cluster-api\";\nexport { detectNano, type NanoInfo } from \"./detect.js\";\nexport { FalconTransport, MalformedFrameError, SubmissionTimeoutError, ConnectTimeoutError } from \"./transport.js\";\nexport { EmbeddedTransport, type EmbeddedHost, type EmbeddedJob } from \"./embedded.js\";\nexport { NanoJobWorker } from \"./nanoWorker.js\";\n\n/** auto: upgrade only on Nano. falcon: force. rest: never upgrade. embedded: in-process μ-nano. */\nexport type NanoTransport = \"auto\" | \"falcon\" | \"rest\" | \"embedded\";\n\ntype AnyOpts = Parameters<typeof createCamundaClientBase>[0] & {\n  config?: Record<string, unknown> & {\n    CAMUNDA_TRANSPORT?: NanoTransport;\n    CAMUNDA_REST_ADDRESS?: string;\n    CAMUNDA_NANO_SUBMIT_TIMEOUT_MS?: string | number;\n    /**\n     * Client-side bound (ms) on the Falcon WebSocket handshake. If the gateway\n     * advertises Falcon but the socket opens without ever completing the\n     * handshake (a proxy that blackholes the upgrade), the SDK gives up after\n     * this deadline and falls back to REST instead of hanging. Defaults to\n     * {@link DEFAULT_CONNECT_TIMEOUT_MS}; set `0` to wait indefinitely.\n     */\n    CAMUNDA_NANO_CONNECT_TIMEOUT_MS?: string | number;\n    /**\n     * Force plain REST even when the gateway advertises Falcon. Useful for\n     * environments where WebSockets are blocked (corporate proxies etc.).\n     * Accepts any truthy string (`1`, `true`, `yes`, `on`).\n     */\n    CAMUNDA_FORCE_REST?: string | boolean | number;\n  };\n  /** Embedded (ADR 0005) in-process engine host; required when transport is \"embedded\". */\n  embeddedHost?: EmbeddedHost;\n};\n\nfunction isTruthyEnv(v: unknown): boolean {\n  if (v === undefined || v === null) return false;\n  const s = String(v).trim().toLowerCase();\n  return s !== \"\" && s !== \"0\" && s !== \"off\" && s !== \"false\" && s !== \"no\";\n}\n\nfunction forceRest(opts?: AnyOpts): boolean {\n  const fromOpts = (opts?.config as any)?.CAMUNDA_FORCE_REST;\n  const fromEnv =\n    typeof process !== \"undefined\" ? process.env?.CAMUNDA_FORCE_REST : undefined;\n  return isTruthyEnv(fromOpts) || isTruthyEnv(fromEnv);\n}\n\nfunction resolveMode(opts?: AnyOpts): NanoTransport {\n  // Explicit escape hatch for environments where WebSockets are blocked\n  // (e.g. corporate proxies). Wins over CAMUNDA_TRANSPORT.\n  if (forceRest(opts)) return \"rest\";\n  const fromOpts = opts?.config?.CAMUNDA_TRANSPORT as NanoTransport | undefined;\n  const fromEnv = (typeof process !== \"undefined\" ? process.env?.CAMUNDA_TRANSPORT : undefined) as NanoTransport | undefined;\n  return fromOpts ?? fromEnv ?? \"auto\";\n}\n\n/**\n * Default client-side submission-timeout (ms) for falcon creates, from\n * `opts.config.CAMUNDA_NANO_SUBMIT_TIMEOUT_MS` or the env var. `undefined`/`0`\n * means wait indefinitely under backpressure. A per-call `submitTimeoutMs` on\n * `createProcessInstance` input overrides this.\n */\nfunction resolveSubmitTimeoutMs(opts?: AnyOpts): number | undefined {\n  const raw =\n    (opts?.config?.CAMUNDA_NANO_SUBMIT_TIMEOUT_MS as string | number | undefined) ??\n    (typeof process !== \"undefined\" ? process.env?.CAMUNDA_NANO_SUBMIT_TIMEOUT_MS : undefined);\n  const n = typeof raw === \"string\" ? Number(raw) : raw;\n  return typeof n === \"number\" && Number.isFinite(n) && n > 0 ? n : undefined;\n}\n\n/**\n * Default Falcon handshake deadline (ms). Unlike the submit timeout (which\n * defaults to indefinite backpressure waiting), the connect timeout defaults to\n * a finite value: a stalled handshake must degrade to REST out of the box, not\n * hang. Generous enough not to trip a slow-but-legitimate handshake.\n */\nexport const DEFAULT_CONNECT_TIMEOUT_MS = 5000;\n\n/**\n * Falcon handshake deadline (ms), from `opts.config.CAMUNDA_NANO_CONNECT_TIMEOUT_MS`\n * or the env var, defaulting to {@link DEFAULT_CONNECT_TIMEOUT_MS}. An explicit\n * `0` (or negative/NaN) disables the deadline — wait indefinitely for `welcome`.\n */\nfunction resolveConnectTimeoutMs(opts?: AnyOpts): number | undefined {\n  const raw =\n    (opts?.config?.CAMUNDA_NANO_CONNECT_TIMEOUT_MS as string | number | undefined) ??\n    (typeof process !== \"undefined\" ? process.env?.CAMUNDA_NANO_CONNECT_TIMEOUT_MS : undefined);\n  if (raw === undefined || raw === null || raw === \"\") return DEFAULT_CONNECT_TIMEOUT_MS;\n  const n = typeof raw === \"string\" ? Number(raw) : raw;\n  if (!Number.isFinite(n) || n <= 0) return undefined; // explicit opt-out\n  return n;\n}\n\n/** restAddress is normalized by the SDK to end with /v2; strip it for base. */\nfunction baseFrom(restAddress: string): string {\n  return normalizeBase(restAddress).replace(/\\/v2$/, \"\");\n}\n\n/**\n * Drop-in replacement for the upstream createCamundaClient. Returns the upstream\n * client wrapped in a Proxy that upgrades createProcessInstance + createJobWorker\n * to the Falcon protocol when connected to a Nano server (overridable via the\n * CAMUNDA_TRANSPORT config: \"auto\" | \"falcon\" | \"rest\").\n */\nexport function createCamundaClient(opts?: AnyOpts): ReturnType<typeof createCamundaClientBase> {\n  const client = createCamundaClientBase(opts as any);\n  const mode = resolveMode(opts);\n  if (mode === \"rest\") return wrapRestPollDefault(client);\n\n  // Embedded (ADR 0005): bind the in-process μ-nano host directly — no detection,\n  // no socket. The host is the loopback \"Nano gateway in the same process\".\n  if (mode === \"embedded\") {\n    if (!opts?.embeddedHost) throw new Error(\"transport 'embedded' requires opts.embeddedHost\");\n    const transport = new EmbeddedTransport(opts.embeddedHost);\n    return wrapClient(client, async () => transport as any, () => transport.close());\n  }\n\n  const restAddress = client.getConfig().restAddress;\n  const base = baseFrom(restAddress);\n  const submitTimeoutMs = resolveSubmitTimeoutMs(opts);\n  const connectTimeoutMs = resolveConnectTimeoutMs(opts);\n\n  let transport: FalconTransport | null = null;\n  let nano: NanoInfo | null | undefined; // undefined = not yet probed\n  // Sticky \"falcon is not reachable\" flag. Once we've seen a WebSocket handshake\n  // fail (e.g. a proxy blocks WS), stop retrying for this client's lifetime and\n  // fall back to REST — matching the escape-hatch semantics of CAMUNDA_FORCE_REST.\n  let falconDead = false;\n  let warnedDead = false;\n  const markDead = (err: unknown) => {\n    falconDead = true;\n    transport = null;\n    if (!warnedDead) {\n      warnedDead = true;\n      // eslint-disable-next-line no-console\n      console.warn(\n        `[@nanobpm/sdk] Falcon transport unavailable (${(err as Error)?.message ?? err}); ` +\n          `falling back to REST. Set CAMUNDA_FORCE_REST=1 to skip Falcon detection.`,\n      );\n    }\n  };\n\n  const ensure = async (): Promise<FalconTransport | null> => {\n    if (falconDead) return null;\n    if (nano === undefined) {\n      nano = mode === \"falcon\"\n        ? { engine: \"nanobpmn\", falconPath: \"/falcon\" }\n        : await detectNano(base);\n    }\n    if (!nano) return null;\n    if (!transport) {\n      const t = new FalconTransport(base, nano.falconPath, submitTimeoutMs, connectTimeoutMs);\n      try {\n        // Eagerly open the socket so a proxy-blocked WebSocket surfaces here\n        // (single failure) rather than on every request.\n        await t.connect();\n        transport = t;\n      } catch (e) {\n        try { t.close(); } catch { /* ignore */ }\n        markDead(e);\n        return null;\n      }\n    }\n    return transport;\n  };\n  return wrapClient(client, ensure, () => transport?.close());\n}\n\n/** Wrap an upstream client, upgrading createProcessInstance + createJobWorker to\n *  a Nano transport (falcon or embedded) when `ensure` yields one. */\nfunction wrapClient(\n  client: ReturnType<typeof createCamundaClientBase>,\n  ensure: () => Promise<{ createInstance: Function; subscribe: Function; close: Function } | null>,\n  onStop?: () => void,\n): ReturnType<typeof createCamundaClientBase> {\n  return new Proxy(client, {\n    get(target, prop, receiver) {\n      if (prop === \"createProcessInstance\") {\n        return (input: any, options?: any) => {\n          return ensure().then(async (t) => {\n            if (!t) return (target as any).createProcessInstance(input, options);\n            try {\n              const r = await (t as any).createInstance({\n                processDefinitionId: input?.processDefinitionId,\n                processDefinitionKey: input?.processDefinitionKey,\n                variables: input?.variables,\n                awaitCompletion: input?.awaitCompletion ?? false,\n                fetchVariables: input?.fetchVariables,\n                submitTimeoutMs: input?.submitTimeoutMs,\n              });\n              if (r.status >= 400) throw new Error(`createInstance failed: ${r.status} ${JSON.stringify(r.body)}`);\n              if (r.completion) {\n                // The completion frame carries the authoritative key/variables;\n                // spread whatever the ack body had underneath it.\n                const base = r.body && typeof r.body === \"object\" ? r.body : {};\n                return { ...base, variables: r.completion.variables, processInstanceKey: r.completion.processInstanceKey };\n              }\n              // No mask: a success status with a missing/non-object body is a\n              // fault (the create returned no result to read a key from), so fail\n              // loud here instead of returning `{}` and deferring to a cryptic\n              // \"missing processInstanceKey/key\" downstream.\n              if (!r.body || typeof r.body !== \"object\") {\n                throw new Error(`createInstance: gateway returned success (${r.status}) with no result body`);\n              }\n              return r.body;\n            } catch (e) {\n              // A late WS failure (proxy severed the socket, reconnect denied) —\n              // fall back to REST for this call.\n              const msg = String((e as Error)?.message ?? e);\n              if (msg.includes(\"falcon closed\") || msg.includes(\"falcon connect failed\")) {\n                // eslint-disable-next-line no-console\n                console.warn(`[@nanobpm/sdk] Falcon call failed (${msg}); retrying via REST.`);\n                return (target as any).createProcessInstance(input, options);\n              }\n              throw e;\n            }\n          });\n        };\n      }\n      if (prop === \"createJobWorker\") {\n        return (cfg: any) => {\n          const w = new NanoJobWorker(null, { ...cfg, autoStart: false });\n          void ensure().then(async (t) => {\n            if (!t) { (target as any).createJobWorker(withRestPollDefault(cfg)); return; }\n            w.bindTransport(t as any);\n            try { await w.start(); }\n            catch (e) {\n              // eslint-disable-next-line no-console\n              console.warn(\n                `[@nanobpm/sdk] Falcon subscribe failed (${(e as Error)?.message ?? e}); ` +\n                  `falling back to REST job worker.`,\n              );\n              (target as any).createJobWorker(withRestPollDefault(cfg));\n            }\n          });\n          return w as any;\n        };\n      }\n      if (prop === \"stopAllWorkers\") {\n        return () => { onStop?.(); return (target as any).stopAllWorkers(); };\n      }\n      return Reflect.get(target, prop, receiver);\n    },\n  });\n}\n\nexport default createCamundaClient;\n","// Nano server detection. A nanobpmn gateway advertises itself on GET /v2/topology\n// via a `nano` block carrying `falconPath`. Detecting it lets the SDK\n// upgrade to the Falcon protocol transparently.\n\nexport interface NanoInfo {\n  engine: string;\n  version?: string;\n  falconPath: string;\n}\n\nconst cache = new Map<string, NanoInfo | null>();\n\n/** Strip trailing slashes for a tidy REST base. */\nexport function normalizeBase(url: string): string {\n  return url.replace(/\\/+$/, \"\");\n}\n\n/**\n * Probes a REST address once (cached) and returns the Nano engine info if the\n * server is a nanobpmn gateway, or null for stock Camunda 8.\n */\nexport async function detectNano(\n  restAddress: string,\n  fetchImpl: typeof fetch = fetch,\n): Promise<NanoInfo | null> {\n  const base = normalizeBase(restAddress);\n  if (cache.has(base)) return cache.get(base) ?? null;\n  let info: NanoInfo | null = null;\n  try {\n    const res = await fetchImpl(`${base}/v2/topology`, {\n      headers: { accept: \"application/json\" },\n    });\n    if (res.ok) {\n      const body: any = await res.json().catch(() => ({}));\n      const nano = body?.nano;\n      if (nano && typeof nano.falconPath === \"string\") {\n        info = {\n          engine: String(nano.engine ?? \"nanobpmn\"),\n          version: nano.version ? String(nano.version) : undefined,\n          falconPath: nano.falconPath,\n        };\n      }\n    }\n  } catch {\n    info = null; // unreachable or non-Nano: stay on REST\n  }\n  cache.set(base, info);\n  return info;\n}\n\n/** Test seam: clear the detection cache. */\nexport function _clearDetectionCache(): void {\n  cache.clear();\n}\n\n/** Build the ws:// or wss:// falcon URL from a REST base + path. */\nexport function falconUrl(restAddress: string, path: string): string {\n  let base = normalizeBase(restAddress);\n  if (base.startsWith(\"https://\")) base = \"wss://\" + base.slice(\"https://\".length);\n  else if (base.startsWith(\"http://\")) base = \"ws://\" + base.slice(\"http://\".length);\n  return base + path;\n}\n","// Picks a WebSocket implementation: native global (browser/Deno/Node>=22) or the\n// `ws` package (Node). Kept tiny so the transport stays runtime-agnostic.\nexport type WSImpl = new (url: string) => WebSocket;\n\nlet cached: WSImpl | null = null;\n\nexport async function getWebSocket(): Promise<WSImpl> {\n  if (cached) return cached;\n  const g = globalThis as { WebSocket?: WSImpl };\n  if (typeof g.WebSocket === \"function\") {\n    cached = g.WebSocket;\n    return cached;\n  }\n  const mod = await import(\"ws\");\n  cached = (mod.default ?? mod.WebSocket) as unknown as WSImpl;\n  return cached;\n}\n","// Falcon transport. One persistent WebSocket per client multiplexes:\n//   * createInstance (corr-correlated CommandResult, awaitCompletion via\n//     InstanceCompleted), and\n//   * job subscriptions (welcome -> subscribe -> job push, replenished by\n//     jobCredits; completeJob/failJob/throwError as unmetered drains).\n// Mirrors the protocol implemented by the engine (server/src/falcon.rs)\n// and the embedded worker SDK (server/src/console/worker_sdk.ts).\nimport { falconUrl } from \"./detect.js\";\nimport { getWebSocket } from \"./ws.js\";\n\ntype Json = Record<string, unknown>;\n\n/**\n * Sleep for `ms`, without keeping the event loop alive on its own — a pending\n * reconnect backoff must never stop the host process from exiting.\n */\nfunction sleep(ms: number): Promise<void> {\n  return new Promise((resolve) => {\n    const t = setTimeout(resolve, ms) as { unref?: () => void };\n    t.unref?.();\n  });\n}\n\n/** Resolved shape of an `awaitCompletion` create. */\ntype Completion = { processCompleted: boolean; variables: unknown; processInstanceKey: string };\n\n/**\n * Thrown when the gateway sends a result frame that omits a field the protocol\n * requires (a `commandResult` with no `status`, an `instanceCompleted` with no\n * `processInstanceKey`). The absence of such an identity/status field always\n * signals an engine or protocol fault, so we surface it loudly here rather than\n * laundering it into a benign-looking default (\"\" / 0) that would pass\n * downstream checks and corrupt state (e.g. completing a job against key \"\").\n */\nexport class MalformedFrameError extends Error {\n  constructor(\n    readonly frameType: string,\n    readonly detail: string,\n    readonly frame: unknown,\n  ) {\n    super(`malformed Falcon '${frameType}' frame: ${detail}`);\n    this.name = \"MalformedFrameError\";\n  }\n}\n\nexport interface JobFrame {\n  jobKey: string;\n  type: string;\n  processInstanceKey: string;\n  variables?: Record<string, unknown>;\n  customHeaders?: Record<string, unknown>;\n  retries?: number;\n  [k: string]: unknown;\n}\n\nexport interface Subscription {\n  jobType: string;\n  worker: string;\n  credits: number;\n  timeoutMs: number;\n  fetchVariables: string[] | null;\n  onJob: (job: JobFrame) => void;\n}\n\ninterface Pending {\n  resolve: (v: { status: number; body: unknown }) => void;\n  reject: (e: unknown) => void;\n}\n\n/** A create queued behind an exhausted submission-credit window. */\ninterface CreditWaiter {\n  /** Grant a credit (decrement + resolve the acquire). */\n  grant: () => void;\n  /** Abandon the wait (reject the acquire), e.g. on disconnect. */\n  fail: (e: unknown) => void;\n}\n\n/**\n * Thrown by {@link FalconTransport.createInstance} when the gateway does\n * not acknowledge a create within `submitTimeoutMs`. On the Falcon protocol,\n * admission backpressure is expressed by the server withholding submission\n * credits (no `503`, no retry), so a create otherwise waits indefinitely for\n * intake capacity. This turns that stall into a typed rejection — treat it as\n * \"the server is backpressured\" and back off; do not tight-loop retry.\n */\nexport class SubmissionTimeoutError extends Error {\n  constructor(readonly timeoutMs: number) {\n    super(\n      `create submission stalled: no gateway ack within ${timeoutMs}ms (server is applying admission backpressure)`,\n    );\n    this.name = \"SubmissionTimeoutError\";\n  }\n}\n\n/**\n * Thrown by {@link FalconTransport.connect} when a WebSocket is opened but the\n * gateway does not complete the Falcon handshake (no `welcome` frame) within\n * `connectTimeoutMs`.\n *\n * Most WebSocket-hostile infrastructure (corporate proxies, HTTP-only ingress,\n * some L7 load balancers / WAFs) *rejects* the upgrade, which surfaces promptly\n * as a socket error and falls back to REST. But a proxy that *blackholes* the\n * upgrade — accepting the TCP connection and forwarding bytes while the\n * handshake never completes and no `welcome` ever arrives — would otherwise\n * leave {@link FalconTransport.connect} pending forever, hanging the first\n * request instead of degrading to REST. A bounded connect deadline turns that\n * silent stall into this typed rejection so the caller can fall back to REST.\n */\nexport class ConnectTimeoutError extends Error {\n  constructor(readonly timeoutMs: number) {\n    super(\n      `falcon handshake stalled: no gateway welcome within ${timeoutMs}ms (WebSocket upgrade may be blocked or blackholed by intervening infra)`,\n    );\n    this.name = \"ConnectTimeoutError\";\n  }\n}\n\nexport class FalconTransport {\n  private url: string;\n  private ws: WebSocket | null = null;\n  private open = false;\n  private corr = 0;\n  private pending = new Map<number, Pending>();\n  private awaits = new Map<number, { resolve: (v: Completion) => void; reject: (e: unknown) => void }>();\n  private subs = new Map<string, Subscription>();\n  private heartbeat: ReturnType<typeof setInterval> | null = null;\n  private connectPromise: Promise<void> | null = null;\n  private closed = false;\n  /** True while the background reconnect loop is running (prevents overlap). */\n  private reconnecting = false;\n  private defaultSubmitTimeoutMs?: number;\n  /**\n   * Client-side bound (ms) on the Falcon handshake: how long a freshly-opened\n   * WebSocket may take to deliver its `welcome` before {@link connect} rejects\n   * with {@link ConnectTimeoutError}. Guards against WebSocket-hostile infra\n   * that blackholes the upgrade (accepts the socket but never completes the\n   * handshake). `undefined`/`0` waits indefinitely (legacy behaviour).\n   */\n  private connectTimeoutMs?: number;\n  /**\n   * Server-granted submission-credit window (seeded by `welcome`, topped up by\n   * `submissionCredits`). Mirrors the engine's intake metering so creates queue\n   * client-side under admission backpressure instead of flooding the gateway.\n   */\n  private credits = 0;\n  private creditWaiters: CreditWaiter[] = [];\n\n  constructor(\n    restAddress: string,\n    path: string,\n    defaultSubmitTimeoutMs?: number,\n    connectTimeoutMs?: number,\n  ) {\n    this.url = falconUrl(restAddress, path);\n    this.defaultSubmitTimeoutMs =\n      defaultSubmitTimeoutMs !== undefined && defaultSubmitTimeoutMs > 0\n        ? defaultSubmitTimeoutMs\n        : undefined;\n    this.connectTimeoutMs =\n      connectTimeoutMs !== undefined && connectTimeoutMs > 0 ? connectTimeoutMs : undefined;\n  }\n\n  private nextCorr(): number {\n    this.corr += 1;\n    return this.corr;\n  }\n\n  /**\n   * Parse the correlation id off an inbound result frame. A missing or invalid\n   * `corr` cannot be routed to its waiting caller, so we drop the frame with a\n   * warning rather than coercing it to 0 (which would silently target — or\n   * mis-target — whatever call happens to hold corr 0).\n   */\n  private frameCorr(frameType: string, f: Json): number | null {\n    const c = f.corr;\n    if (typeof c !== \"number\" || !Number.isInteger(c) || c <= 0) {\n      // eslint-disable-next-line no-console\n      console.warn(`[@nanobpm/sdk] dropping malformed Falcon '${frameType}' frame: missing/invalid 'corr'`);\n      return null;\n    }\n    return c;\n  }\n\n  private send(frame: Json): void {\n    if (this.ws && this.open) this.ws.send(JSON.stringify(frame));\n  }\n\n  async connect(): Promise<void> {\n    if (this.closed) throw new Error(`falcon closed: ${this.url}`);\n    if (this.open) return;\n    if (this.connectPromise) {\n      // A connect/reconnect is already in flight. Wait for it, then verify we are\n      // actually open: the shared reconnect promise also resolves when the loop\n      // exits because close() was called, so callers must not proceed (e.g. into\n      // acquireCredit) against a transport that never came back up.\n      await this.connectPromise;\n      if (!this.open) throw new Error(`falcon not connected (closed or unavailable): ${this.url}`);\n      return;\n    }\n    // Initial connect fails fast (reject-once) so the caller sees a prompt error,\n    // but we clear the shared promise on failure so a later call can retry cleanly.\n    // Only *background* reconnects (after an unexpected drop) loop; see\n    // {@link reconnectLoop}.\n    const p = this.openSocket();\n    this.connectPromise = p;\n    try {\n      await p;\n    } catch (e) {\n      if (this.connectPromise === p) this.connectPromise = null;\n      throw e;\n    }\n  }\n\n  /**\n   * Open exactly one WebSocket and settle the returned promise once: resolve on\n   * the server's `welcome`, reject if the socket errors or closes before it\n   * arrives. Also wires the persistent handlers, including the unexpected-close\n   * hook that drives the background reconnect loop.\n   *\n   * The returned promise's rejection is *owned by the caller*: the initial\n   * {@link connect} surfaces it (fast first failure), while {@link reconnectLoop}\n   * catches and retries it. The close handler itself never lets a failure escape\n   * as an unhandled rejection — that was the crash bug this fixes.\n   */\n  private openSocket(): Promise<void> {\n    return new Promise<void>((resolve, reject) => {\n      let settled = false;\n      // Whether this socket reached `welcome`. Only an *established* connection\n      // dropping should start the background reconnect loop — a failed initial\n      // dial rejects once (see {@link connect}) and must not spin up a hidden\n      // loop the caller never asked for.\n      let reachedWelcome = false;\n      // Bound the handshake: if `welcome` never arrives (a blackholing proxy\n      // that upgrades the socket but stalls), reject with ConnectTimeoutError so\n      // the caller can fall back to REST instead of hanging forever. Cleared on\n      // the first settle (either outcome).\n      let connectTimer: ReturnType<typeof setTimeout> | undefined;\n      const clearConnectTimer = () => {\n        if (connectTimer !== undefined) {\n          clearTimeout(connectTimer);\n          connectTimer = undefined;\n        }\n      };\n      const settleResolve = () => {\n        reachedWelcome = true;\n        clearConnectTimer();\n        if (!settled) {\n          settled = true;\n          resolve();\n        }\n      };\n      const settleReject = (e: unknown) => {\n        clearConnectTimer();\n        if (!settled) {\n          settled = true;\n          reject(e);\n        }\n      };\n      void getWebSocket()\n        .then((WS) => {\n          const ws = new WS(this.url);\n          this.ws = ws;\n          // Arm the handshake deadline now that a socket exists. On expiry,\n          // reject this attempt and close the socket; `reachedWelcome` stays\n          // false, so the onclose handler will NOT start a background reconnect\n          // (a stalled handshake is treated like a failed initial dial).\n          if (this.connectTimeoutMs !== undefined) {\n            connectTimer = setTimeout(() => {\n              if (this.ws !== ws || settled) return;\n              settleReject(new ConnectTimeoutError(this.connectTimeoutMs!));\n              // De-select this socket before closing so any late frame (e.g. a\n              // `welcome` that arrives between reject and the socket actually\n              // closing) is ignored by `onmessage`/`onclose` (`this.ws !== ws`)\n              // and cannot flip `this.open`/arm the heartbeat after connect has\n              // already failed with ConnectTimeoutError.\n              this.ws = null;\n              try {\n                ws.close(4408, \"handshake timeout\");\n              } catch {\n                /* ignore */\n              }\n            }, this.connectTimeoutMs);\n            (connectTimer as { unref?: () => void }).unref?.();\n          }\n          ws.onopen = () => {};\n          ws.onmessage = (ev: MessageEvent) => {\n            // Ignore frames from a socket that has already been superseded by a\n            // newer attempt — a stale message must not mutate current-connection\n            // state (e.g. a late `welcome` flipping `open`).\n            if (this.ws !== ws) return;\n            let f: Json;\n            try {\n              f = JSON.parse(typeof ev.data === \"string\" ? ev.data : String(ev.data));\n            } catch {\n              return;\n            }\n            this.handle(f, settleResolve);\n          };\n          ws.onerror = () => {\n            if (this.ws !== ws) return;\n            // A pre-`welcome` error means this attempt failed. Once we are open,\n            // errors surface via `onclose`, which drives the reconnect instead.\n            if (!this.open) settleReject(new Error(`falcon connect failed: ${this.url}`));\n          };\n          ws.onclose = () => {\n            // A superseded socket (replaced by a newer attempt) may close late —\n            // after the error→close gap. It must NOT reset `open`, clear credits,\n            // or reject the pending work of the *current* connection. Settle only\n            // its own (already-settled) attempt promise and bail.\n            if (this.ws !== ws) {\n              settleReject(new Error(`falcon closed before ready: ${this.url}`));\n              return;\n            }\n            this.open = false;\n            // The credit window is per-connection; a reconnect's welcome re-grants\n            // a fresh one. Reset it and fail any queued creates so callers can retry\n            // on the new connection rather than hang against a dead window.\n            this.credits = 0;\n            const waiters = this.creditWaiters;\n            this.creditWaiters = [];\n            for (const w of waiters) w.fail(new Error(\"falcon closed\"));\n            if (this.heartbeat) clearInterval(this.heartbeat);\n            // Reject in-flight commands; subscriptions resubscribe on reconnect.\n            for (const p of this.pending.values()) p.reject(new Error(\"falcon closed\"));\n            this.pending.clear();\n            // Also fail any awaited completions: a dropped socket must surface as\n            // a rejection, not leave the completion promise hanging forever.\n            for (const a of this.awaits.values()) a.reject(new Error(\"falcon closed\"));\n            this.awaits.clear();\n            // If this socket never reached `welcome`, settle the attempt as failed\n            // so its awaiter (initial connect or the reconnect loop) moves on.\n            settleReject(new Error(`falcon closed before ready: ${this.url}`));\n            // Only recover from the drop of an *established* connection. A failed\n            // initial dial already rejected once and must not start a background\n            // loop; the reconnect loop drives its own retries via its while-loop,\n            // so its failed attempts don't need to re-trigger here either.\n            if (!this.closed && reachedWelcome) this.beginReconnect();\n          };\n        })\n        .catch((e) => settleReject(e));\n    });\n  }\n\n  /**\n   * Start the background reconnect loop after an unexpected close — unless one is\n   * already running or the client was explicitly closed. Publishes a shared\n   * `connectPromise` so any {@link connect} awaiter during the outage latches onto\n   * the in-flight reconnect instead of racing a second socket.\n   */\n  private beginReconnect(): void {\n    if (this.reconnecting || this.closed) return;\n    this.reconnecting = true;\n    let settleShared!: () => void;\n    this.connectPromise = new Promise<void>((res) => {\n      settleShared = res;\n    });\n    void this.reconnectLoop(settleShared);\n  }\n\n  /**\n   * Reconnect with bounded exponential backoff + jitter until the engine returns\n   * (or {@link close} is called). Each attempt's failure is caught here, so a\n   * transient failure while the engine is down never escapes as an unhandled\n   * rejection. On success the server's `welcome` re-subscribes every active\n   * worker (see {@link handle}).\n   */\n  private async reconnectLoop(settleShared: () => void): Promise<void> {\n    let attempt = 0;\n    while (!this.closed && !this.open) {\n      await sleep(this.backoffDelayMs(attempt));\n      attempt += 1;\n      if (this.closed || this.open) break;\n      try {\n        await this.openSocket();\n      } catch {\n        // Transient (engine still down): back off and retry. Swallowed on purpose.\n        continue;\n      }\n    }\n    this.reconnecting = false;\n    if (!this.open) this.connectPromise = null;\n    settleShared();\n  }\n\n  /**\n   * Full-jitter (AWS-style) exponential backoff: a uniform random delay in\n   * `[0, ceiling)`, where the ceiling grows exponentially from `base` (100ms)\n   * and is clamped to `cap` (2s). Full jitter spreads retries out to avoid a\n   * thundering herd while still keeping a long outage retrying at a steady\n   * ceiling rather than backing off forever.\n   */\n  private backoffDelayMs(attempt: number): number {\n    const base = 100;\n    const cap = 2000;\n    const ceiling = Math.min(cap, base * 2 ** attempt);\n    return Math.floor(Math.random() * ceiling);\n  }\n\n  private handle(f: Json, resolveConnect: () => void): void {\n    switch (f.type) {\n      case \"welcome\": {\n        this.open = true;\n        this.credits = Number(f.submissionCredits ?? 0);\n        const hb = Number(f.heartbeatMs ?? 0);\n        if (hb > 0) {\n          if (this.heartbeat) clearInterval(this.heartbeat);\n          this.heartbeat = setInterval(() => this.send({ type: \"heartbeat\" }), hb);\n        }\n        for (const sub of this.subs.values()) this.sendSubscribe(sub);\n        this.releaseCreditWaiters();\n        resolveConnect();\n        break;\n      }\n      case \"job\": {\n        const job = (f.job as JobFrame) ?? ({} as JobFrame);\n        const sub = this.subs.get(job.type);\n        if (sub) sub.onJob(job);\n        break;\n      }\n      case \"commandResult\": {\n        const corr = this.frameCorr(\"commandResult\", f);\n        if (corr === null) break;\n        const p = this.pending.get(corr);\n        if (!p) break;\n        this.pending.delete(corr);\n        // `status` is required: a missing status must never be read as 0 (a\n        // success code below 400), which would surface an empty result as a\n        // completed create.\n        if (typeof f.status !== \"number\" || !Number.isFinite(f.status)) {\n          p.reject(new MalformedFrameError(\"commandResult\", \"missing numeric 'status'\", f));\n          break;\n        }\n        // Pass the body through unchanged (no `?? null` mask): an absent body on\n        // a success status is a fault the caller must see, not paper over.\n        p.resolve({ status: f.status, body: f.body });\n        break;\n      }\n      case \"instanceCompleted\": {\n        const corr = this.frameCorr(\"instanceCompleted\", f);\n        if (corr === null) break;\n        const cb = this.awaits.get(corr);\n        if (!cb) break;\n        this.awaits.delete(corr);\n        // `processInstanceKey` is required. Tolerate a numeric key (coerce to the\n        // string the protocol uses) but reject an absent/empty one rather than\n        // masking it as \"\" — a completion with no key is a protocol violation.\n        const rawKey = f.processInstanceKey;\n        const key = typeof rawKey === \"number\" ? String(rawKey) : rawKey;\n        if (typeof key !== \"string\" || key.length === 0) {\n          cb.reject(new MalformedFrameError(\"instanceCompleted\", \"missing 'processInstanceKey'\", f));\n          break;\n        }\n        cb.resolve({\n          processCompleted: Boolean(f.processCompleted),\n          variables: f.variables ?? {},\n          processInstanceKey: key,\n        });\n        break;\n      }\n      case \"submissionCredits\": {\n        // Server topped up the create-side window (intake headroom returned).\n        this.credits += Number(f.n ?? 0);\n        this.releaseCreditWaiters();\n        break;\n      }\n      // pressure / heartbeat: no client action needed yet.\n    }\n  }\n\n  private sendSubscribe(sub: Subscription): void {\n    this.send({\n      type: \"subscribe\",\n      jobType: sub.jobType,\n      jobCredits: sub.credits,\n      worker: sub.worker,\n      timeout: sub.timeoutMs,\n      fetchVariable: sub.fetchVariables,\n    });\n  }\n\n  // ----- submission-credit gating ------------------------------------------\n\n  /**\n   * Take one submission credit, waiting for the gateway to replenish when the\n   * window is exhausted (admission backpressure — no 503, no retry). When\n   * `timeoutMs` elapses first, rejects with {@link SubmissionTimeoutError} and\n   * removes the queued waiter so no credit slot leaks.\n   */\n  private acquireCredit(timeoutMs?: number): Promise<void> {\n    if (this.credits > 0) {\n      this.credits -= 1;\n      return Promise.resolve();\n    }\n    return new Promise<void>((resolve, reject) => {\n      let timer: ReturnType<typeof setTimeout> | undefined;\n      const waiter: CreditWaiter = {\n        grant: () => {\n          if (timer !== undefined) clearTimeout(timer);\n          this.credits -= 1;\n          resolve();\n        },\n        fail: (e) => {\n          if (timer !== undefined) clearTimeout(timer);\n          reject(e);\n        },\n      };\n      this.creditWaiters.push(waiter);\n      if (timeoutMs !== undefined && timeoutMs > 0) {\n        timer = setTimeout(() => {\n          const i = this.creditWaiters.indexOf(waiter);\n          if (i >= 0) this.creditWaiters.splice(i, 1);\n          reject(new SubmissionTimeoutError(timeoutMs));\n        }, timeoutMs);\n        timer.unref?.();\n      }\n    });\n  }\n\n  private releaseCreditWaiters(): void {\n    while (this.credits > 0 && this.creditWaiters.length > 0) {\n      this.creditWaiters.shift()!.grant();\n    }\n  }\n\n  /** Create a process instance over the stream. */\n  async createInstance(input: {\n    processDefinitionId?: string;\n    processDefinitionKey?: string;\n    variables?: Record<string, unknown>;\n    awaitCompletion?: boolean;\n    fetchVariables?: string[];\n    requestTimeoutMs?: number;\n    /**\n     * Client-side bound (ms) on how long to wait for a submission credit before\n     * rejecting with {@link SubmissionTimeoutError}. Overrides the transport-wide\n     * default. Omit to wait indefinitely under backpressure. Never sent on the\n     * wire — the server is unaware of it.\n     */\n    submitTimeoutMs?: number;\n  }): Promise<{ status: number; body: unknown; completion?: { processCompleted: boolean; variables: unknown; processInstanceKey: string } }> {\n    await this.connect();\n    // Gate intake on the server's submission-credit window: block here (bounded\n    // by submitTimeoutMs) rather than flooding the gateway with creates it has\n    // not granted capacity for. Job completions stay unmetered.\n    await this.acquireCredit(input.submitTimeoutMs ?? this.defaultSubmitTimeoutMs);\n    const corr = this.nextCorr();\n    let completionResolve: ((v: Completion) => void) | null = null;\n    let completionReject: ((e: unknown) => void) | null = null;\n    const completion = input.awaitCompletion\n      ? new Promise<Completion>((resolve, reject) => {\n          completionResolve = resolve;\n          completionReject = reject;\n        })\n      : null;\n    if (completionResolve && completionReject) {\n      this.awaits.set(corr, { resolve: completionResolve, reject: completionReject });\n      // If the completion waiter is rejected (command failure or socket close)\n      // before we reach the `await completion` below, that rejection must not be\n      // an unhandled rejection — attach a benign handler now. `await completion`\n      // still throws on rejection as normal.\n      completion?.catch(() => {});\n    }\n    let result: { status: number; body: unknown };\n    try {\n      result = await new Promise<{ status: number; body: unknown }>((resolve, reject) => {\n        this.pending.set(corr, { resolve, reject });\n        this.send({\n          type: \"createInstance\",\n          corr,\n          processDefinitionId: input.processDefinitionId ?? null,\n          processDefinitionKey: input.processDefinitionKey ?? null,\n          variables: input.variables ?? null,\n          awaitCompletion: input.awaitCompletion ?? false,\n          fetchVariables: input.fetchVariables ?? null,\n          requestTimeout: input.requestTimeoutMs ?? null,\n        });\n      });\n    } catch (err) {\n      // The command failed at or before ack (malformed frame, or socket close).\n      // Tear down the completion waiter so it neither leaks in `this.awaits` nor\n      // is left to reject on a later socket close.\n      const waiter = this.awaits.get(corr);\n      if (waiter) {\n        this.awaits.delete(corr);\n        waiter.reject(err);\n      }\n      throw err;\n    }\n    const done = completion ? await completion : undefined;\n    return { ...result, completion: done };\n  }\n\n  /** Subscribe a worker; jobs arrive via sub.onJob, credits replenished by ack helpers. */\n  async subscribe(sub: Subscription): Promise<void> {\n    // If the transport is not yet open, connect() resolves via the `welcome`\n    // handler, which already sends a subscribe frame for every active sub\n    // (including this one, since we register it below). Sending again here would\n    // duplicate the frame, so only send ourselves when the transport was already\n    // open — where no `welcome` fires for this call.\n    const wasOpen = this.open;\n    this.subs.set(sub.jobType, sub);\n    await this.connect();\n    // If unsubscribe()/stop() ran while we awaited connect(), this sub was\n    // removed and the caller no longer wants it. Sending the subscribe now would\n    // make the gateway activate/push jobs the client has no handler for and would\n    // silently drop (unsubscribe() only clears the local handler), causing\n    // avoidable job timeouts. Only send while this sub is still the active one.\n    if (this.subs.get(sub.jobType) !== sub) return;\n    if (wasOpen) this.sendSubscribe(sub);\n  }\n\n  unsubscribe(jobType: string): void {\n    this.subs.delete(jobType);\n  }\n\n  completeJob(jobKey: string, variables?: Record<string, unknown>): void {\n    this.send({ type: \"completeJob\", corr: this.nextCorr(), jobKey, variables: variables ?? null });\n  }\n  failJob(jobKey: string, retries?: number, errorMessage?: string): void {\n    this.send({ type: \"failJob\", corr: this.nextCorr(), jobKey, retries: retries ?? 0, errorMessage: errorMessage ?? null });\n  }\n  throwError(jobKey: string, errorCode: string, errorMessage?: string): void {\n    this.send({ type: \"throwError\", corr: this.nextCorr(), jobKey, errorCode, errorMessage: errorMessage ?? null });\n  }\n  grantCredits(jobType: string, n: number): void {\n    this.send({ type: \"jobCredits\", jobType, n });\n  }\n\n  close(): void {\n    this.closed = true;\n    if (this.heartbeat) clearInterval(this.heartbeat);\n    try {\n      this.ws?.close(1000, \"shutdown\");\n    } catch {\n      /* ignore */\n    }\n  }\n}\n","// Embedded transport (ADR 0005, realization (a) \"in-process direct\"). Satisfies\n// the same operation seam the falcon transport does — createInstance,\n// subscribe/job-push, complete/fail/throw — but calls an in-process EmbeddedHost\n// directly: no sockets, no frames. The host wraps μ-nano.wasm and owns the wall\n// clock, durable journal, and timer/dispatch tick. Bind it via\n// createCamundaClient({ config: { CAMUNDA_TRANSPORT: \"embedded\" }, embeddedHost }).\n\n/** A job handed to a worker handler. */\nexport interface EmbeddedJob {\n  jobKey: string;\n  type: string;\n  processInstanceKey: string;\n  elementId: string;\n  retries: number;\n  variables: Record<string, unknown>;\n}\n\n/** The engine bridge an embedded app supplies (implemented by the app template's\n *  Deno host around μ-nano.wasm). Engine-core stays clock-free; the host injects now(). */\nexport interface EmbeddedHost {\n  deploy(xml: string): Promise<{ processIds: string[] }>;\n  createInstance(input: { processDefinitionId?: string; variables?: Record<string, unknown> }): Promise<{ processInstanceKey: string }>;\n  activateJobs(type: string, max: number, timeoutMs: number, worker: string): Promise<EmbeddedJob[]>;\n  completeJob(jobKey: string, variables?: Record<string, unknown>): Promise<void>;\n  failJob(jobKey: string, retries: number, errorMessage?: string): Promise<void>;\n  instanceCompleted(key: string): boolean;\n  instanceVariables(key: string): Record<string, unknown>;\n  /** Drive timers + dispatch at wall-clock now. Returns true if anything changed. */\n  tick(): void;\n}\n\ninterface Sub {\n  worker: string;\n  credits: number;\n  timeoutMs: number;\n  onJob: (j: EmbeddedJob) => void;\n}\n\n/** Pull-based dispatch over an in-process host. Structurally compatible with the\n *  FalconTransport surface used by NanoJobWorker + the client proxy. */\nexport class EmbeddedTransport {\n  private subs = new Map<string, Sub>();\n  private timer: ReturnType<typeof setInterval> | null = null;\n  private closed = false;\n  constructor(private host: EmbeddedHost, private pollMs = 25) {}\n\n  async createInstance(input: {\n    processDefinitionId?: string;\n    variables?: Record<string, unknown>;\n    awaitCompletion?: boolean;\n  }): Promise<{ status: number; body: unknown; completion?: { processCompleted: boolean; variables: unknown; processInstanceKey: string } }> {\n    const { processInstanceKey } = await this.host.createInstance(input);\n    this.host.tick();\n    if (!input.awaitCompletion) return { status: 200, body: { processInstanceKey } };\n    while (!this.host.instanceCompleted(processInstanceKey)) {\n      this.host.tick();\n      await new Promise((r) => setTimeout(r, this.pollMs));\n    }\n    const variables = this.host.instanceVariables(processInstanceKey);\n    return { status: 200, body: { processInstanceKey, variables }, completion: { processCompleted: true, variables, processInstanceKey } };\n  }\n\n  async subscribe(sub: { jobType: string; worker: string; credits: number; timeoutMs: number; onJob: (j: any) => void }): Promise<void> {\n    this.subs.set(sub.jobType, { worker: sub.worker, credits: sub.credits, timeoutMs: sub.timeoutMs, onJob: sub.onJob });\n    if (!this.timer) this.timer = setInterval(() => this.pump(), this.pollMs);\n  }\n  unsubscribe(jobType: string): void { this.subs.delete(jobType); }\n\n  completeJob(jobKey: string, variables?: Record<string, unknown>): void { void this.host.completeJob(jobKey, variables); }\n  failJob(jobKey: string, retries?: number, errorMessage?: string): void { void this.host.failJob(jobKey, retries ?? 0, errorMessage); }\n  throwError(jobKey: string, _errorCode: string, errorMessage?: string): void { void this.host.failJob(jobKey, 0, errorMessage); }\n  grantCredits(jobType: string, n: number): void { const s = this.subs.get(jobType); if (s) s.credits += n; }\n\n  private pump(): void {\n    if (this.closed) return;\n    this.host.tick();\n    for (const [type, s] of this.subs) {\n      if (s.credits <= 0) continue;\n      void this.host.activateJobs(type, s.credits, s.timeoutMs, s.worker).then((jobs) => {\n        for (const j of jobs) { if (s.credits <= 0) break; s.credits--; s.onJob(j); }\n      });\n    }\n  }\n  close(): void { this.closed = true; if (this.timer) clearInterval(this.timer); }\n}\n","// A drop-in JobWorker that drains jobs over the Falcon protocol. Mirrors the\n// upstream JobWorker surface (close/stop/start) but receives pushed jobs instead\n// of long-polling. Each handled job replenishes one credit, keeping demand at\n// maxParallelJobs.\nimport { JobActionReceiptSymbol as JobActionReceipt } from \"@camunda8/orchestration-cluster-api\";\nimport type { FalconTransport, JobFrame } from \"./transport.js\";\n\nexport interface NanoJobWorkerConfig {\n  jobType: string;\n  workerName?: string;\n  maxParallelJobs?: number;\n  jobTimeoutMs?: number;\n  fetchVariables?: string[];\n  jobHandler: (job: any) => Promise<typeof JobActionReceipt> | typeof JobActionReceipt;\n  autoStart?: boolean;\n}\n\nexport class NanoJobWorker {\n  private transport: FalconTransport | null;\n  private cfg: NanoJobWorkerConfig;\n  private credits: number;\n  private stopped = false;\n  /** A start() was requested while transport was still null; replay it on bind. */\n  private startRequested = false;\n  /**\n   * Shared in-flight (then settled) subscription attempt. Racing start() /\n   * bindTransport() calls await this same promise so they subscribe at most\n   * once, yet all observe a rejection if transport.subscribe() fails. Reset to\n   * null on failure (so a later start() can retry) and on stop().\n   */\n  private subscribePromise: Promise<void> | null = null;\n  readonly jobType: string;\n  readonly name: string;\n\n  constructor(transport: FalconTransport | null, cfg: NanoJobWorkerConfig) {\n    this.transport = transport;\n    this.cfg = cfg;\n    this.jobType = cfg.jobType;\n    this.credits = cfg.maxParallelJobs ?? 10;\n    this.name = cfg.workerName || `worker-${cfg.jobType}`;\n    if (cfg.autoStart !== false) void this.start();\n  }\n\n  /**\n   * Bind the transport after async Nano detection. If a start() was requested\n   * before the transport was available, honour it now (subscribing exactly once).\n   */\n  bindTransport(transport: FalconTransport): void {\n    this.transport = transport;\n    // Fire-and-forget replay: swallow the rejection on this path so it is not an\n    // unhandled rejection. The shared subscribePromise still rejects, so any\n    // start() awaiting it observes the failure (and resets for a later retry).\n    if (this.startRequested && !this.stopped) void this.subscribe().catch(() => {});\n  }\n\n  /**\n   * Begin draining jobs. Null-safe and idempotent:\n   * - If the transport is not yet bound, the request is deferred and replayed by\n   *   bindTransport() once it is — no null dereference.\n   * - Duplicate calls (e.g. proxy self-start plus an eager caller start) result\n   *   in a single subscription.\n   */\n  async start(): Promise<void> {\n    this.stopped = false;\n    this.startRequested = true;\n    if (this.transport === null) return; // deferred until bindTransport()\n    await this.subscribe();\n  }\n\n  private subscribe(): Promise<void> {\n    if (this.stopped || this.transport === null) return Promise.resolve();\n    // Share a single in-flight promise so racing callers coalesce onto one\n    // attempt instead of relying on a synchronous flag set before the await.\n    if (this.subscribePromise !== null) return this.subscribePromise;\n    this.subscribePromise = this.transport\n      .subscribe({\n        jobType: this.jobType,\n        worker: this.name,\n        credits: this.credits,\n        timeoutMs: this.cfg.jobTimeoutMs ?? 60_000,\n        fetchVariables: this.cfg.fetchVariables ?? null,\n        onJob: (raw) => void this.dispatch(raw),\n      })\n      .catch((err) => {\n        // Do not wedge the worker as \"subscribed\" on failure: clear the shared\n        // promise so a later start() can retry, and rethrow so awaiting callers\n        // (including the proxy's fallback-to-REST catch) observe the error.\n        this.subscribePromise = null;\n        throw err;\n      });\n    return this.subscribePromise;\n  }\n\n  private enrich(raw: JobFrame) {\n    let acted = false;\n    const ack = (replenish = true) => {\n      if (!acted) {\n        acted = true;\n        if (replenish) this.transport?.grantCredits(this.jobType, 1);\n      }\n    };\n    return {\n      ...raw,\n      complete: async (variables: Record<string, unknown> = {}) => {\n        this.transport?.completeJob(raw.jobKey, variables);\n        ack();\n        return JobActionReceipt;\n      },\n      fail: async (reason: { errorMessage?: string; retries?: number } = {}) => {\n        this.transport?.failJob(raw.jobKey, reason.retries, reason.errorMessage);\n        ack();\n        return JobActionReceipt;\n      },\n      error: async (e: { errorCode: string; errorMessage?: string }) => {\n        this.transport?.throwError(raw.jobKey, e.errorCode, e.errorMessage);\n        ack();\n        return JobActionReceipt;\n      },\n      ignore: async () => {\n        ack();\n        return JobActionReceipt;\n      },\n    };\n  }\n\n  private async dispatch(raw: JobFrame): Promise<void> {\n    if (this.stopped) return;\n    const job = this.enrich(raw);\n    try {\n      await this.cfg.jobHandler(job);\n    } catch (err) {\n      try {\n        this.transport?.failJob(raw.jobKey, 0, String(err));\n      } finally {\n        this.transport?.grantCredits(this.jobType, 1);\n      }\n    }\n  }\n\n  stop(): void {\n    this.stopped = true;\n    this.startRequested = false;\n    this.subscribePromise = null;\n    this.transport?.unsubscribe(this.jobType);\n  }\n  close(): void {\n    this.stop();\n  }\n}\n","// Internal helpers for the pure-REST job-worker long-poll default.\n//\n// These are intentionally NOT re-exported from the package entrypoint: they are\n// implementation details of the REST path, exercised directly by the unit tests\n// (which import this module). Keeping them here keeps the public API surface\n// focused on the SDK it re-exports and avoids leaking helper types into the\n// published `.d.ts`.\nimport type { createCamundaClient as createCamundaClientBase } from \"@camunda8/orchestration-cluster-api\";\n\n/** The `createJobWorker` config accepted by the upstream client, derived from\n *  the SDK entrypoint so we keep its full job-worker config type (including\n *  `pollTimeoutMs`) without importing extra types. */\nexport type RestJobWorkerConfig = Parameters<\n  ReturnType<typeof createCamundaClientBase>[\"createJobWorker\"]\n>[0];\n\n/**\n * Default broker long-poll window (ms) applied to a REST job worker when the\n * caller doesn't set one. Without it the upstream REST worker activates jobs\n * with `pollTimeoutMs` unset, so the Nano gateway falls back to its own short\n * default window (~5s) and an idle worker reconnects every few seconds. A 30s\n * window keeps an idle REST worker on one held connection ~6x longer, cutting\n * reconnect churn (and the chances of a transient connect failure) while the\n * broker still returns immediately the moment a job arrives. Falcon workers are\n * push-based and are unaffected. An explicit `pollTimeoutMs` (including\n * `0`/negative) always wins.\n */\nexport const REST_DEFAULT_POLL_TIMEOUT_MS = 30_000;\n\n/** Inject the default REST long-poll window unless the caller set one. */\nexport function withRestPollDefault(cfg: RestJobWorkerConfig): RestJobWorkerConfig {\n  if (cfg && cfg.pollTimeoutMs === undefined) {\n    return { ...cfg, pollTimeoutMs: REST_DEFAULT_POLL_TIMEOUT_MS };\n  }\n  return cfg;\n}\n\n/** Wrap a base client so REST `createJobWorker` gets the default long-poll\n *  window, leaving every other method untouched. Used on the pure-REST path. */\nexport function wrapRestPollDefault(\n  client: ReturnType<typeof createCamundaClientBase>,\n): ReturnType<typeof createCamundaClientBase> {\n  // Cache bound methods per property key so repeated reads return the same\n  // function reference (stable identity), rather than allocating a fresh bound\n  // function on every access.\n  const boundCache = new Map<PropertyKey, unknown>();\n  return new Proxy(client, {\n    get(target, prop) {\n      if (prop === \"createJobWorker\") {\n        let wrapper = boundCache.get(prop);\n        if (wrapper === undefined) {\n          wrapper = (cfg: RestJobWorkerConfig) =>\n            (target as any).createJobWorker(withRestPollDefault(cfg));\n          boundCache.set(prop, wrapper);\n        }\n        return wrapper;\n      }\n      // Preserve original call semantics: read (and later bind) against the\n      // underlying client so `this` is the real client (not this Proxy),\n      // which matters for methods — and getter accessors — that touch private\n      // fields. Passing `target` as the receiver ensures accessors run with\n      // `this === target` rather than the Proxy.\n      const value = Reflect.get(target, prop, target);\n      if (typeof value !== \"function\") {\n        return value;\n      }\n      let bound = boundCache.get(prop);\n      if (bound === undefined) {\n        bound = value.bind(target);\n        boundCache.set(prop, bound);\n      }\n      return bound;\n    },\n  });\n}\n"],"mappings":";;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAOA,IAAAA,oCAEO;;;ACCP,IAAM,QAAQ,oBAAI,IAA6B;AAGxC,SAAS,cAAc,KAAqB;AACjD,SAAO,IAAI,QAAQ,QAAQ,EAAE;AAC/B;AAMA,eAAsB,WACpB,aACA,YAA0B,OACA;AAC1B,QAAM,OAAO,cAAc,WAAW;AACtC,MAAI,MAAM,IAAI,IAAI,EAAG,QAAO,MAAM,IAAI,IAAI,KAAK;AAC/C,MAAI,OAAwB;AAC5B,MAAI;AACF,UAAM,MAAM,MAAM,UAAU,GAAG,IAAI,gBAAgB;AAAA,MACjD,SAAS,EAAE,QAAQ,mBAAmB;AAAA,IACxC,CAAC;AACD,QAAI,IAAI,IAAI;AACV,YAAM,OAAY,MAAM,IAAI,KAAK,EAAE,MAAM,OAAO,CAAC,EAAE;AACnD,YAAM,OAAO,MAAM;AACnB,UAAI,QAAQ,OAAO,KAAK,eAAe,UAAU;AAC/C,eAAO;AAAA,UACL,QAAQ,OAAO,KAAK,UAAU,UAAU;AAAA,UACxC,SAAS,KAAK,UAAU,OAAO,KAAK,OAAO,IAAI;AAAA,UAC/C,YAAY,KAAK;AAAA,QACnB;AAAA,MACF;AAAA,IACF;AAAA,EACF,QAAQ;AACN,WAAO;AAAA,EACT;AACA,QAAM,IAAI,MAAM,IAAI;AACpB,SAAO;AACT;AAQO,SAAS,UAAU,aAAqB,MAAsB;AACnE,MAAI,OAAO,cAAc,WAAW;AACpC,MAAI,KAAK,WAAW,UAAU,EAAG,QAAO,WAAW,KAAK,MAAM,WAAW,MAAM;AAAA,WACtE,KAAK,WAAW,SAAS,EAAG,QAAO,UAAU,KAAK,MAAM,UAAU,MAAM;AACjF,SAAO,OAAO;AAChB;;;ACzDA,IAAI,SAAwB;AAE5B,eAAsB,eAAgC;AACpD,MAAI,OAAQ,QAAO;AACnB,QAAM,IAAI;AACV,MAAI,OAAO,EAAE,cAAc,YAAY;AACrC,aAAS,EAAE;AACX,WAAO;AAAA,EACT;AACA,QAAM,MAAM,MAAM,OAAO,IAAI;AAC7B,WAAU,IAAI,WAAW,IAAI;AAC7B,SAAO;AACT;;;ACAA,SAAS,MAAM,IAA2B;AACxC,SAAO,IAAI,QAAQ,CAAC,YAAY;AAC9B,UAAM,IAAI,WAAW,SAAS,EAAE;AAChC,MAAE,QAAQ;AAAA,EACZ,CAAC;AACH;AAaO,IAAM,sBAAN,cAAkC,MAAM;AAAA,EAC7C,YACW,WACA,QACA,OACT;AACA,UAAM,qBAAqB,SAAS,YAAY,MAAM,EAAE;AAJ/C;AACA;AACA;AAGT,SAAK,OAAO;AAAA,EACd;AAAA,EANW;AAAA,EACA;AAAA,EACA;AAKb;AA0CO,IAAM,yBAAN,cAAqC,MAAM;AAAA,EAChD,YAAqB,WAAmB;AACtC;AAAA,MACE,oDAAoD,SAAS;AAAA,IAC/D;AAHmB;AAInB,SAAK,OAAO;AAAA,EACd;AAAA,EALqB;AAMvB;AAgBO,IAAM,sBAAN,cAAkC,MAAM;AAAA,EAC7C,YAAqB,WAAmB;AACtC;AAAA,MACE,uDAAuD,SAAS;AAAA,IAClE;AAHmB;AAInB,SAAK,OAAO;AAAA,EACd;AAAA,EALqB;AAMvB;AAEO,IAAM,kBAAN,MAAsB;AAAA,EACnB;AAAA,EACA,KAAuB;AAAA,EACvB,OAAO;AAAA,EACP,OAAO;AAAA,EACP,UAAU,oBAAI,IAAqB;AAAA,EACnC,SAAS,oBAAI,IAAgF;AAAA,EAC7F,OAAO,oBAAI,IAA0B;AAAA,EACrC,YAAmD;AAAA,EACnD,iBAAuC;AAAA,EACvC,SAAS;AAAA;AAAA,EAET,eAAe;AAAA,EACf;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA,EAQA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA,EAMA,UAAU;AAAA,EACV,gBAAgC,CAAC;AAAA,EAEzC,YACE,aACA,MACA,wBACA,kBACA;AACA,SAAK,MAAM,UAAU,aAAa,IAAI;AACtC,SAAK,yBACH,2BAA2B,UAAa,yBAAyB,IAC7D,yBACA;AACN,SAAK,mBACH,qBAAqB,UAAa,mBAAmB,IAAI,mBAAmB;AAAA,EAChF;AAAA,EAEQ,WAAmB;AACzB,SAAK,QAAQ;AACb,WAAO,KAAK;AAAA,EACd;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA,EAQQ,UAAU,WAAmB,GAAwB;AAC3D,UAAM,IAAI,EAAE;AACZ,QAAI,OAAO,MAAM,YAAY,CAAC,OAAO,UAAU,CAAC,KAAK,KAAK,GAAG;AAE3D,cAAQ,KAAK,6CAA6C,SAAS,iCAAiC;AACpG,aAAO;AAAA,IACT;AACA,WAAO;AAAA,EACT;AAAA,EAEQ,KAAK,OAAmB;AAC9B,QAAI,KAAK,MAAM,KAAK,KAAM,MAAK,GAAG,KAAK,KAAK,UAAU,KAAK,CAAC;AAAA,EAC9D;AAAA,EAEA,MAAM,UAAyB;AAC7B,QAAI,KAAK,OAAQ,OAAM,IAAI,MAAM,kBAAkB,KAAK,GAAG,EAAE;AAC7D,QAAI,KAAK,KAAM;AACf,QAAI,KAAK,gBAAgB;AAKvB,YAAM,KAAK;AACX,UAAI,CAAC,KAAK,KAAM,OAAM,IAAI,MAAM,iDAAiD,KAAK,GAAG,EAAE;AAC3F;AAAA,IACF;AAKA,UAAM,IAAI,KAAK,WAAW;AAC1B,SAAK,iBAAiB;AACtB,QAAI;AACF,YAAM;AAAA,IACR,SAAS,GAAG;AACV,UAAI,KAAK,mBAAmB,EAAG,MAAK,iBAAiB;AACrD,YAAM;AAAA,IACR;AAAA,EACF;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA,EAaQ,aAA4B;AAClC,WAAO,IAAI,QAAc,CAAC,SAAS,WAAW;AAC5C,UAAI,UAAU;AAKd,UAAI,iBAAiB;AAKrB,UAAI;AACJ,YAAM,oBAAoB,MAAM;AAC9B,YAAI,iBAAiB,QAAW;AAC9B,uBAAa,YAAY;AACzB,yBAAe;AAAA,QACjB;AAAA,MACF;AACA,YAAM,gBAAgB,MAAM;AAC1B,yBAAiB;AACjB,0BAAkB;AAClB,YAAI,CAAC,SAAS;AACZ,oBAAU;AACV,kBAAQ;AAAA,QACV;AAAA,MACF;AACA,YAAM,eAAe,CAAC,MAAe;AACnC,0BAAkB;AAClB,YAAI,CAAC,SAAS;AACZ,oBAAU;AACV,iBAAO,CAAC;AAAA,QACV;AAAA,MACF;AACA,WAAK,aAAa,EACf,KAAK,CAAC,OAAO;AACZ,cAAM,KAAK,IAAI,GAAG,KAAK,GAAG;AAC1B,aAAK,KAAK;AAKV,YAAI,KAAK,qBAAqB,QAAW;AACvC,yBAAe,WAAW,MAAM;AAC9B,gBAAI,KAAK,OAAO,MAAM,QAAS;AAC/B,yBAAa,IAAI,oBAAoB,KAAK,gBAAiB,CAAC;AAM5D,iBAAK,KAAK;AACV,gBAAI;AACF,iBAAG,MAAM,MAAM,mBAAmB;AAAA,YACpC,QAAQ;AAAA,YAER;AAAA,UACF,GAAG,KAAK,gBAAgB;AACxB,UAAC,aAAwC,QAAQ;AAAA,QACnD;AACA,WAAG,SAAS,MAAM;AAAA,QAAC;AACnB,WAAG,YAAY,CAAC,OAAqB;AAInC,cAAI,KAAK,OAAO,GAAI;AACpB,cAAI;AACJ,cAAI;AACF,gBAAI,KAAK,MAAM,OAAO,GAAG,SAAS,WAAW,GAAG,OAAO,OAAO,GAAG,IAAI,CAAC;AAAA,UACxE,QAAQ;AACN;AAAA,UACF;AACA,eAAK,OAAO,GAAG,aAAa;AAAA,QAC9B;AACA,WAAG,UAAU,MAAM;AACjB,cAAI,KAAK,OAAO,GAAI;AAGpB,cAAI,CAAC,KAAK,KAAM,cAAa,IAAI,MAAM,0BAA0B,KAAK,GAAG,EAAE,CAAC;AAAA,QAC9E;AACA,WAAG,UAAU,MAAM;AAKjB,cAAI,KAAK,OAAO,IAAI;AAClB,yBAAa,IAAI,MAAM,+BAA+B,KAAK,GAAG,EAAE,CAAC;AACjE;AAAA,UACF;AACA,eAAK,OAAO;AAIZ,eAAK,UAAU;AACf,gBAAM,UAAU,KAAK;AACrB,eAAK,gBAAgB,CAAC;AACtB,qBAAW,KAAK,QAAS,GAAE,KAAK,IAAI,MAAM,eAAe,CAAC;AAC1D,cAAI,KAAK,UAAW,eAAc,KAAK,SAAS;AAEhD,qBAAW,KAAK,KAAK,QAAQ,OAAO,EAAG,GAAE,OAAO,IAAI,MAAM,eAAe,CAAC;AAC1E,eAAK,QAAQ,MAAM;AAGnB,qBAAW,KAAK,KAAK,OAAO,OAAO,EAAG,GAAE,OAAO,IAAI,MAAM,eAAe,CAAC;AACzE,eAAK,OAAO,MAAM;AAGlB,uBAAa,IAAI,MAAM,+BAA+B,KAAK,GAAG,EAAE,CAAC;AAKjE,cAAI,CAAC,KAAK,UAAU,eAAgB,MAAK,eAAe;AAAA,QAC1D;AAAA,MACF,CAAC,EACA,MAAM,CAAC,MAAM,aAAa,CAAC,CAAC;AAAA,IACjC,CAAC;AAAA,EACH;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA,EAQQ,iBAAuB;AAC7B,QAAI,KAAK,gBAAgB,KAAK,OAAQ;AACtC,SAAK,eAAe;AACpB,QAAI;AACJ,SAAK,iBAAiB,IAAI,QAAc,CAAC,QAAQ;AAC/C,qBAAe;AAAA,IACjB,CAAC;AACD,SAAK,KAAK,cAAc,YAAY;AAAA,EACtC;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA,EASA,MAAc,cAAc,cAAyC;AACnE,QAAI,UAAU;AACd,WAAO,CAAC,KAAK,UAAU,CAAC,KAAK,MAAM;AACjC,YAAM,MAAM,KAAK,eAAe,OAAO,CAAC;AACxC,iBAAW;AACX,UAAI,KAAK,UAAU,KAAK,KAAM;AAC9B,UAAI;AACF,cAAM,KAAK,WAAW;AAAA,MACxB,QAAQ;AAEN;AAAA,MACF;AAAA,IACF;AACA,SAAK,eAAe;AACpB,QAAI,CAAC,KAAK,KAAM,MAAK,iBAAiB;AACtC,iBAAa;AAAA,EACf;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA,EASQ,eAAe,SAAyB;AAC9C,UAAM,OAAO;AACb,UAAM,MAAM;AACZ,UAAM,UAAU,KAAK,IAAI,KAAK,OAAO,KAAK,OAAO;AACjD,WAAO,KAAK,MAAM,KAAK,OAAO,IAAI,OAAO;AAAA,EAC3C;AAAA,EAEQ,OAAO,GAAS,gBAAkC;AACxD,YAAQ,EAAE,MAAM;AAAA,MACd,KAAK,WAAW;AACd,aAAK,OAAO;AACZ,aAAK,UAAU,OAAO,EAAE,qBAAqB,CAAC;AAC9C,cAAM,KAAK,OAAO,EAAE,eAAe,CAAC;AACpC,YAAI,KAAK,GAAG;AACV,cAAI,KAAK,UAAW,eAAc,KAAK,SAAS;AAChD,eAAK,YAAY,YAAY,MAAM,KAAK,KAAK,EAAE,MAAM,YAAY,CAAC,GAAG,EAAE;AAAA,QACzE;AACA,mBAAW,OAAO,KAAK,KAAK,OAAO,EAAG,MAAK,cAAc,GAAG;AAC5D,aAAK,qBAAqB;AAC1B,uBAAe;AACf;AAAA,MACF;AAAA,MACA,KAAK,OAAO;AACV,cAAM,MAAO,EAAE,OAAqB,CAAC;AACrC,cAAM,MAAM,KAAK,KAAK,IAAI,IAAI,IAAI;AAClC,YAAI,IAAK,KAAI,MAAM,GAAG;AACtB;AAAA,MACF;AAAA,MACA,KAAK,iBAAiB;AACpB,cAAM,OAAO,KAAK,UAAU,iBAAiB,CAAC;AAC9C,YAAI,SAAS,KAAM;AACnB,cAAM,IAAI,KAAK,QAAQ,IAAI,IAAI;AAC/B,YAAI,CAAC,EAAG;AACR,aAAK,QAAQ,OAAO,IAAI;AAIxB,YAAI,OAAO,EAAE,WAAW,YAAY,CAAC,OAAO,SAAS,EAAE,MAAM,GAAG;AAC9D,YAAE,OAAO,IAAI,oBAAoB,iBAAiB,4BAA4B,CAAC,CAAC;AAChF;AAAA,QACF;AAGA,UAAE,QAAQ,EAAE,QAAQ,EAAE,QAAQ,MAAM,EAAE,KAAK,CAAC;AAC5C;AAAA,MACF;AAAA,MACA,KAAK,qBAAqB;AACxB,cAAM,OAAO,KAAK,UAAU,qBAAqB,CAAC;AAClD,YAAI,SAAS,KAAM;AACnB,cAAM,KAAK,KAAK,OAAO,IAAI,IAAI;AAC/B,YAAI,CAAC,GAAI;AACT,aAAK,OAAO,OAAO,IAAI;AAIvB,cAAM,SAAS,EAAE;AACjB,cAAM,MAAM,OAAO,WAAW,WAAW,OAAO,MAAM,IAAI;AAC1D,YAAI,OAAO,QAAQ,YAAY,IAAI,WAAW,GAAG;AAC/C,aAAG,OAAO,IAAI,oBAAoB,qBAAqB,gCAAgC,CAAC,CAAC;AACzF;AAAA,QACF;AACA,WAAG,QAAQ;AAAA,UACT,kBAAkB,QAAQ,EAAE,gBAAgB;AAAA,UAC5C,WAAW,EAAE,aAAa,CAAC;AAAA,UAC3B,oBAAoB;AAAA,QACtB,CAAC;AACD;AAAA,MACF;AAAA,MACA,KAAK,qBAAqB;AAExB,aAAK,WAAW,OAAO,EAAE,KAAK,CAAC;AAC/B,aAAK,qBAAqB;AAC1B;AAAA,MACF;AAAA,IAEF;AAAA,EACF;AAAA,EAEQ,cAAc,KAAyB;AAC7C,SAAK,KAAK;AAAA,MACR,MAAM;AAAA,MACN,SAAS,IAAI;AAAA,MACb,YAAY,IAAI;AAAA,MAChB,QAAQ,IAAI;AAAA,MACZ,SAAS,IAAI;AAAA,MACb,eAAe,IAAI;AAAA,IACrB,CAAC;AAAA,EACH;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA,EAUQ,cAAc,WAAmC;AACvD,QAAI,KAAK,UAAU,GAAG;AACpB,WAAK,WAAW;AAChB,aAAO,QAAQ,QAAQ;AAAA,IACzB;AACA,WAAO,IAAI,QAAc,CAAC,SAAS,WAAW;AAC5C,UAAI;AACJ,YAAM,SAAuB;AAAA,QAC3B,OAAO,MAAM;AACX,cAAI,UAAU,OAAW,cAAa,KAAK;AAC3C,eAAK,WAAW;AAChB,kBAAQ;AAAA,QACV;AAAA,QACA,MAAM,CAAC,MAAM;AACX,cAAI,UAAU,OAAW,cAAa,KAAK;AAC3C,iBAAO,CAAC;AAAA,QACV;AAAA,MACF;AACA,WAAK,cAAc,KAAK,MAAM;AAC9B,UAAI,cAAc,UAAa,YAAY,GAAG;AAC5C,gBAAQ,WAAW,MAAM;AACvB,gBAAM,IAAI,KAAK,cAAc,QAAQ,MAAM;AAC3C,cAAI,KAAK,EAAG,MAAK,cAAc,OAAO,GAAG,CAAC;AAC1C,iBAAO,IAAI,uBAAuB,SAAS,CAAC;AAAA,QAC9C,GAAG,SAAS;AACZ,cAAM,QAAQ;AAAA,MAChB;AAAA,IACF,CAAC;AAAA,EACH;AAAA,EAEQ,uBAA6B;AACnC,WAAO,KAAK,UAAU,KAAK,KAAK,cAAc,SAAS,GAAG;AACxD,WAAK,cAAc,MAAM,EAAG,MAAM;AAAA,IACpC;AAAA,EACF;AAAA;AAAA,EAGA,MAAM,eAAe,OAcsH;AACzI,UAAM,KAAK,QAAQ;AAInB,UAAM,KAAK,cAAc,MAAM,mBAAmB,KAAK,sBAAsB;AAC7E,UAAM,OAAO,KAAK,SAAS;AAC3B,QAAI,oBAAsD;AAC1D,QAAI,mBAAkD;AACtD,UAAM,aAAa,MAAM,kBACrB,IAAI,QAAoB,CAAC,SAAS,WAAW;AAC3C,0BAAoB;AACpB,yBAAmB;AAAA,IACrB,CAAC,IACD;AACJ,QAAI,qBAAqB,kBAAkB;AACzC,WAAK,OAAO,IAAI,MAAM,EAAE,SAAS,mBAAmB,QAAQ,iBAAiB,CAAC;AAK9E,kBAAY,MAAM,MAAM;AAAA,MAAC,CAAC;AAAA,IAC5B;AACA,QAAI;AACJ,QAAI;AACF,eAAS,MAAM,IAAI,QAA2C,CAAC,SAAS,WAAW;AACjF,aAAK,QAAQ,IAAI,MAAM,EAAE,SAAS,OAAO,CAAC;AAC1C,aAAK,KAAK;AAAA,UACR,MAAM;AAAA,UACN;AAAA,UACA,qBAAqB,MAAM,uBAAuB;AAAA,UAClD,sBAAsB,MAAM,wBAAwB;AAAA,UACpD,WAAW,MAAM,aAAa;AAAA,UAC9B,iBAAiB,MAAM,mBAAmB;AAAA,UAC1C,gBAAgB,MAAM,kBAAkB;AAAA,UACxC,gBAAgB,MAAM,oBAAoB;AAAA,QAC5C,CAAC;AAAA,MACH,CAAC;AAAA,IACH,SAAS,KAAK;AAIZ,YAAM,SAAS,KAAK,OAAO,IAAI,IAAI;AACnC,UAAI,QAAQ;AACV,aAAK,OAAO,OAAO,IAAI;AACvB,eAAO,OAAO,GAAG;AAAA,MACnB;AACA,YAAM;AAAA,IACR;AACA,UAAM,OAAO,aAAa,MAAM,aAAa;AAC7C,WAAO,EAAE,GAAG,QAAQ,YAAY,KAAK;AAAA,EACvC;AAAA;AAAA,EAGA,MAAM,UAAU,KAAkC;AAMhD,UAAM,UAAU,KAAK;AACrB,SAAK,KAAK,IAAI,IAAI,SAAS,GAAG;AAC9B,UAAM,KAAK,QAAQ;AAMnB,QAAI,KAAK,KAAK,IAAI,IAAI,OAAO,MAAM,IAAK;AACxC,QAAI,QAAS,MAAK,cAAc,GAAG;AAAA,EACrC;AAAA,EAEA,YAAY,SAAuB;AACjC,SAAK,KAAK,OAAO,OAAO;AAAA,EAC1B;AAAA,EAEA,YAAY,QAAgB,WAA2C;AACrE,SAAK,KAAK,EAAE,MAAM,eAAe,MAAM,KAAK,SAAS,GAAG,QAAQ,WAAW,aAAa,KAAK,CAAC;AAAA,EAChG;AAAA,EACA,QAAQ,QAAgB,SAAkB,cAA6B;AACrE,SAAK,KAAK,EAAE,MAAM,WAAW,MAAM,KAAK,SAAS,GAAG,QAAQ,SAAS,WAAW,GAAG,cAAc,gBAAgB,KAAK,CAAC;AAAA,EACzH;AAAA,EACA,WAAW,QAAgB,WAAmB,cAA6B;AACzE,SAAK,KAAK,EAAE,MAAM,cAAc,MAAM,KAAK,SAAS,GAAG,QAAQ,WAAW,cAAc,gBAAgB,KAAK,CAAC;AAAA,EAChH;AAAA,EACA,aAAa,SAAiB,GAAiB;AAC7C,SAAK,KAAK,EAAE,MAAM,cAAc,SAAS,EAAE,CAAC;AAAA,EAC9C;AAAA,EAEA,QAAc;AACZ,SAAK,SAAS;AACd,QAAI,KAAK,UAAW,eAAc,KAAK,SAAS;AAChD,QAAI;AACF,WAAK,IAAI,MAAM,KAAM,UAAU;AAAA,IACjC,QAAQ;AAAA,IAER;AAAA,EACF;AACF;;;ACrlBO,IAAM,oBAAN,MAAwB;AAAA,EAI7B,YAAoB,MAA4B,SAAS,IAAI;AAAzC;AAA4B;AAAA,EAAc;AAAA,EAA1C;AAAA,EAA4B;AAAA,EAHxC,OAAO,oBAAI,IAAiB;AAAA,EAC5B,QAA+C;AAAA,EAC/C,SAAS;AAAA,EAGjB,MAAM,eAAe,OAIsH;AACzI,UAAM,EAAE,mBAAmB,IAAI,MAAM,KAAK,KAAK,eAAe,KAAK;AACnE,SAAK,KAAK,KAAK;AACf,QAAI,CAAC,MAAM,gBAAiB,QAAO,EAAE,QAAQ,KAAK,MAAM,EAAE,mBAAmB,EAAE;AAC/E,WAAO,CAAC,KAAK,KAAK,kBAAkB,kBAAkB,GAAG;AACvD,WAAK,KAAK,KAAK;AACf,YAAM,IAAI,QAAQ,CAAC,MAAM,WAAW,GAAG,KAAK,MAAM,CAAC;AAAA,IACrD;AACA,UAAM,YAAY,KAAK,KAAK,kBAAkB,kBAAkB;AAChE,WAAO,EAAE,QAAQ,KAAK,MAAM,EAAE,oBAAoB,UAAU,GAAG,YAAY,EAAE,kBAAkB,MAAM,WAAW,mBAAmB,EAAE;AAAA,EACvI;AAAA,EAEA,MAAM,UAAU,KAAsH;AACpI,SAAK,KAAK,IAAI,IAAI,SAAS,EAAE,QAAQ,IAAI,QAAQ,SAAS,IAAI,SAAS,WAAW,IAAI,WAAW,OAAO,IAAI,MAAM,CAAC;AACnH,QAAI,CAAC,KAAK,MAAO,MAAK,QAAQ,YAAY,MAAM,KAAK,KAAK,GAAG,KAAK,MAAM;AAAA,EAC1E;AAAA,EACA,YAAY,SAAuB;AAAE,SAAK,KAAK,OAAO,OAAO;AAAA,EAAG;AAAA,EAEhE,YAAY,QAAgB,WAA2C;AAAE,SAAK,KAAK,KAAK,YAAY,QAAQ,SAAS;AAAA,EAAG;AAAA,EACxH,QAAQ,QAAgB,SAAkB,cAA6B;AAAE,SAAK,KAAK,KAAK,QAAQ,QAAQ,WAAW,GAAG,YAAY;AAAA,EAAG;AAAA,EACrI,WAAW,QAAgB,YAAoB,cAA6B;AAAE,SAAK,KAAK,KAAK,QAAQ,QAAQ,GAAG,YAAY;AAAA,EAAG;AAAA,EAC/H,aAAa,SAAiB,GAAiB;AAAE,UAAM,IAAI,KAAK,KAAK,IAAI,OAAO;AAAG,QAAI,EAAG,GAAE,WAAW;AAAA,EAAG;AAAA,EAElG,OAAa;AACnB,QAAI,KAAK,OAAQ;AACjB,SAAK,KAAK,KAAK;AACf,eAAW,CAAC,MAAM,CAAC,KAAK,KAAK,MAAM;AACjC,UAAI,EAAE,WAAW,EAAG;AACpB,WAAK,KAAK,KAAK,aAAa,MAAM,EAAE,SAAS,EAAE,WAAW,EAAE,MAAM,EAAE,KAAK,CAAC,SAAS;AACjF,mBAAW,KAAK,MAAM;AAAE,cAAI,EAAE,WAAW,EAAG;AAAO,YAAE;AAAW,YAAE,MAAM,CAAC;AAAA,QAAG;AAAA,MAC9E,CAAC;AAAA,IACH;AAAA,EACF;AAAA,EACA,QAAc;AAAE,SAAK,SAAS;AAAM,QAAI,KAAK,MAAO,eAAc,KAAK,KAAK;AAAA,EAAG;AACjF;;;AChFA,uCAA2D;AAapD,IAAM,gBAAN,MAAoB;AAAA,EACjB;AAAA,EACA;AAAA,EACA;AAAA,EACA,UAAU;AAAA;AAAA,EAEV,iBAAiB;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA,EAOjB,mBAAyC;AAAA,EACxC;AAAA,EACA;AAAA,EAET,YAAY,WAAmC,KAA0B;AACvE,SAAK,YAAY;AACjB,SAAK,MAAM;AACX,SAAK,UAAU,IAAI;AACnB,SAAK,UAAU,IAAI,mBAAmB;AACtC,SAAK,OAAO,IAAI,cAAc,UAAU,IAAI,OAAO;AACnD,QAAI,IAAI,cAAc,MAAO,MAAK,KAAK,MAAM;AAAA,EAC/C;AAAA;AAAA;AAAA;AAAA;AAAA,EAMA,cAAc,WAAkC;AAC9C,SAAK,YAAY;AAIjB,QAAI,KAAK,kBAAkB,CAAC,KAAK,QAAS,MAAK,KAAK,UAAU,EAAE,MAAM,MAAM;AAAA,IAAC,CAAC;AAAA,EAChF;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA,EASA,MAAM,QAAuB;AAC3B,SAAK,UAAU;AACf,SAAK,iBAAiB;AACtB,QAAI,KAAK,cAAc,KAAM;AAC7B,UAAM,KAAK,UAAU;AAAA,EACvB;AAAA,EAEQ,YAA2B;AACjC,QAAI,KAAK,WAAW,KAAK,cAAc,KAAM,QAAO,QAAQ,QAAQ;AAGpE,QAAI,KAAK,qBAAqB,KAAM,QAAO,KAAK;AAChD,SAAK,mBAAmB,KAAK,UAC1B,UAAU;AAAA,MACT,SAAS,KAAK;AAAA,MACd,QAAQ,KAAK;AAAA,MACb,SAAS,KAAK;AAAA,MACd,WAAW,KAAK,IAAI,gBAAgB;AAAA,MACpC,gBAAgB,KAAK,IAAI,kBAAkB;AAAA,MAC3C,OAAO,CAAC,QAAQ,KAAK,KAAK,SAAS,GAAG;AAAA,IACxC,CAAC,EACA,MAAM,CAAC,QAAQ;AAId,WAAK,mBAAmB;AACxB,YAAM;AAAA,IACR,CAAC;AACH,WAAO,KAAK;AAAA,EACd;AAAA,EAEQ,OAAO,KAAe;AAC5B,QAAI,QAAQ;AACZ,UAAM,MAAM,CAAC,YAAY,SAAS;AAChC,UAAI,CAAC,OAAO;AACV,gBAAQ;AACR,YAAI,UAAW,MAAK,WAAW,aAAa,KAAK,SAAS,CAAC;AAAA,MAC7D;AAAA,IACF;AACA,WAAO;AAAA,MACL,GAAG;AAAA,MACH,UAAU,OAAO,YAAqC,CAAC,MAAM;AAC3D,aAAK,WAAW,YAAY,IAAI,QAAQ,SAAS;AACjD,YAAI;AACJ,eAAO,iCAAAC;AAAA,MACT;AAAA,MACA,MAAM,OAAO,SAAsD,CAAC,MAAM;AACxE,aAAK,WAAW,QAAQ,IAAI,QAAQ,OAAO,SAAS,OAAO,YAAY;AACvE,YAAI;AACJ,eAAO,iCAAAA;AAAA,MACT;AAAA,MACA,OAAO,OAAO,MAAoD;AAChE,aAAK,WAAW,WAAW,IAAI,QAAQ,EAAE,WAAW,EAAE,YAAY;AAClE,YAAI;AACJ,eAAO,iCAAAA;AAAA,MACT;AAAA,MACA,QAAQ,YAAY;AAClB,YAAI;AACJ,eAAO,iCAAAA;AAAA,MACT;AAAA,IACF;AAAA,EACF;AAAA,EAEA,MAAc,SAAS,KAA8B;AACnD,QAAI,KAAK,QAAS;AAClB,UAAM,MAAM,KAAK,OAAO,GAAG;AAC3B,QAAI;AACF,YAAM,KAAK,IAAI,WAAW,GAAG;AAAA,IAC/B,SAAS,KAAK;AACZ,UAAI;AACF,aAAK,WAAW,QAAQ,IAAI,QAAQ,GAAG,OAAO,GAAG,CAAC;AAAA,MACpD,UAAE;AACA,aAAK,WAAW,aAAa,KAAK,SAAS,CAAC;AAAA,MAC9C;AAAA,IACF;AAAA,EACF;AAAA,EAEA,OAAa;AACX,SAAK,UAAU;AACf,SAAK,iBAAiB;AACtB,SAAK,mBAAmB;AACxB,SAAK,WAAW,YAAY,KAAK,OAAO;AAAA,EAC1C;AAAA,EACA,QAAc;AACZ,SAAK,KAAK;AAAA,EACZ;AACF;;;ACzHO,IAAM,+BAA+B;AAGrC,SAAS,oBAAoB,KAA+C;AACjF,MAAI,OAAO,IAAI,kBAAkB,QAAW;AAC1C,WAAO,EAAE,GAAG,KAAK,eAAe,6BAA6B;AAAA,EAC/D;AACA,SAAO;AACT;AAIO,SAAS,oBACd,QAC4C;AAI5C,QAAM,aAAa,oBAAI,IAA0B;AACjD,SAAO,IAAI,MAAM,QAAQ;AAAA,IACvB,IAAI,QAAQ,MAAM;AAChB,UAAI,SAAS,mBAAmB;AAC9B,YAAI,UAAU,WAAW,IAAI,IAAI;AACjC,YAAI,YAAY,QAAW;AACzB,oBAAU,CAAC,QACR,OAAe,gBAAgB,oBAAoB,GAAG,CAAC;AAC1D,qBAAW,IAAI,MAAM,OAAO;AAAA,QAC9B;AACA,eAAO;AAAA,MACT;AAMA,YAAM,QAAQ,QAAQ,IAAI,QAAQ,MAAM,MAAM;AAC9C,UAAI,OAAO,UAAU,YAAY;AAC/B,eAAO;AAAA,MACT;AACA,UAAI,QAAQ,WAAW,IAAI,IAAI;AAC/B,UAAI,UAAU,QAAW;AACvB,gBAAQ,MAAM,KAAK,MAAM;AACzB,mBAAW,IAAI,MAAM,KAAK;AAAA,MAC5B;AACA,aAAO;AAAA,IACT;AAAA,EACF,CAAC;AACH;;;ANxDA,0BAAc,gDAlBd;AAmDA,SAAS,YAAY,GAAqB;AACxC,MAAI,MAAM,UAAa,MAAM,KAAM,QAAO;AAC1C,QAAM,IAAI,OAAO,CAAC,EAAE,KAAK,EAAE,YAAY;AACvC,SAAO,MAAM,MAAM,MAAM,OAAO,MAAM,SAAS,MAAM,WAAW,MAAM;AACxE;AAEA,SAAS,UAAU,MAAyB;AAC1C,QAAM,WAAY,MAAM,QAAgB;AACxC,QAAM,UACJ,OAAO,YAAY,cAAc,QAAQ,KAAK,qBAAqB;AACrE,SAAO,YAAY,QAAQ,KAAK,YAAY,OAAO;AACrD;AAEA,SAAS,YAAY,MAA+B;AAGlD,MAAI,UAAU,IAAI,EAAG,QAAO;AAC5B,QAAM,WAAW,MAAM,QAAQ;AAC/B,QAAM,UAAW,OAAO,YAAY,cAAc,QAAQ,KAAK,oBAAoB;AACnF,SAAO,YAAY,WAAW;AAChC;AAQA,SAAS,uBAAuB,MAAoC;AAClE,QAAM,MACH,MAAM,QAAQ,mCACd,OAAO,YAAY,cAAc,QAAQ,KAAK,iCAAiC;AAClF,QAAM,IAAI,OAAO,QAAQ,WAAW,OAAO,GAAG,IAAI;AAClD,SAAO,OAAO,MAAM,YAAY,OAAO,SAAS,CAAC,KAAK,IAAI,IAAI,IAAI;AACpE;AAQO,IAAM,6BAA6B;AAO1C,SAAS,wBAAwB,MAAoC;AACnE,QAAM,MACH,MAAM,QAAQ,oCACd,OAAO,YAAY,cAAc,QAAQ,KAAK,kCAAkC;AACnF,MAAI,QAAQ,UAAa,QAAQ,QAAQ,QAAQ,GAAI,QAAO;AAC5D,QAAM,IAAI,OAAO,QAAQ,WAAW,OAAO,GAAG,IAAI;AAClD,MAAI,CAAC,OAAO,SAAS,CAAC,KAAK,KAAK,EAAG,QAAO;AAC1C,SAAO;AACT;AAGA,SAAS,SAAS,aAA6B;AAC7C,SAAO,cAAc,WAAW,EAAE,QAAQ,SAAS,EAAE;AACvD;AAQO,SAAS,oBAAoB,MAA4D;AAC9F,QAAM,aAAS,kCAAAC,qBAAwB,IAAW;AAClD,QAAM,OAAO,YAAY,IAAI;AAC7B,MAAI,SAAS,OAAQ,QAAO,oBAAoB,MAAM;AAItD,MAAI,SAAS,YAAY;AACvB,QAAI,CAAC,MAAM,aAAc,OAAM,IAAI,MAAM,iDAAiD;AAC1F,UAAMC,aAAY,IAAI,kBAAkB,KAAK,YAAY;AACzD,WAAO,WAAW,QAAQ,YAAYA,YAAkB,MAAMA,WAAU,MAAM,CAAC;AAAA,EACjF;AAEA,QAAM,cAAc,OAAO,UAAU,EAAE;AACvC,QAAM,OAAO,SAAS,WAAW;AACjC,QAAM,kBAAkB,uBAAuB,IAAI;AACnD,QAAM,mBAAmB,wBAAwB,IAAI;AAErD,MAAI,YAAoC;AACxC,MAAI;AAIJ,MAAI,aAAa;AACjB,MAAI,aAAa;AACjB,QAAM,WAAW,CAAC,QAAiB;AACjC,iBAAa;AACb,gBAAY;AACZ,QAAI,CAAC,YAAY;AACf,mBAAa;AAEb,cAAQ;AAAA,QACN,gDAAiD,KAAe,WAAW,GAAG;AAAA,MAEhF;AAAA,IACF;AAAA,EACF;AAEA,QAAM,SAAS,YAA6C;AAC1D,QAAI,WAAY,QAAO;AACvB,QAAI,SAAS,QAAW;AACtB,aAAO,SAAS,WACZ,EAAE,QAAQ,YAAY,YAAY,UAAU,IAC5C,MAAM,WAAW,IAAI;AAAA,IAC3B;AACA,QAAI,CAAC,KAAM,QAAO;AAClB,QAAI,CAAC,WAAW;AACd,YAAM,IAAI,IAAI,gBAAgB,MAAM,KAAK,YAAY,iBAAiB,gBAAgB;AACtF,UAAI;AAGF,cAAM,EAAE,QAAQ;AAChB,oBAAY;AAAA,MACd,SAAS,GAAG;AACV,YAAI;AAAE,YAAE,MAAM;AAAA,QAAG,QAAQ;AAAA,QAAe;AACxC,iBAAS,CAAC;AACV,eAAO;AAAA,MACT;AAAA,IACF;AACA,WAAO;AAAA,EACT;AACA,SAAO,WAAW,QAAQ,QAAQ,MAAM,WAAW,MAAM,CAAC;AAC5D;AAIA,SAAS,WACP,QACA,QACA,QAC4C;AAC5C,SAAO,IAAI,MAAM,QAAQ;AAAA,IACvB,IAAI,QAAQ,MAAM,UAAU;AAC1B,UAAI,SAAS,yBAAyB;AACpC,eAAO,CAAC,OAAY,YAAkB;AACpC,iBAAO,OAAO,EAAE,KAAK,OAAO,MAAM;AAChC,gBAAI,CAAC,EAAG,QAAQ,OAAe,sBAAsB,OAAO,OAAO;AACnE,gBAAI;AACF,oBAAM,IAAI,MAAO,EAAU,eAAe;AAAA,gBACxC,qBAAqB,OAAO;AAAA,gBAC5B,sBAAsB,OAAO;AAAA,gBAC7B,WAAW,OAAO;AAAA,gBAClB,iBAAiB,OAAO,mBAAmB;AAAA,gBAC3C,gBAAgB,OAAO;AAAA,gBACvB,iBAAiB,OAAO;AAAA,cAC1B,CAAC;AACD,kBAAI,EAAE,UAAU,IAAK,OAAM,IAAI,MAAM,0BAA0B,EAAE,MAAM,IAAI,KAAK,UAAU,EAAE,IAAI,CAAC,EAAE;AACnG,kBAAI,EAAE,YAAY;AAGhB,sBAAM,OAAO,EAAE,QAAQ,OAAO,EAAE,SAAS,WAAW,EAAE,OAAO,CAAC;AAC9D,uBAAO,EAAE,GAAG,MAAM,WAAW,EAAE,WAAW,WAAW,oBAAoB,EAAE,WAAW,mBAAmB;AAAA,cAC3G;AAKA,kBAAI,CAAC,EAAE,QAAQ,OAAO,EAAE,SAAS,UAAU;AACzC,sBAAM,IAAI,MAAM,6CAA6C,EAAE,MAAM,uBAAuB;AAAA,cAC9F;AACA,qBAAO,EAAE;AAAA,YACX,SAAS,GAAG;AAGV,oBAAM,MAAM,OAAQ,GAAa,WAAW,CAAC;AAC7C,kBAAI,IAAI,SAAS,eAAe,KAAK,IAAI,SAAS,uBAAuB,GAAG;AAE1E,wBAAQ,KAAK,sCAAsC,GAAG,uBAAuB;AAC7E,uBAAQ,OAAe,sBAAsB,OAAO,OAAO;AAAA,cAC7D;AACA,oBAAM;AAAA,YACR;AAAA,UACF,CAAC;AAAA,QACH;AAAA,MACF;AACA,UAAI,SAAS,mBAAmB;AAC9B,eAAO,CAAC,QAAa;AACnB,gBAAM,IAAI,IAAI,cAAc,MAAM,EAAE,GAAG,KAAK,WAAW,MAAM,CAAC;AAC9D,eAAK,OAAO,EAAE,KAAK,OAAO,MAAM;AAC9B,gBAAI,CAAC,GAAG;AAAE,cAAC,OAAe,gBAAgB,oBAAoB,GAAG,CAAC;AAAG;AAAA,YAAQ;AAC7E,cAAE,cAAc,CAAQ;AACxB,gBAAI;AAAE,oBAAM,EAAE,MAAM;AAAA,YAAG,SAChB,GAAG;AAER,sBAAQ;AAAA,gBACN,2CAA4C,GAAa,WAAW,CAAC;AAAA,cAEvE;AACA,cAAC,OAAe,gBAAgB,oBAAoB,GAAG,CAAC;AAAA,YAC1D;AAAA,UACF,CAAC;AACD,iBAAO;AAAA,QACT;AAAA,MACF;AACA,UAAI,SAAS,kBAAkB;AAC7B,eAAO,MAAM;AAAE,mBAAS;AAAG,iBAAQ,OAAe,eAAe;AAAA,QAAG;AAAA,MACtE;AACA,aAAO,QAAQ,IAAI,QAAQ,MAAM,QAAQ;AAAA,IAC3C;AAAA,EACF,CAAC;AACH;AAEA,IAAO,gBAAQ;","names":["import_orchestration_cluster_api","JobActionReceipt","createCamundaClientBase","transport"]}