{"version":3,"file":"index.cjs","names":["DefaultGeneratedFile","DefaultGeneratedFileWithType","DefaultStepResult","PubSub","#maxRemoteClientQueuedBytes","#isBroker","#brokerClients","#ensureStarted","#deliverLocal","#publishFromBroker","#clientSocket","#sendToBroker","#callbacks","#sendSubscribeToBroker","#pendingWrites","#closed","#rejectSubscribeWaiters","#removeBrokerClient","#server","#starting","#start","#throwIfClosed","#listen","#connectClient","#electBroker","#handleBrokerClient","#handleServerFrame","#handleClientDisconnect","#resubscribeClient","#recoverClientConnection","#recovering","#recoverClientConnectionLoop","#isElectionLockStale","#subscribeWaiters","#removeSubscribeWaiter","#settleSubscribeWaiters","#enqueueBrokerClientWrite","#invokeLocalCallback","#sendToActiveBroker","#handlePromotedBrokerFrame"],"sources":["../../src/events/codec/error.ts","../../src/events/codec/registry.ts","../../src/events/codec/registrations.ts","../../src/events/codec/tags.ts","../../src/events/codec/codec.ts","../../src/events/unix-socket-pubsub.ts"],"sourcesContent":["import type { SerializedError } from '../../error/utils';\n\nconst MAX_CAUSE_DEPTH = 5;\n\n/**\n * Serializes an Error instance to a plain SerializedError object without\n * mutating the original. Unlike `getErrorFromUnknown(...).toJSON()`, this does\n * not attach a non-enumerable `toJSON` to the live Error.\n */\nexport function serializeError(err: Error, depth = 0): SerializedError {\n  const json: SerializedError = {\n    name: err.name || 'Error',\n    message: err.message,\n  };\n  if (err.stack !== undefined) json.stack = err.stack;\n\n  if (err.cause !== undefined) {\n    if (err.cause instanceof Error && depth < MAX_CAUSE_DEPTH) {\n      json.cause = serializeError(err.cause, depth + 1);\n    } else {\n      json.cause = err.cause;\n    }\n  }\n\n  // Copy enumerable own properties (custom error fields)\n  for (const key in err) {\n    if (!Object.prototype.hasOwnProperty.call(err, key)) continue;\n    if (key === 'message' || key === 'name' || key === 'stack' || key === 'cause') continue;\n    json[key] = (err as unknown as Record<string, unknown>)[key];\n  }\n\n  return json;\n}\n\n/**\n * Rehydrates a SerializedError into a vanilla Error instance. We never\n * instantiate user-controlled prototypes — name is preserved as a string field.\n */\nexport function rehydrateError(s: SerializedError): Error {\n  const cause =\n    s.cause !== undefined\n      ? s.cause && typeof s.cause === 'object' && 'message' in (s.cause as object) && 'name' in (s.cause as object)\n        ? rehydrateError(s.cause as SerializedError)\n        : s.cause\n      : undefined;\n\n  const err = cause !== undefined ? new Error(s.message, { cause }) : new Error(s.message);\n  if (s.name) err.name = s.name;\n  if (s.stack !== undefined) err.stack = s.stack;\n\n  for (const key in s) {\n    if (!Object.prototype.hasOwnProperty.call(s, key)) continue;\n    if (key === 'message' || key === 'name' || key === 'stack' || key === 'cause') continue;\n    (err as unknown as Record<string, unknown>)[key] = s[key];\n  }\n\n  return err;\n}\n","/**\n * Class registry for the pubsub frame codec. Lets non-plain class instances\n * (e.g. GeneratedFile) survive a JSON round-trip across the unix-socket\n * pubsub transport.\n *\n * The registry is closed: only classes registered here can be reconstructed.\n * Unknown class names decode to plain data so we never instantiate\n * user-controlled prototypes.\n */\nexport interface ClassCodec<T = unknown, D = unknown> {\n  /** Convert an instance to JSON-safe data. */\n  toData(instance: T): D;\n  /** Reconstruct an instance from previously serialized data. */\n  fromData(data: D): T;\n}\n\nconst classRegistry = new Map<string, ClassCodec>();\n\nexport function registerClass<T, D>(name: string, codec: ClassCodec<T, D>): void {\n  classRegistry.set(name, codec as ClassCodec);\n}\n\nexport function unregisterClass(name: string): void {\n  classRegistry.delete(name);\n}\n\nexport function getClassCodec(name: string): ClassCodec | undefined {\n  return classRegistry.get(name);\n}\n\nexport function hasClassCodec(name: string): boolean {\n  return classRegistry.has(name);\n}\n","import type { ToolSet } from '@internal/ai-sdk-v5';\nimport { DefaultGeneratedFile, DefaultGeneratedFileWithType } from '../../stream/aisdk/v5/file';\nimport { DefaultStepResult } from '../../stream/aisdk/v5/output-helpers';\nimport { registerClass } from './registry';\n\ntype GeneratedFileData = { data: string; mediaType: string };\n\ntype DefaultStepResultData = {\n  content: DefaultStepResult<ToolSet>['content'];\n  finishReason: DefaultStepResult<ToolSet>['finishReason'];\n  usage: DefaultStepResult<ToolSet>['usage'];\n  warnings: DefaultStepResult<ToolSet>['warnings'];\n  request: DefaultStepResult<ToolSet>['request'];\n  response: DefaultStepResult<ToolSet>['response'];\n  providerMetadata: DefaultStepResult<ToolSet>['providerMetadata'];\n  tripwire?: DefaultStepResult<ToolSet>['tripwire'];\n};\n\n/**\n * Built-in class codecs.\n *\n * IMPORTANT: This module must be loaded for `instanceof` to survive\n * serialization across the unix-socket pubsub. The `BUILTIN_CODECS_REGISTERED`\n * constant exists solely so consumers can import it as a *named* import (see\n * `codec.ts`). A `import './registrations'` side-effect import would be\n * tree-shaken out of dist consumers because `packages/core/package.json` has\n * `\"sideEffects\": false`. Named imports are kept by every bundler.\n *\n * `DefaultStepResult` flows through workflow output schemas typed as\n * `z.any()`, so it crosses the unix-socket pubsub on the evented engine.\n * Without a class registration, consumers that rely on\n * `instanceof DefaultStepResult` (e.g. the in-memory storage path) would\n * receive plain data. The generic `TOOLS` type parameter is erased at\n * runtime, so we register the constructor name once.\n */\nexport const BUILTIN_CODECS_REGISTERED = (() => {\n  registerClass<DefaultGeneratedFile, GeneratedFileData>('DefaultGeneratedFile', {\n    toData: f => ({ data: f.base64, mediaType: f.mediaType }),\n    fromData: d => new DefaultGeneratedFile({ data: d.data, mediaType: d.mediaType }),\n  });\n\n  registerClass<DefaultGeneratedFileWithType, GeneratedFileData>('DefaultGeneratedFileWithType', {\n    toData: f => ({ data: f.base64, mediaType: f.mediaType }),\n    fromData: d => new DefaultGeneratedFileWithType({ data: d.data, mediaType: d.mediaType }),\n  });\n\n  registerClass<DefaultStepResult<ToolSet>, DefaultStepResultData>('DefaultStepResult', {\n    toData: s => ({\n      content: s.content,\n      finishReason: s.finishReason,\n      usage: s.usage,\n      warnings: s.warnings,\n      request: s.request,\n      response: s.response,\n      providerMetadata: s.providerMetadata,\n      tripwire: s.tripwire,\n    }),\n    fromData: d => new DefaultStepResult(d),\n  });\n\n  return true as const;\n})();\n","import type { SerializedError } from '../../error/utils';\n\n/**\n * Discriminator key for codec-tagged envelopes. Long, namespaced, and unlikely\n * to collide with user data. Plain objects that happen to carry this key but\n * do not match an envelope shape are preserved as-is by the decoder.\n *\n * **Reservation contract:** `__m_codec__` is reserved on the wire. Do not use\n * it as a property name on user objects that cross the pubsub boundary — if\n * the surrounding shape also happens to match an envelope (e.g.\n * `{ __m_codec__: 'Date', v: '...' }`), the decoder will reconstruct it as\n * the tagged type. The conservative shape check in `isEnvelope` keeps the\n * blast radius narrow, and `toJSON()` is skipped for objects already carrying\n * this key, but the safest path is to avoid the namespace entirely.\n */\nexport const CODEC_TAG = '__m_codec__';\n\n/**\n * Upper bound on the `source` length of a serialized `RegExp` envelope. Real\n * serialized regexes are tens of characters at most; this is generous headroom\n * for unusual patterns while bounding the input to `new RegExp(...)` on decode\n * so a hostile peer cannot push an unbounded pattern through the constructor.\n */\nexport const MAX_REGEXP_SOURCE_LENGTH = 1024;\n\nexport type EnvelopeTag = 'Date' | 'Error' | 'Map' | 'Set' | 'RegExp' | 'URL' | 'BigInt' | 'Undefined' | 'Class';\n\nexport type Envelope =\n  | { [CODEC_TAG]: 'Date'; v: string }\n  | { [CODEC_TAG]: 'Error'; v: SerializedError }\n  | { [CODEC_TAG]: 'Map'; v: Array<[unknown, unknown]> }\n  | { [CODEC_TAG]: 'Set'; v: Array<unknown> }\n  | { [CODEC_TAG]: 'RegExp'; v: { source: string; flags: string } }\n  | { [CODEC_TAG]: 'URL'; v: string }\n  | { [CODEC_TAG]: 'BigInt'; v: string }\n  | { [CODEC_TAG]: 'Undefined' }\n  | { [CODEC_TAG]: 'Class'; n: string; v: unknown };\n\n/**\n * Returns true when `value` looks like a codec envelope. The check is\n * conservative — an object with the tag key but an unknown tag value, or a\n * shape that does not match any envelope variant, is treated as user data.\n */\nexport function isEnvelope(value: object): value is Envelope {\n  const tag = (value as Record<string, unknown>)[CODEC_TAG];\n  if (typeof tag !== 'string') return false;\n  switch (tag as EnvelopeTag) {\n    case 'Undefined':\n      return true;\n    case 'Date':\n    case 'BigInt':\n    case 'URL':\n      return typeof (value as { v?: unknown }).v === 'string';\n    case 'RegExp': {\n      const v = (value as { v?: unknown }).v as { source?: unknown; flags?: unknown } | undefined;\n      if (!v || typeof v.source !== 'string' || typeof v.flags !== 'string') return false;\n      // Bound `source` to a generous-but-finite length. Real-world serialized\n      // regexes are tiny; a multi-KB pattern is either corrupted or a ReDoS\n      // attempt, and we'd rather treat it as user data than feed it into\n      // `new RegExp(...)`.\n      if (v.source.length > MAX_REGEXP_SOURCE_LENGTH) return false;\n      // Bound flags to the spec-defined RegExp flags. Anything else means the\n      // payload is either corrupted or hostile — treat as user data and let\n      // the decoder skip envelope reconstruction.\n      return /^[dgimsuvy]*$/.test(v.flags) && new Set(v.flags).size === v.flags.length;\n    }\n    case 'Map':\n    case 'Set':\n      return Array.isArray((value as { v?: unknown }).v);\n    case 'Error':\n      return typeof (value as { v?: unknown }).v === 'object' && (value as { v?: unknown }).v !== null;\n    case 'Class':\n      return typeof (value as { n?: unknown }).n === 'string';\n    default:\n      return false;\n  }\n}\n","import type { SerializedError } from '../../error/utils';\nimport { rehydrateError, serializeError } from './error';\nimport { BUILTIN_CODECS_REGISTERED } from './registrations';\nimport { getClassCodec } from './registry';\n// Named import of built-in registrations. Importing for side-effects only\n// (`import './registrations'`) gets tree-shaken by bundlers honoring\n// `\"sideEffects\": false` in `packages/core/package.json`. We import a value\n// and *use it* (inside `encode`) so the bundler must keep the module and\n// run its IIFE.\nimport { CODEC_TAG, MAX_REGEXP_SOURCE_LENGTH, isEnvelope } from './tags';\nimport type { Envelope } from './tags';\n\n/**\n * Encode a value into a JSON-safe shape. Non-JSON-safe types (Date, Error,\n * Map, Set, RegExp, URL, BigInt, undefined, registered classes) are wrapped\n * in tagged envelopes that the decoder can reconstruct.\n *\n * Functions and symbols are dropped (parity with JSON.stringify). Cycles are\n * replaced with null at the second visit. NaN/Infinity become null. Honors\n * user `toJSON()` methods on plain objects.\n */\nexport function encode(value: unknown): unknown {\n  // Reference `BUILTIN_CODECS_REGISTERED` so bundlers (which honor\n  // `\"sideEffects\": false` in package.json) cannot drop the registrations\n  // module. The IIFE there runs at module evaluation, registering built-in\n  // class codecs (DefaultGeneratedFile, DefaultStepResult, etc.). Without\n  // this reference, those `registerClass` calls get tree-shaken from dist\n  // and `instanceof` checks fail after a roundtrip across UnixSocketPubSub.\n  if (!BUILTIN_CODECS_REGISTERED) {\n    throw new Error('Built-in codec registrations failed to load');\n  }\n  return walk(value, new WeakSet());\n}\n\nfunction walk(v: unknown, seen: WeakSet<object>): unknown {\n  if (v === undefined) return { [CODEC_TAG]: 'Undefined' } satisfies Envelope;\n  if (v === null) return null;\n\n  const t = typeof v;\n  if (t === 'string' || t === 'boolean') return v;\n  if (t === 'number') return Number.isFinite(v as number) ? v : null;\n  if (t === 'bigint') return { [CODEC_TAG]: 'BigInt', v: (v as bigint).toString() } satisfies Envelope;\n  if (t === 'function' || t === 'symbol') return undefined;\n\n  if (v instanceof Date) {\n    return { [CODEC_TAG]: 'Date', v: v.toISOString() } satisfies Envelope;\n  }\n  if (v instanceof RegExp) {\n    return { [CODEC_TAG]: 'RegExp', v: { source: v.source, flags: v.flags } } satisfies Envelope;\n  }\n  if (v instanceof URL) {\n    return { [CODEC_TAG]: 'URL', v: v.toString() } satisfies Envelope;\n  }\n  if (v instanceof Error) {\n    return { [CODEC_TAG]: 'Error', v: serializeError(v) } satisfies Envelope;\n  }\n  if (v instanceof Map) {\n    if (seen.has(v)) return null;\n    seen.add(v);\n    const entries: Array<[unknown, unknown]> = [];\n    for (const [k, val] of v.entries()) {\n      entries.push([walk(k, seen), walk(val, seen)]);\n    }\n    return { [CODEC_TAG]: 'Map', v: entries } satisfies Envelope;\n  }\n  if (v instanceof Set) {\n    if (seen.has(v)) return null;\n    seen.add(v);\n    const values: unknown[] = [];\n    for (const x of v) values.push(walk(x, seen));\n    return { [CODEC_TAG]: 'Set', v: values } satisfies Envelope;\n  }\n\n  if (Array.isArray(v)) {\n    if (seen.has(v)) return null;\n    seen.add(v);\n    return v.map(x => walk(x, seen));\n  }\n\n  if (t === 'object') {\n    if (seen.has(v as object)) return null;\n    seen.add(v as object);\n\n    // Honor toJSON() — matches JSON.stringify behavior. Skip for plain objects\n    // that already carry a literal CODEC_TAG key (defensive: do not let user\n    // data masquerade as an envelope through toJSON).\n    const maybeToJSON = (v as { toJSON?: unknown }).toJSON;\n    if (typeof maybeToJSON === 'function') {\n      return walk((maybeToJSON as () => unknown).call(v), seen);\n    }\n\n    // Class registry lookup by exact constructor name. Falls back to plain\n    // object walk for unknown classes — they decode to plain data.\n    const ctor = (v as object).constructor;\n    const ctorName = ctor?.name;\n    if (ctorName && ctorName !== 'Object') {\n      const reg = getClassCodec(ctorName);\n      if (reg) {\n        return {\n          [CODEC_TAG]: 'Class',\n          n: ctorName,\n          v: walk(reg.toData(v), seen),\n        } satisfies Envelope;\n      }\n    }\n\n    const out: Record<string, unknown> = {};\n    for (const k of Object.keys(v as Record<string, unknown>)) {\n      // Skip prototype-pollution vectors.\n      if (k === '__proto__') continue;\n      const raw = (v as Record<string, unknown>)[k];\n      if (raw === undefined) {\n        // Preserve explicit undefined keys (JSON.stringify drops them).\n        out[k] = { [CODEC_TAG]: 'Undefined' } satisfies Envelope;\n        continue;\n      }\n      const encoded = walk(raw, seen);\n      if (encoded === undefined) continue; // function/symbol fields are dropped\n      out[k] = encoded;\n    }\n    return out;\n  }\n\n  return v;\n}\n\n/**\n * Reconstruct a `RegExp` from a decoded envelope payload. Re-validates the\n * payload locally (independent of `isEnvelope`) so the constructor input is\n * narrowed at this single call site: bounded `source` length, spec-defined\n * flag whitelist, no duplicate flags. A malformed or hostile envelope yields\n * an empty regex (`/(?:)/`) rather than throwing, keeping frame decoding\n * resilient.\n *\n * NOTE: `source` is intentionally NOT escaped — a `RegExp` envelope's whole\n * purpose is to round-trip a pattern, so metacharacters must reach\n * `new RegExp(...)` verbatim. Safety comes from the bounded length, the\n * flag whitelist, and the `try/catch` fallback, not from escaping.\n */\nfunction decodeRegExpEnvelope(v: unknown): RegExp {\n  if (!v || typeof v !== 'object') return /(?:)/;\n  const candidate = v as { source?: unknown; flags?: unknown };\n  if (typeof candidate.source !== 'string') return /(?:)/;\n  if (typeof candidate.flags !== 'string') return /(?:)/;\n  if (candidate.source.length > MAX_REGEXP_SOURCE_LENGTH) return /(?:)/;\n  const flags = candidate.flags;\n  if (!/^[dgimsuvy]*$/.test(flags)) return /(?:)/;\n  if (new Set(flags).size !== flags.length) return /(?:)/;\n  try {\n    // lgtm[js/regex-injection] -- safe because of three structural gates that\n    // run before this line: (1) `isEnvelope` rejects payloads whose `source`\n    // exceeds `MAX_REGEXP_SOURCE_LENGTH` and whose `flags` fall outside the\n    // spec flag whitelist `[dgimsuvy]` (no duplicates); (2) the local checks\n    // above re-narrow `source`/`flags` so the constructor input is bounded at\n    // this single call site; (3) this `try/catch` swallows any constructor\n    // failure and returns the empty regex. Do NOT \"fix\" this by escaping\n    // `source` — see the helper docstring above for why that would break the\n    // round-trip contract.\n    return new RegExp(candidate.source, flags);\n  } catch {\n    return /(?:)/;\n  }\n}\n\n/**\n * Decode a value previously produced by `encode`. Reconstructs envelope-tagged\n * types and recursively decodes nested values. Plain objects that happen to\n * carry a `CODEC_TAG` key but do not match an envelope shape are preserved.\n */\nexport function decode(value: unknown): unknown {\n  if (value === null) return null;\n  if (typeof value !== 'object') return value;\n\n  if (Array.isArray(value)) return value.map(decode);\n\n  if (CODEC_TAG in value && isEnvelope(value)) {\n    const env = value as Envelope;\n    switch (env[CODEC_TAG]) {\n      case 'Undefined':\n        return undefined;\n      case 'Date':\n        return new Date(env.v);\n      case 'BigInt':\n        return BigInt(env.v);\n      case 'RegExp':\n        return decodeRegExpEnvelope(env.v);\n      case 'URL':\n        return new URL(env.v);\n      case 'Map':\n        return new Map(env.v.map(([k, val]) => [decode(k), decode(val)]));\n      case 'Set':\n        return new Set(env.v.map(decode));\n      case 'Error':\n        return rehydrateError(decodeSerializedError(env.v));\n      case 'Class': {\n        const reg = getClassCodec(env.n);\n        const data = decode(env.v);\n        return reg ? reg.fromData(data) : data;\n      }\n    }\n  }\n\n  const out: Record<string, unknown> = {};\n  for (const k of Object.keys(value as Record<string, unknown>)) {\n    if (k === '__proto__') continue;\n    out[k] = decode((value as Record<string, unknown>)[k]);\n  }\n  return out;\n}\n\n/**\n * Recursively decode any envelopes embedded inside a SerializedError's custom\n * fields (e.g. an error with a `details: Map` field). The top-level shape\n * stays a SerializedError so `rehydrateError` can consume it.\n */\nfunction decodeSerializedError(s: SerializedError): SerializedError {\n  const decoded = decode(s) as SerializedError;\n  return decoded;\n}\n","import { randomUUID } from 'node:crypto';\nimport { mkdir, open, stat, unlink } from 'node:fs/promises';\nimport type { FileHandle } from 'node:fs/promises';\nimport net from 'node:net';\nimport { dirname } from 'node:path';\n\nimport { decode, encode } from './codec';\nimport { PubSub } from './pubsub';\nimport type { PubSubDeliveryMode } from './pubsub';\nimport type { Event, EventCallback, SubscribeOptions } from './types';\n\ntype ClientFrame =\n  | { type: 'subscribe'; topic: string }\n  | { type: 'unsubscribe'; topic: string }\n  | { type: 'publish'; topic: string; event: Omit<Event, 'id' | 'createdAt'>; localOnly?: boolean }\n  | { type: 'ack'; id?: string }\n  | { type: 'nack'; id?: string };\n\ntype ServerFrame = { type: 'event'; topic: string; event: Event } | { type: 'subscribed'; topic: string };\n\ntype UnixSocketPubSubOptions = {\n  maxRemoteClientQueuedBytes?: number;\n};\n\ntype BrokerClient = {\n  socket: net.Socket;\n  subscriptions: Set<string>;\n  writeChain: Promise<void>;\n  queuedBytes: number;\n};\n\ntype SubscribeWaiter = {\n  resolve: () => void;\n  reject: (error: Error) => void;\n};\n\nconst DEFAULT_MAX_REMOTE_CLIENT_QUEUED_BYTES = 64 * 1024 * 1024;\n\n/**\n * Max number of times a local subscriber callback may be redelivered after a\n * nack. MUST be >= the consumer-side retry budget\n * (`WorkflowEventProcessor.MAX_DELIVERY_ATTEMPTS`) — otherwise the transport\n * gives up before the consumer can exhaust its budget and surface the terminal\n * failure, which would leave the run silently hung.\n *\n * An invariant test in\n * `packages/core/src/events/unix-socket-pubsub-redelivery-budget.test.ts`\n * pins this ordering against the consumer constant so the two constants stay\n * in sync as the consumer budget changes.\n */\nexport const MAX_LOCAL_REDELIVERIES = 6;\nconst REDELIVERY_DELAY_MS = 100;\n\nfunction serializeFrame(frame: ClientFrame | ServerFrame): string {\n  // Encode through the codec so non-JSON-safe values (Date, Error, Map, Set,\n  // RegExp, URL, BigInt, undefined, registered classes) survive the wire\n  // round-trip via tagged envelopes.\n  return `${JSON.stringify(encode(frame))}\\n`;\n}\n\nfunction writeSerializedFrame(socket: net.Socket, serializedFrame: string): Promise<void> {\n  return new Promise((resolve, reject) => {\n    let writeCompleted = false;\n    let drainCompleted = true;\n    let settled = false;\n\n    const cleanup = () => {\n      socket.off('error', onError);\n      socket.off('close', onClose);\n      socket.off('drain', onDrain);\n    };\n    const settle = (error?: Error) => {\n      if (settled) return;\n      settled = true;\n      cleanup();\n      if (error) {\n        reject(error);\n        return;\n      }\n      resolve();\n    };\n    const maybeResolve = () => {\n      if (writeCompleted && drainCompleted) {\n        settle();\n      }\n    };\n    const onError = (error: Error) => settle(error);\n    // NOTE: keep this exact message in sync with the transient-error classifier\n    // in #sendToBroker (search for 'socket closed before write completed').\n    const onClose = () => settle(new Error('UnixSocketPubSub socket closed before write completed'));\n    const onDrain = () => {\n      drainCompleted = true;\n      maybeResolve();\n    };\n\n    socket.once('error', onError);\n    socket.once('close', onClose);\n    let drained: boolean;\n    try {\n      drained = socket.write(serializedFrame, error => {\n        if (error) {\n          settle(error);\n          return;\n        }\n        writeCompleted = true;\n        maybeResolve();\n      });\n    } catch (error) {\n      settle(error as Error);\n      return;\n    }\n    if (!drained) {\n      drainCompleted = false;\n      socket.once('drain', onDrain);\n    }\n  });\n}\n\nfunction writeFrame(socket: net.Socket, frame: ClientFrame | ServerFrame): Promise<void> {\n  return writeSerializedFrame(socket, serializeFrame(frame));\n}\n\nfunction nextTick(): Promise<void> {\n  return new Promise(resolve => setImmediate(resolve));\n}\n\nfunction readFrames(socket: net.Socket, onFrame: (frame: any) => void) {\n  let buffer = '';\n  socket.setEncoding('utf8');\n  socket.on('data', chunk => {\n    buffer += chunk;\n    while (true) {\n      const newlineIndex = buffer.indexOf('\\n');\n      if (newlineIndex === -1) break;\n      const line = buffer.slice(0, newlineIndex);\n      buffer = buffer.slice(newlineIndex + 1);\n      if (!line.trim()) continue;\n      try {\n        onFrame(decode(JSON.parse(line)));\n      } catch {\n        // Ignore malformed frames. The transport is local IPC and callers can retry.\n      }\n    }\n  });\n}\n\nexport class UnixSocketPubSub extends PubSub {\n  readonly socketPath: string;\n  #server?: net.Server;\n  #clientSocket?: net.Socket;\n  #isBroker = false;\n  #closed = false;\n  #starting?: Promise<void>;\n  #callbacks = new Map<string, Set<EventCallback>>();\n  #subscribeWaiters = new Map<string, SubscribeWaiter[]>();\n  #brokerClients = new Map<net.Socket, BrokerClient>();\n  #pendingWrites = new Set<Promise<void>>();\n  #recovering?: Promise<void>;\n  #maxRemoteClientQueuedBytes: number;\n\n  constructor(socketPath: string, options: UnixSocketPubSubOptions = {}) {\n    super();\n    this.socketPath = socketPath;\n    this.#maxRemoteClientQueuedBytes = options.maxRemoteClientQueuedBytes ?? DEFAULT_MAX_REMOTE_CLIENT_QUEUED_BYTES;\n  }\n\n  override get supportedModes(): ReadonlyArray<PubSubDeliveryMode> {\n    return ['push'];\n  }\n\n  get isBroker(): boolean {\n    return this.#isBroker;\n  }\n\n  /** Number of remote clients currently connected to this broker. Always 0 for non-broker instances. */\n  get remoteClientCount(): number {\n    return this.#isBroker ? this.#brokerClients.size : 0;\n  }\n\n  async publish(\n    topic: string,\n    event: Omit<Event, 'id' | 'createdAt'>,\n    options?: { localOnly?: boolean },\n  ): Promise<void> {\n    await this.#ensureStarted();\n\n    // `localOnly` events stay entirely within the publishing process. They are\n    // never serialized over a unix socket, so live methods on payload values\n    // (e.g. `MastraModelOutput.getFullOutput`, `step.condition` functions on\n    // serialized step graphs) survive intact. This is the semantic the agent's\n    // execution-workflow relies on: the run result is delivered via\n    // `workflows-finish` and includes the `MastraModelOutput` instance —\n    // round-tripping it through the broker would strip its methods.\n    if (options?.localOnly) {\n      const localEvent: Event = {\n        ...event,\n        id: randomUUID(),\n        createdAt: new Date(),\n        deliveryAttempt: 1,\n      };\n      this.#deliverLocal(topic, localEvent);\n      return;\n    }\n\n    if (this.#isBroker) {\n      await this.#publishFromBroker(topic, event, undefined, options?.localOnly);\n      return;\n    }\n\n    const socket = this.#clientSocket;\n    if (!socket || socket.destroyed) {\n      await this.#ensureStarted(true);\n    }\n    await this.#sendToBroker({ type: 'publish', topic, event, localOnly: options?.localOnly });\n  }\n\n  async subscribe(topic: string, cb: EventCallback, options?: SubscribeOptions): Promise<void> {\n    if (options?.group) {\n      throw new Error('UnixSocketPubSub does not support grouped subscriptions yet');\n    }\n\n    const callbacks = this.#callbacks.get(topic) ?? new Set<EventCallback>();\n    const hadCallback = callbacks.has(cb);\n    const wasConnected = Boolean(this.#clientSocket && !this.#clientSocket.destroyed);\n    callbacks.add(cb);\n    this.#callbacks.set(topic, callbacks);\n\n    try {\n      await this.#ensureStarted();\n      if (!this.#isBroker && !hadCallback && wasConnected) {\n        await this.#sendSubscribeToBroker(topic);\n      }\n    } catch (error) {\n      if (!hadCallback) {\n        callbacks.delete(cb);\n        if (callbacks.size === 0) {\n          this.#callbacks.delete(topic);\n        }\n      }\n      throw error;\n    }\n  }\n\n  async unsubscribe(topic: string, cb: EventCallback): Promise<void> {\n    const callbacks = this.#callbacks.get(topic);\n    callbacks?.delete(cb);\n    if (callbacks?.size === 0) {\n      this.#callbacks.delete(topic);\n      if (!this.#isBroker && this.#clientSocket && !this.#clientSocket.destroyed) {\n        await this.#sendToBroker({ type: 'unsubscribe', topic });\n        await nextTick();\n      }\n    }\n  }\n\n  async flush(): Promise<void> {\n    await Promise.allSettled([...this.#pendingWrites]);\n  }\n\n  async close(): Promise<void> {\n    this.#closed = true;\n    this.#callbacks.clear();\n\n    this.#clientSocket?.destroy();\n    this.#clientSocket = undefined;\n    this.#rejectSubscribeWaiters(new Error('UnixSocketPubSub is closed'));\n\n    for (const client of [...this.#brokerClients.values()]) {\n      this.#removeBrokerClient(client);\n    }\n\n    if (this.#server) {\n      await new Promise<void>(resolve => this.#server?.close(() => resolve()));\n      this.#server = undefined;\n    }\n\n    if (this.#isBroker) {\n      await unlink(this.socketPath).catch(() => {});\n    }\n    this.#isBroker = false;\n  }\n\n  async #ensureStarted(forceReconnect = false): Promise<void> {\n    if (this.#closed) {\n      throw new Error('UnixSocketPubSub is closed');\n    }\n    if (!forceReconnect && (this.#isBroker || (this.#clientSocket && !this.#clientSocket.destroyed))) {\n      return;\n    }\n    if (this.#starting) {\n      return this.#starting;\n    }\n\n    this.#starting = this.#start(forceReconnect).finally(() => {\n      this.#starting = undefined;\n    });\n    return this.#starting;\n  }\n\n  async #start(forceReconnect: boolean): Promise<void> {\n    if (forceReconnect) {\n      this.#clientSocket?.destroy();\n      this.#clientSocket = undefined;\n      this.#isBroker = false;\n    }\n\n    this.#throwIfClosed();\n    await mkdir(dirname(this.socketPath), { recursive: true });\n    this.#throwIfClosed();\n\n    try {\n      await this.#listen();\n      this.#throwIfClosed();\n      this.#isBroker = true;\n      return;\n    } catch (error) {\n      if (this.#closed) {\n        await this.close();\n        throw new Error('UnixSocketPubSub is closed');\n      }\n      const code = (error as NodeJS.ErrnoException).code;\n      // EADDRINUSE: another broker bound the socket. EEXIST: another process\n      // created the socket file but hasn't bound yet (macOS race). Both mean\n      // \"fall through and try to connect as a client\".\n      if (code !== 'EADDRINUSE' && code !== 'EEXIST') throw error;\n    }\n\n    try {\n      await this.#connectClient();\n      this.#throwIfClosed();\n    } catch (error) {\n      if (this.#closed) {\n        await this.close();\n        throw new Error('UnixSocketPubSub is closed');\n      }\n      const code = (error as NodeJS.ErrnoException).code;\n      if (code === 'ECONNREFUSED' || code === 'ENOENT' || code === 'ENOTSOCK') {\n        this.#throwIfClosed();\n        await this.#electBroker();\n        return;\n      }\n      throw error;\n    }\n  }\n\n  #throwIfClosed() {\n    if (this.#closed) {\n      throw new Error('UnixSocketPubSub is closed');\n    }\n  }\n\n  #listen(): Promise<void> {\n    return new Promise((resolve, reject) => {\n      const server = net.createServer(socket => this.#handleBrokerClient(socket));\n      const onError = (error: Error) => {\n        server.off('listening', onListening);\n        reject(error);\n      };\n      const onListening = () => {\n        server.off('error', onError);\n        this.#server = server;\n        resolve();\n      };\n\n      server.once('error', onError);\n      server.once('listening', onListening);\n      server.listen(this.socketPath);\n    });\n  }\n\n  #connectClient(): Promise<void> {\n    return new Promise((resolve, reject) => {\n      const socket = net.createConnection(this.socketPath);\n      const onError = (error: Error) => {\n        socket.off('connect', onConnect);\n        reject(error);\n      };\n      const onConnect = () => {\n        socket.off('error', onError);\n        this.#clientSocket = socket;\n        this.#isBroker = false;\n        readFrames(socket, frame => this.#handleServerFrame(frame));\n        // NOTE: keep this exact message in sync with the transient-error\n        // classifier in #sendToBroker (search for 'broker connection closed').\n        socket.on('close', () =>\n          this.#handleClientDisconnect(socket, new Error('UnixSocketPubSub broker connection closed')),\n        );\n        socket.on('error', error => this.#handleClientDisconnect(socket, error));\n        void this.#resubscribeClient().then(resolve, reject);\n      };\n\n      socket.once('error', onError);\n      socket.once('connect', onConnect);\n    });\n  }\n\n  async #resubscribeClient() {\n    for (const topic of this.#callbacks.keys()) {\n      await this.#sendSubscribeToBroker(topic);\n    }\n  }\n\n  #handleClientDisconnect(socket: net.Socket, error: Error) {\n    if (this.#clientSocket !== socket) return;\n    this.#clientSocket = undefined;\n    this.#rejectSubscribeWaiters(error);\n    if (!this.#closed) {\n      void this.#recoverClientConnection();\n    }\n  }\n\n  async #recoverClientConnection(): Promise<void> {\n    if (this.#recovering) return this.#recovering;\n    this.#recovering = this.#recoverClientConnectionLoop().finally(() => {\n      this.#recovering = undefined;\n    });\n    return this.#recovering;\n  }\n\n  async #recoverClientConnectionLoop(): Promise<void> {\n    while (!this.#closed && !this.#isBroker && !(this.#clientSocket && !this.#clientSocket.destroyed)) {\n      try {\n        await this.#ensureStarted(true);\n        return;\n      } catch {\n        if (this.#closed) return;\n        await new Promise(resolve => setTimeout(resolve, 10));\n      }\n    }\n  }\n\n  /**\n   * Serializes broker election across processes using an exclusive lock file.\n   * Only the lock winner unlinks the stale socket and listens; losers wait\n   * then connect as clients to the newly elected broker.\n   */\n  async #electBroker(): Promise<void> {\n    const lockPath = this.socketPath + '.elect';\n    let lockFd: FileHandle | undefined;\n    try {\n      lockFd = await open(lockPath, 'wx');\n    } catch (e) {\n      if ((e as NodeJS.ErrnoException).code === 'EEXIST') {\n        if (await this.#isElectionLockStale(lockPath)) {\n          await unlink(lockPath).catch(() => {});\n          throw new Error('Stale broker election lock removed');\n        }\n        await new Promise(resolve => setTimeout(resolve, 150));\n        try {\n          await this.#connectClient();\n          this.#throwIfClosed();\n          return;\n        } catch {\n          throw new Error('Broker election in progress by another process');\n        }\n      }\n      throw e;\n    }\n\n    try {\n      // Re-check: a previous election round may have installed a broker\n      // between our initial connectClient() and acquiring this lock.\n      try {\n        await this.#connectClient();\n        this.#throwIfClosed();\n        return;\n      } catch {\n        // Still no live broker — proceed with election.\n      }\n      await unlink(this.socketPath).catch(() => {});\n      this.#throwIfClosed();\n      await this.#listen();\n      this.#throwIfClosed();\n      this.#isBroker = true;\n    } finally {\n      await lockFd.close().catch(() => {});\n      await unlink(lockPath).catch(() => {});\n    }\n  }\n\n  async #isElectionLockStale(lockPath: string): Promise<boolean> {\n    try {\n      const lockStat = await stat(lockPath);\n      return Date.now() - lockStat.mtimeMs > 2000;\n    } catch {\n      return true;\n    }\n  }\n\n  async #sendSubscribeToBroker(topic: string): Promise<void> {\n    let waiter: SubscribeWaiter | undefined;\n    const subscribed = new Promise<void>((resolve, reject) => {\n      waiter = { resolve, reject };\n      const waiters = this.#subscribeWaiters.get(topic) ?? [];\n      waiters.push(waiter);\n      this.#subscribeWaiters.set(topic, waiters);\n    });\n    try {\n      await this.#sendToBroker({ type: 'subscribe', topic });\n    } catch (error) {\n      this.#removeSubscribeWaiter(topic, waiter);\n      throw error;\n    }\n    await subscribed;\n  }\n\n  #removeSubscribeWaiter(topic: string, waiter: SubscribeWaiter | undefined) {\n    if (!waiter) return;\n    const waiters = this.#subscribeWaiters.get(topic);\n    if (!waiters) return;\n    const nextWaiters = waiters.filter(item => item !== waiter);\n    if (nextWaiters.length === 0) {\n      this.#subscribeWaiters.delete(topic);\n      return;\n    }\n    this.#subscribeWaiters.set(topic, nextWaiters);\n  }\n\n  #settleSubscribeWaiters(topic: string, error?: Error) {\n    const waiters = this.#subscribeWaiters.get(topic);\n    this.#subscribeWaiters.delete(topic);\n    if (error) {\n      waiters?.forEach(waiter => waiter.reject(error));\n      return;\n    }\n    waiters?.forEach(waiter => waiter.resolve());\n  }\n\n  #rejectSubscribeWaiters(error: Error) {\n    for (const topic of this.#subscribeWaiters.keys()) {\n      this.#settleSubscribeWaiters(topic, error);\n    }\n  }\n\n  #handleBrokerClient(socket: net.Socket) {\n    const client: BrokerClient = {\n      socket,\n      subscriptions: new Set(),\n      writeChain: Promise.resolve(),\n      queuedBytes: 0,\n    };\n    this.#brokerClients.set(socket, client);\n    readFrames(socket, frame => {\n      const clientFrame = frame as ClientFrame;\n      if (clientFrame.type === 'subscribe') {\n        client.subscriptions.add(clientFrame.topic);\n        this.#enqueueBrokerClientWrite(client, { type: 'subscribed', topic: clientFrame.topic });\n      } else if (clientFrame.type === 'unsubscribe') {\n        client.subscriptions.delete(clientFrame.topic);\n      } else if (clientFrame.type === 'publish') {\n        void this.#publishFromBroker(clientFrame.topic, clientFrame.event, client, clientFrame.localOnly);\n      }\n    });\n    socket.on('close', () => this.#removeBrokerClient(client));\n    socket.on('error', () => this.#removeBrokerClient(client));\n  }\n\n  #enqueueBrokerClientWrite(client: BrokerClient, frame: ServerFrame) {\n    if (this.#brokerClients.get(client.socket) !== client || client.socket.destroyed) return;\n\n    const serializedFrame = serializeFrame(frame);\n    const queuedBytes = Buffer.byteLength(serializedFrame);\n    if (client.queuedBytes + queuedBytes > this.#maxRemoteClientQueuedBytes) {\n      this.#removeBrokerClient(client);\n      return;\n    }\n\n    client.queuedBytes += queuedBytes;\n\n    const write = client.writeChain\n      .catch(() => {})\n      .then(async () => {\n        if (this.#brokerClients.get(client.socket) !== client || client.socket.destroyed) return;\n        await writeSerializedFrame(client.socket, serializedFrame);\n      })\n      .catch(() => {\n        this.#removeBrokerClient(client);\n      })\n      .finally(() => {\n        client.queuedBytes = Math.max(0, client.queuedBytes - queuedBytes);\n      });\n\n    client.writeChain = write;\n    this.#pendingWrites.add(write);\n    void write.finally(() => this.#pendingWrites.delete(write));\n  }\n\n  #removeBrokerClient(client: BrokerClient) {\n    if (this.#brokerClients.get(client.socket) !== client) return;\n    this.#brokerClients.delete(client.socket);\n    client.subscriptions.clear();\n    client.queuedBytes = 0;\n    client.writeChain = Promise.resolve();\n    if (!client.socket.destroyed) {\n      client.socket.destroy();\n    }\n  }\n\n  #handleServerFrame(frame: ServerFrame) {\n    if (frame.type === 'subscribed') {\n      this.#settleSubscribeWaiters(frame.topic);\n      return;\n    }\n    if (frame.type !== 'event') return;\n    // `createdAt` is already a Date — the codec rehydrates it during JSON.parse\n    // in `readFrames`. No ad-hoc conversion needed.\n    this.#deliverLocal(frame.topic, frame.event);\n  }\n\n  async #publishFromBroker(\n    topic: string,\n    event: Omit<Event, 'id' | 'createdAt'>,\n    sourceClient?: BrokerClient,\n    localOnly?: boolean,\n  ) {\n    const brokerEvent: Event = {\n      ...event,\n      id: randomUUID(),\n      createdAt: new Date(),\n      deliveryAttempt: 1,\n    };\n\n    this.#deliverLocal(topic, brokerEvent);\n\n    // Skip serialization entirely when no remote clients could receive the event.\n    if (this.#brokerClients.size === 0) return;\n\n    // `localOnly` events are scoped to the publishing instance.\n    // When the publisher is the broker, the `#deliverLocal` above is enough.\n    // When the publisher is a remote client, relay the event back ONLY to\n    // that client so its subscription callback fires, but do NOT fan out to\n    // other clients — their WEP would just drop the event via `#ownsWorkflow`\n    // and the multi-MB payload would waste socket/kernel buffer for nothing.\n    if (localOnly) {\n      if (sourceClient && sourceClient.subscriptions.has(topic) && !sourceClient.socket.destroyed) {\n        this.#enqueueBrokerClientWrite(sourceClient, { type: 'event', topic, event: brokerEvent });\n      }\n      return;\n    }\n\n    let frame: ServerFrame | undefined;\n    for (const client of this.#brokerClients.values()) {\n      if (!client.subscriptions.has(topic) || client.socket.destroyed) continue;\n      // Lazily build the frame only when we know at least one client needs it.\n      frame ??= { type: 'event', topic, event: brokerEvent };\n      this.#enqueueBrokerClientWrite(client, frame);\n    }\n  }\n\n  #deliverLocal(topic: string, event: Event) {\n    const callbacks = this.#callbacks.get(topic);\n    if (!callbacks) return;\n    for (const cb of callbacks) {\n      this.#invokeLocalCallback(topic, event, cb, 0);\n    }\n  }\n\n  #invokeLocalCallback(topic: string, event: Event, cb: EventCallback, attempt: number) {\n    let nacked = false;\n    const nack = async () => {\n      if (nacked || this.#closed) return;\n      nacked = true;\n      if (attempt >= MAX_LOCAL_REDELIVERIES) return;\n      const stillSubscribed = this.#callbacks.get(topic)?.has(cb);\n      if (!stillSubscribed) return;\n      const timer = setTimeout(\n        () => {\n          if (this.#closed) return;\n          if (!this.#callbacks.get(topic)?.has(cb)) return;\n          const redeliveredEvent: Event = {\n            ...event,\n            deliveryAttempt: (event.deliveryAttempt ?? 1) + 1,\n          };\n          this.#invokeLocalCallback(topic, redeliveredEvent, cb, attempt + 1);\n        },\n        REDELIVERY_DELAY_MS * (attempt + 1),\n      );\n      // Unrefed so a queued redelivery never holds the event loop open at\n      // shutdown. The trade-off: an in-flight redelivery during process exit\n      // is silently dropped. That's acceptable because the consumer (WEP)\n      // is itself shutting down and the workflow will be re-driven from\n      // durable state on the next start.\n      timer.unref?.();\n    };\n    try {\n      const result = (cb as (event: Event, ack: () => Promise<void>, nack: () => Promise<void>) => unknown)(\n        event,\n        async () => {},\n        nack,\n      );\n      if (result && typeof (result as Promise<void>).catch === 'function') {\n        void (result as Promise<void>).catch(() => {});\n      }\n    } catch {\n      // Ignore subscriber failures so one callback cannot poison topic delivery.\n    }\n  }\n\n  async #sendToBroker(frame: ClientFrame) {\n    // If the broker died mid-write (EPIPE) or while election is rotating, we\n    // reconnect and retry. The first attempt is the normal path. Each retry\n    // forces a fresh broker resolution. Retry budget is bounded so a truly\n    // unreachable broker still errors instead of looping forever.\n    const maxRetries = 3;\n    let lastError: unknown;\n    for (let attempt = 0; attempt <= maxRetries; attempt++) {\n      try {\n        if (attempt === 0) {\n          await this.#sendToActiveBroker(frame);\n        } else {\n          if (this.#closed) throw lastError;\n          const failedSocket = this.#clientSocket;\n          this.#clientSocket = undefined;\n          failedSocket?.destroy();\n          await this.#ensureStarted(true);\n          await this.#sendToActiveBroker(frame);\n        }\n        return;\n      } catch (error) {\n        lastError = error;\n        if (this.#closed) throw error;\n        const code = (error as NodeJS.ErrnoException)?.code;\n        // EPIPE/ECONNRESET/ENOTCONN: broker died mid-write — retry against a\n        // fresh broker. Anything else (e.g. closed pubsub, validation error)\n        // is not safe to retry blindly. The string-message checks cover three\n        // internal errors thrown from within this file that don't carry an\n        // ErrnoException-style `code` — keep them in lockstep with those\n        // throw sites:\n        //   - \"socket closed before write completed\" (writeSerializedFrame,\n        //     when the broker dies mid-write before the drain settles)\n        //   - \"broker connection closed\" (#handleClientDisconnect)\n        //   - \"not connected to a broker\" (#sendToActiveBroker)\n        const transient =\n          code === 'EPIPE' ||\n          code === 'ECONNRESET' ||\n          code === 'ENOTCONN' ||\n          (error as Error)?.message?.includes('socket closed before write completed') ||\n          (error as Error)?.message?.includes('broker connection closed') ||\n          (error as Error)?.message?.includes('not connected to a broker');\n        if (!transient || attempt === maxRetries) throw error;\n        // Tiny backoff so concurrent senders don't dogpile re-election.\n        await new Promise(resolve => setTimeout(resolve, 10 * (attempt + 1)));\n      }\n    }\n  }\n\n  async #sendToActiveBroker(frame: ClientFrame) {\n    const socket = this.#clientSocket;\n    if (!socket || socket.destroyed) {\n      await this.#ensureStarted(true);\n    }\n    if (this.#isBroker) {\n      await this.#handlePromotedBrokerFrame(frame);\n      return;\n    }\n    const activeSocket = this.#clientSocket;\n    if (!activeSocket || activeSocket.destroyed) {\n      // NOTE: keep this exact message in sync with the transient-error\n      // classifier in #sendToBroker (search for 'not connected to a broker').\n      throw new Error('UnixSocketPubSub is not connected to a broker');\n    }\n    await writeFrame(activeSocket, frame);\n  }\n\n  async #handlePromotedBrokerFrame(frame: ClientFrame) {\n    if (frame.type === 'subscribe') {\n      this.#settleSubscribeWaiters(frame.topic);\n    } else if (frame.type === 'publish') {\n      await this.#publishFromBroker(frame.topic, frame.event);\n    }\n  }\n}\n"],"mappings":";;;;;;;;;;;;AAEA,MAAM,kBAAkB;;;;;;AAOxB,SAAgB,eAAe,KAAY,QAAQ,GAAoB;CACrE,MAAM,OAAwB;EAC5B,MAAM,IAAI,QAAQ;EAClB,SAAS,IAAI;CACf;CACA,IAAI,IAAI,UAAU,KAAA,GAAW,KAAK,QAAQ,IAAI;CAE9C,IAAI,IAAI,UAAU,KAAA,GAChB,IAAI,IAAI,iBAAiB,SAAS,QAAQ,iBACxC,KAAK,QAAQ,eAAe,IAAI,OAAO,QAAQ,CAAC;MAEhD,KAAK,QAAQ,IAAI;CAKrB,KAAK,MAAM,OAAO,KAAK;EACrB,IAAI,CAAC,OAAO,UAAU,eAAe,KAAK,KAAK,GAAG,GAAG;EACrD,IAAI,QAAQ,aAAa,QAAQ,UAAU,QAAQ,WAAW,QAAQ,SAAS;EAC/E,KAAK,OAAQ,IAA2C;CAC1D;CAEA,OAAO;AACT;;;;;AAMA,SAAgB,eAAe,GAA2B;CACxD,MAAM,QACJ,EAAE,UAAU,KAAA,IACR,EAAE,SAAS,OAAO,EAAE,UAAU,YAAY,aAAc,EAAE,SAAoB,UAAW,EAAE,QACzF,eAAe,EAAE,KAAwB,IACzC,EAAE,QACJ,KAAA;CAEN,MAAM,MAAM,UAAU,KAAA,IAAY,IAAI,MAAM,EAAE,SAAS,EAAE,MAAM,CAAC,IAAI,IAAI,MAAM,EAAE,OAAO;CACvF,IAAI,EAAE,MAAM,IAAI,OAAO,EAAE;CACzB,IAAI,EAAE,UAAU,KAAA,GAAW,IAAI,QAAQ,EAAE;CAEzC,KAAK,MAAM,OAAO,GAAG;EACnB,IAAI,CAAC,OAAO,UAAU,eAAe,KAAK,GAAG,GAAG,GAAG;EACnD,IAAI,QAAQ,aAAa,QAAQ,UAAU,QAAQ,WAAW,QAAQ,SAAS;EAC/E,IAA4C,OAAO,EAAE;CACvD;CAEA,OAAO;AACT;;;ACzCA,MAAM,gCAAgB,IAAI,IAAwB;AAElD,SAAgB,cAAoB,MAAc,OAA+B;CAC/E,cAAc,IAAI,MAAM,KAAmB;AAC7C;AAMA,SAAgB,cAAc,MAAsC;CAClE,OAAO,cAAc,IAAI,IAAI;AAC/B;;;;;;;;;;;;;;;;;;;;ACOA,MAAa,mCAAmC;CAC9C,cAAuD,wBAAwB;EAC7E,SAAQ,OAAM;GAAE,MAAM,EAAE;GAAQ,WAAW,EAAE;EAAU;EACvD,WAAU,MAAK,IAAIA,aAAAA,qBAAqB;GAAE,MAAM,EAAE;GAAM,WAAW,EAAE;EAAU,CAAC;CAClF,CAAC;CAED,cAA+D,gCAAgC;EAC7F,SAAQ,OAAM;GAAE,MAAM,EAAE;GAAQ,WAAW,EAAE;EAAU;EACvD,WAAU,MAAK,IAAIC,aAAAA,6BAA6B;GAAE,MAAM,EAAE;GAAM,WAAW,EAAE;EAAU,CAAC;CAC1F,CAAC;CAED,cAAiE,qBAAqB;EACpF,SAAQ,OAAM;GACZ,SAAS,EAAE;GACX,cAAc,EAAE;GAChB,OAAO,EAAE;GACT,UAAU,EAAE;GACZ,SAAS,EAAE;GACX,UAAU,EAAE;GACZ,kBAAkB,EAAE;GACpB,UAAU,EAAE;EACd;EACA,WAAU,MAAK,IAAIC,uBAAAA,kBAAkB,CAAC;CACxC,CAAC;CAED,OAAO;AACT,EAAA,CAAG;;;;;;;;;;;;;;;;AC9CH,MAAa,YAAY;;;;;;AA4BzB,SAAgB,WAAW,OAAkC;CAC3D,MAAM,MAAO,MAAkC;CAC/C,IAAI,OAAO,QAAQ,UAAU,OAAO;CACpC,QAAQ,KAAR;EACE,KAAK,aACH,OAAO;EACT,KAAK;EACL,KAAK;EACL,KAAK,OACH,OAAO,OAAQ,MAA0B,MAAM;EACjD,KAAK,UAAU;GACb,MAAM,IAAK,MAA0B;GACrC,IAAI,CAAC,KAAK,OAAO,EAAE,WAAW,YAAY,OAAO,EAAE,UAAU,UAAU,OAAO;GAK9E,IAAI,EAAE,OAAO,SAAA,MAAmC,OAAO;GAIvD,OAAO,gBAAgB,KAAK,EAAE,KAAK,KAAK,IAAI,IAAI,EAAE,KAAK,CAAC,CAAC,SAAS,EAAE,MAAM;EAC5E;EACA,KAAK;EACL,KAAK,OACH,OAAO,MAAM,QAAS,MAA0B,CAAC;EACnD,KAAK,SACH,OAAO,OAAQ,MAA0B,MAAM,YAAa,MAA0B,MAAM;EAC9F,KAAK,SACH,OAAO,OAAQ,MAA0B,MAAM;EACjD,SACE,OAAO;CACX;AACF;;;;;;;;;;;;ACvDA,SAAgB,OAAO,OAAyB;CAO9C,IAAI,CAAC,2BACH,MAAM,IAAI,MAAM,6CAA6C;CAE/D,OAAO,KAAK,uBAAO,IAAI,QAAQ,CAAC;AAClC;AAEA,SAAS,KAAK,GAAY,MAAgC;CACxD,IAAI,MAAM,KAAA,GAAW,OAAO,GAAG,YAAY,YAAY;CACvD,IAAI,MAAM,MAAM,OAAO;CAEvB,MAAM,IAAI,OAAO;CACjB,IAAI,MAAM,YAAY,MAAM,WAAW,OAAO;CAC9C,IAAI,MAAM,UAAU,OAAO,OAAO,SAAS,CAAW,IAAI,IAAI;CAC9D,IAAI,MAAM,UAAU,OAAO;GAAG,YAAY;EAAU,GAAI,EAAa,SAAS;CAAE;CAChF,IAAI,MAAM,cAAc,MAAM,UAAU,OAAO,KAAA;CAE/C,IAAI,aAAa,MACf,OAAO;GAAG,YAAY;EAAQ,GAAG,EAAE,YAAY;CAAE;CAEnD,IAAI,aAAa,QACf,OAAO;GAAG,YAAY;EAAU,GAAG;GAAE,QAAQ,EAAE;GAAQ,OAAO,EAAE;EAAM;CAAE;CAE1E,IAAI,aAAa,KACf,OAAO;GAAG,YAAY;EAAO,GAAG,EAAE,SAAS;CAAE;CAE/C,IAAI,aAAa,OACf,OAAO;GAAG,YAAY;EAAS,GAAG,eAAe,CAAC;CAAE;CAEtD,IAAI,aAAa,KAAK;EACpB,IAAI,KAAK,IAAI,CAAC,GAAG,OAAO;EACxB,KAAK,IAAI,CAAC;EACV,MAAM,UAAqC,CAAC;EAC5C,KAAK,MAAM,CAAC,GAAG,QAAQ,EAAE,QAAQ,GAC/B,QAAQ,KAAK,CAAC,KAAK,GAAG,IAAI,GAAG,KAAK,KAAK,IAAI,CAAC,CAAC;EAE/C,OAAO;IAAG,YAAY;GAAO,GAAG;EAAQ;CAC1C;CACA,IAAI,aAAa,KAAK;EACpB,IAAI,KAAK,IAAI,CAAC,GAAG,OAAO;EACxB,KAAK,IAAI,CAAC;EACV,MAAM,SAAoB,CAAC;EAC3B,KAAK,MAAM,KAAK,GAAG,OAAO,KAAK,KAAK,GAAG,IAAI,CAAC;EAC5C,OAAO;IAAG,YAAY;GAAO,GAAG;EAAO;CACzC;CAEA,IAAI,MAAM,QAAQ,CAAC,GAAG;EACpB,IAAI,KAAK,IAAI,CAAC,GAAG,OAAO;EACxB,KAAK,IAAI,CAAC;EACV,OAAO,EAAE,KAAI,MAAK,KAAK,GAAG,IAAI,CAAC;CACjC;CAEA,IAAI,MAAM,UAAU;EAClB,IAAI,KAAK,IAAI,CAAW,GAAG,OAAO;EAClC,KAAK,IAAI,CAAW;EAKpB,MAAM,cAAe,EAA2B;EAChD,IAAI,OAAO,gBAAgB,YACzB,OAAO,KAAM,YAA8B,KAAK,CAAC,GAAG,IAAI;EAM1D,MAAM,WADQ,EAAa,aACJ;EACvB,IAAI,YAAY,aAAa,UAAU;GACrC,MAAM,MAAM,cAAc,QAAQ;GAClC,IAAI,KACF,OAAO;KACJ,YAAY;IACb,GAAG;IACH,GAAG,KAAK,IAAI,OAAO,CAAC,GAAG,IAAI;GAC7B;EAEJ;EAEA,MAAM,MAA+B,CAAC;EACtC,KAAK,MAAM,KAAK,OAAO,KAAK,CAA4B,GAAG;GAEzD,IAAI,MAAM,aAAa;GACvB,MAAM,MAAO,EAA8B;GAC3C,IAAI,QAAQ,KAAA,GAAW;IAErB,IAAI,KAAK,GAAG,YAAY,YAAY;IACpC;GACF;GACA,MAAM,UAAU,KAAK,KAAK,IAAI;GAC9B,IAAI,YAAY,KAAA,GAAW;GAC3B,IAAI,KAAK;EACX;EACA,OAAO;CACT;CAEA,OAAO;AACT;;;;;;;;;;;;;;AAeA,SAAS,qBAAqB,GAAoB;CAChD,IAAI,CAAC,KAAK,OAAO,MAAM,UAAU,OAAO;CACxC,MAAM,YAAY;CAClB,IAAI,OAAO,UAAU,WAAW,UAAU,OAAO;CACjD,IAAI,OAAO,UAAU,UAAU,UAAU,OAAO;CAChD,IAAI,UAAU,OAAO,SAAA,MAAmC,OAAO;CAC/D,MAAM,QAAQ,UAAU;CACxB,IAAI,CAAC,gBAAgB,KAAK,KAAK,GAAG,OAAO;CACzC,IAAI,IAAI,IAAI,KAAK,CAAC,CAAC,SAAS,MAAM,QAAQ,OAAO;CACjD,IAAI;EAUF,OAAO,IAAI,OAAO,UAAU,QAAQ,KAAK;CAC3C,QAAQ;EACN,OAAO;CACT;AACF;;;;;;AAOA,SAAgB,OAAO,OAAyB;CAC9C,IAAI,UAAU,MAAM,OAAO;CAC3B,IAAI,OAAO,UAAU,UAAU,OAAO;CAEtC,IAAI,MAAM,QAAQ,KAAK,GAAG,OAAO,MAAM,IAAI,MAAM;CAEjD,IAAA,iBAAiB,SAAS,WAAW,KAAK,GAAG;EAC3C,MAAM,MAAM;EACZ,QAAQ,IAAI,YAAZ;GACE,KAAK,aACH;GACF,KAAK,QACH,OAAO,IAAI,KAAK,IAAI,CAAC;GACvB,KAAK,UACH,OAAO,OAAO,IAAI,CAAC;GACrB,KAAK,UACH,OAAO,qBAAqB,IAAI,CAAC;GACnC,KAAK,OACH,OAAO,IAAI,IAAI,IAAI,CAAC;GACtB,KAAK,OACH,OAAO,IAAI,IAAI,IAAI,EAAE,KAAK,CAAC,GAAG,SAAS,CAAC,OAAO,CAAC,GAAG,OAAO,GAAG,CAAC,CAAC,CAAC;GAClE,KAAK,OACH,OAAO,IAAI,IAAI,IAAI,EAAE,IAAI,MAAM,CAAC;GAClC,KAAK,SACH,OAAO,eAAe,sBAAsB,IAAI,CAAC,CAAC;GACpD,KAAK,SAAS;IACZ,MAAM,MAAM,cAAc,IAAI,CAAC;IAC/B,MAAM,OAAO,OAAO,IAAI,CAAC;IACzB,OAAO,MAAM,IAAI,SAAS,IAAI,IAAI;GACpC;EACF;CACF;CAEA,MAAM,MAA+B,CAAC;CACtC,KAAK,MAAM,KAAK,OAAO,KAAK,KAAgC,GAAG;EAC7D,IAAI,MAAM,aAAa;EACvB,IAAI,KAAK,OAAQ,MAAkC,EAAE;CACvD;CACA,OAAO;AACT;;;;;;AAOA,SAAS,sBAAsB,GAAqC;CAElE,OADgB,OAAO,CACV;AACf;;;ACtLA,MAAM,yCAAyC,KAAK,OAAO;;;;;;;;;;;;;AAc3D,MAAa,yBAAyB;AACtC,MAAM,sBAAsB;AAE5B,SAAS,eAAe,OAA0C;CAIhE,OAAO,GAAG,KAAK,UAAU,OAAO,KAAK,CAAC,EAAE;AAC1C;AAEA,SAAS,qBAAqB,QAAoB,iBAAwC;CACxF,OAAO,IAAI,SAAS,SAAS,WAAW;EACtC,IAAI,iBAAiB;EACrB,IAAI,iBAAiB;EACrB,IAAI,UAAU;EAEd,MAAM,gBAAgB;GACpB,OAAO,IAAI,SAAS,OAAO;GAC3B,OAAO,IAAI,SAAS,OAAO;GAC3B,OAAO,IAAI,SAAS,OAAO;EAC7B;EACA,MAAM,UAAU,UAAkB;GAChC,IAAI,SAAS;GACb,UAAU;GACV,QAAQ;GACR,IAAI,OAAO;IACT,OAAO,KAAK;IACZ;GACF;GACA,QAAQ;EACV;EACA,MAAM,qBAAqB;GACzB,IAAI,kBAAkB,gBACpB,OAAO;EAEX;EACA,MAAM,WAAW,UAAiB,OAAO,KAAK;EAG9C,MAAM,gBAAgB,uBAAO,IAAI,MAAM,uDAAuD,CAAC;EAC/F,MAAM,gBAAgB;GACpB,iBAAiB;GACjB,aAAa;EACf;EAEA,OAAO,KAAK,SAAS,OAAO;EAC5B,OAAO,KAAK,SAAS,OAAO;EAC5B,IAAI;EACJ,IAAI;GACF,UAAU,OAAO,MAAM,kBAAiB,UAAS;IAC/C,IAAI,OAAO;KACT,OAAO,KAAK;KACZ;IACF;IACA,iBAAiB;IACjB,aAAa;GACf,CAAC;EACH,SAAS,OAAO;GACd,OAAO,KAAc;GACrB;EACF;EACA,IAAI,CAAC,SAAS;GACZ,iBAAiB;GACjB,OAAO,KAAK,SAAS,OAAO;EAC9B;CACF,CAAC;AACH;AAEA,SAAS,WAAW,QAAoB,OAAiD;CACvF,OAAO,qBAAqB,QAAQ,eAAe,KAAK,CAAC;AAC3D;AAEA,SAAS,WAA0B;CACjC,OAAO,IAAI,SAAQ,YAAW,aAAa,OAAO,CAAC;AACrD;AAEA,SAAS,WAAW,QAAoB,SAA+B;CACrE,IAAI,SAAS;CACb,OAAO,YAAY,MAAM;CACzB,OAAO,GAAG,SAAQ,UAAS;EACzB,UAAU;EACV,OAAO,MAAM;GACX,MAAM,eAAe,OAAO,QAAQ,IAAI;GACxC,IAAI,iBAAiB,IAAI;GACzB,MAAM,OAAO,OAAO,MAAM,GAAG,YAAY;GACzC,SAAS,OAAO,MAAM,eAAe,CAAC;GACtC,IAAI,CAAC,KAAK,KAAK,GAAG;GAClB,IAAI;IACF,QAAQ,OAAO,KAAK,MAAM,IAAI,CAAC,CAAC;GAClC,QAAQ,CAER;EACF;CACF,CAAC;AACH;AAEA,IAAa,mBAAb,cAAsCC,sBAAAA,OAAO;CAC3C;CACA;CACA;CACA,YAAY;CACZ,UAAU;CACV;CACA,6BAAa,IAAI,IAAgC;CACjD,oCAAoB,IAAI,IAA+B;CACvD,iCAAiB,IAAI,IAA8B;CACnD,iCAAiB,IAAI,IAAmB;CACxC;CACA;CAEA,YAAY,YAAoB,UAAmC,CAAC,GAAG;EACrE,MAAM;EACN,KAAK,aAAa;EAClB,KAAKC,8BAA8B,QAAQ,8BAA8B;CAC3E;CAEA,IAAa,iBAAoD;EAC/D,OAAO,CAAC,MAAM;CAChB;CAEA,IAAI,WAAoB;EACtB,OAAO,KAAKC;CACd;;CAGA,IAAI,oBAA4B;EAC9B,OAAO,KAAKA,YAAY,KAAKC,eAAe,OAAO;CACrD;CAEA,MAAM,QACJ,OACA,OACA,SACe;EACf,MAAM,KAAKC,eAAe;EAS1B,IAAI,SAAS,WAAW;GACtB,MAAM,aAAoB;IACxB,GAAG;IACH,KAAA,GAAA,OAAA,WAAA,CAAe;IACf,2BAAW,IAAI,KAAK;IACpB,iBAAiB;GACnB;GACA,KAAKC,cAAc,OAAO,UAAU;GACpC;EACF;EAEA,IAAI,KAAKH,WAAW;GAClB,MAAM,KAAKI,mBAAmB,OAAO,OAAO,KAAA,GAAW,SAAS,SAAS;GACzE;EACF;EAEA,MAAM,SAAS,KAAKC;EACpB,IAAI,CAAC,UAAU,OAAO,WACpB,MAAM,KAAKH,eAAe,IAAI;EAEhC,MAAM,KAAKI,cAAc;GAAE,MAAM;GAAW;GAAO;GAAO,WAAW,SAAS;EAAU,CAAC;CAC3F;CAEA,MAAM,UAAU,OAAe,IAAmB,SAA2C;EAC3F,IAAI,SAAS,OACX,MAAM,IAAI,MAAM,6DAA6D;EAG/E,MAAM,YAAY,KAAKC,WAAW,IAAI,KAAK,qBAAK,IAAI,IAAmB;EACvE,MAAM,cAAc,UAAU,IAAI,EAAE;EACpC,MAAM,eAAe,QAAQ,KAAKF,iBAAiB,CAAC,KAAKA,cAAc,SAAS;EAChF,UAAU,IAAI,EAAE;EAChB,KAAKE,WAAW,IAAI,OAAO,SAAS;EAEpC,IAAI;GACF,MAAM,KAAKL,eAAe;GAC1B,IAAI,CAAC,KAAKF,aAAa,CAAC,eAAe,cACrC,MAAM,KAAKQ,uBAAuB,KAAK;EAE3C,SAAS,OAAO;GACd,IAAI,CAAC,aAAa;IAChB,UAAU,OAAO,EAAE;IACnB,IAAI,UAAU,SAAS,GACrB,KAAKD,WAAW,OAAO,KAAK;GAEhC;GACA,MAAM;EACR;CACF;CAEA,MAAM,YAAY,OAAe,IAAkC;EACjE,MAAM,YAAY,KAAKA,WAAW,IAAI,KAAK;EAC3C,WAAW,OAAO,EAAE;EACpB,IAAI,WAAW,SAAS,GAAG;GACzB,KAAKA,WAAW,OAAO,KAAK;GAC5B,IAAI,CAAC,KAAKP,aAAa,KAAKK,iBAAiB,CAAC,KAAKA,cAAc,WAAW;IAC1E,MAAM,KAAKC,cAAc;KAAE,MAAM;KAAe;IAAM,CAAC;IACvD,MAAM,SAAS;GACjB;EACF;CACF;CAEA,MAAM,QAAuB;EAC3B,MAAM,QAAQ,WAAW,CAAC,GAAG,KAAKG,cAAc,CAAC;CACnD;CAEA,MAAM,QAAuB;EAC3B,KAAKC,UAAU;EACf,KAAKH,WAAW,MAAM;EAEtB,KAAKF,eAAe,QAAQ;EAC5B,KAAKA,gBAAgB,KAAA;EACrB,KAAKM,wCAAwB,IAAI,MAAM,4BAA4B,CAAC;EAEpE,KAAK,MAAM,UAAU,CAAC,GAAG,KAAKV,eAAe,OAAO,CAAC,GACnD,KAAKW,oBAAoB,MAAM;EAGjC,IAAI,KAAKC,SAAS;GAChB,MAAM,IAAI,SAAc,YAAW,KAAKA,SAAS,YAAY,QAAQ,CAAC,CAAC;GACvE,KAAKA,UAAU,KAAA;EACjB;EAEA,IAAI,KAAKb,WACP,OAAA,GAAA,YAAA,OAAA,CAAa,KAAK,UAAU,CAAC,CAAC,YAAY,CAAC,CAAC;EAE9C,KAAKA,YAAY;CACnB;CAEA,MAAME,eAAe,iBAAiB,OAAsB;EAC1D,IAAI,KAAKQ,SACP,MAAM,IAAI,MAAM,4BAA4B;EAE9C,IAAI,CAAC,mBAAmB,KAAKV,aAAc,KAAKK,iBAAiB,CAAC,KAAKA,cAAc,YACnF;EAEF,IAAI,KAAKS,WACP,OAAO,KAAKA;EAGd,KAAKA,YAAY,KAAKC,OAAO,cAAc,CAAC,CAAC,cAAc;GACzD,KAAKD,YAAY,KAAA;EACnB,CAAC;EACD,OAAO,KAAKA;CACd;CAEA,MAAMC,OAAO,gBAAwC;EACnD,IAAI,gBAAgB;GAClB,KAAKV,eAAe,QAAQ;GAC5B,KAAKA,gBAAgB,KAAA;GACrB,KAAKL,YAAY;EACnB;EAEA,KAAKgB,eAAe;EACpB,OAAA,GAAA,YAAA,MAAA,EAAA,GAAA,KAAA,QAAA,CAAoB,KAAK,UAAU,GAAG,EAAE,WAAW,KAAK,CAAC;EACzD,KAAKA,eAAe;EAEpB,IAAI;GACF,MAAM,KAAKC,QAAQ;GACnB,KAAKD,eAAe;GACpB,KAAKhB,YAAY;GACjB;EACF,SAAS,OAAO;GACd,IAAI,KAAKU,SAAS;IAChB,MAAM,KAAK,MAAM;IACjB,MAAM,IAAI,MAAM,4BAA4B;GAC9C;GACA,MAAM,OAAQ,MAAgC;GAI9C,IAAI,SAAS,gBAAgB,SAAS,UAAU,MAAM;EACxD;EAEA,IAAI;GACF,MAAM,KAAKQ,eAAe;GAC1B,KAAKF,eAAe;EACtB,SAAS,OAAO;GACd,IAAI,KAAKN,SAAS;IAChB,MAAM,KAAK,MAAM;IACjB,MAAM,IAAI,MAAM,4BAA4B;GAC9C;GACA,MAAM,OAAQ,MAAgC;GAC9C,IAAI,SAAS,kBAAkB,SAAS,YAAY,SAAS,YAAY;IACvE,KAAKM,eAAe;IACpB,MAAM,KAAKG,aAAa;IACxB;GACF;GACA,MAAM;EACR;CACF;CAEA,iBAAiB;EACf,IAAI,KAAKT,SACP,MAAM,IAAI,MAAM,4BAA4B;CAEhD;CAEA,UAAyB;EACvB,OAAO,IAAI,SAAS,SAAS,WAAW;GACtC,MAAM,SAAS,IAAA,QAAI,cAAa,WAAU,KAAKU,oBAAoB,MAAM,CAAC;GAC1E,MAAM,WAAW,UAAiB;IAChC,OAAO,IAAI,aAAa,WAAW;IACnC,OAAO,KAAK;GACd;GACA,MAAM,oBAAoB;IACxB,OAAO,IAAI,SAAS,OAAO;IAC3B,KAAKP,UAAU;IACf,QAAQ;GACV;GAEA,OAAO,KAAK,SAAS,OAAO;GAC5B,OAAO,KAAK,aAAa,WAAW;GACpC,OAAO,OAAO,KAAK,UAAU;EAC/B,CAAC;CACH;CAEA,iBAAgC;EAC9B,OAAO,IAAI,SAAS,SAAS,WAAW;GACtC,MAAM,SAAS,IAAA,QAAI,iBAAiB,KAAK,UAAU;GACnD,MAAM,WAAW,UAAiB;IAChC,OAAO,IAAI,WAAW,SAAS;IAC/B,OAAO,KAAK;GACd;GACA,MAAM,kBAAkB;IACtB,OAAO,IAAI,SAAS,OAAO;IAC3B,KAAKR,gBAAgB;IACrB,KAAKL,YAAY;IACjB,WAAW,SAAQ,UAAS,KAAKqB,mBAAmB,KAAK,CAAC;IAG1D,OAAO,GAAG,eACR,KAAKC,wBAAwB,wBAAQ,IAAI,MAAM,2CAA2C,CAAC,CAC7F;IACA,OAAO,GAAG,UAAS,UAAS,KAAKA,wBAAwB,QAAQ,KAAK,CAAC;IACvE,KAAUC,mBAAmB,CAAC,CAAC,KAAK,SAAS,MAAM;GACrD;GAEA,OAAO,KAAK,SAAS,OAAO;GAC5B,OAAO,KAAK,WAAW,SAAS;EAClC,CAAC;CACH;CAEA,MAAMA,qBAAqB;EACzB,KAAK,MAAM,SAAS,KAAKhB,WAAW,KAAK,GACvC,MAAM,KAAKC,uBAAuB,KAAK;CAE3C;CAEA,wBAAwB,QAAoB,OAAc;EACxD,IAAI,KAAKH,kBAAkB,QAAQ;EACnC,KAAKA,gBAAgB,KAAA;EACrB,KAAKM,wBAAwB,KAAK;EAClC,IAAI,CAAC,KAAKD,SACR,KAAUc,yBAAyB;CAEvC;CAEA,MAAMA,2BAA0C;EAC9C,IAAI,KAAKC,aAAa,OAAO,KAAKA;EAClC,KAAKA,cAAc,KAAKC,6BAA6B,CAAC,CAAC,cAAc;GACnE,KAAKD,cAAc,KAAA;EACrB,CAAC;EACD,OAAO,KAAKA;CACd;CAEA,MAAMC,+BAA8C;EAClD,OAAO,CAAC,KAAKhB,WAAW,CAAC,KAAKV,aAAa,EAAE,KAAKK,iBAAiB,CAAC,KAAKA,cAAc,YACrF,IAAI;GACF,MAAM,KAAKH,eAAe,IAAI;GAC9B;EACF,QAAQ;GACN,IAAI,KAAKQ,SAAS;GAClB,MAAM,IAAI,SAAQ,YAAW,WAAW,SAAS,EAAE,CAAC;EACtD;CAEJ;;;;;;CAOA,MAAMS,eAA8B;EAClC,MAAM,WAAW,KAAK,aAAa;EACnC,IAAI;EACJ,IAAI;GACF,SAAS,OAAA,GAAA,YAAA,KAAA,CAAW,UAAU,IAAI;EACpC,SAAS,GAAG;GACV,IAAK,EAA4B,SAAS,UAAU;IAClD,IAAI,MAAM,KAAKQ,qBAAqB,QAAQ,GAAG;KAC7C,OAAA,GAAA,YAAA,OAAA,CAAa,QAAQ,CAAC,CAAC,YAAY,CAAC,CAAC;KACrC,MAAM,IAAI,MAAM,oCAAoC;IACtD;IACA,MAAM,IAAI,SAAQ,YAAW,WAAW,SAAS,GAAG,CAAC;IACrD,IAAI;KACF,MAAM,KAAKT,eAAe;KAC1B,KAAKF,eAAe;KACpB;IACF,QAAQ;KACN,MAAM,IAAI,MAAM,gDAAgD;IAClE;GACF;GACA,MAAM;EACR;EAEA,IAAI;GAGF,IAAI;IACF,MAAM,KAAKE,eAAe;IAC1B,KAAKF,eAAe;IACpB;GACF,QAAQ,CAER;GACA,OAAA,GAAA,YAAA,OAAA,CAAa,KAAK,UAAU,CAAC,CAAC,YAAY,CAAC,CAAC;GAC5C,KAAKA,eAAe;GACpB,MAAM,KAAKC,QAAQ;GACnB,KAAKD,eAAe;GACpB,KAAKhB,YAAY;EACnB,UAAU;GACR,MAAM,OAAO,MAAM,CAAC,CAAC,YAAY,CAAC,CAAC;GACnC,OAAA,GAAA,YAAA,OAAA,CAAa,QAAQ,CAAC,CAAC,YAAY,CAAC,CAAC;EACvC;CACF;CAEA,MAAM2B,qBAAqB,UAAoC;EAC7D,IAAI;GACF,MAAM,WAAW,OAAA,GAAA,YAAA,KAAA,CAAW,QAAQ;GACpC,OAAO,KAAK,IAAI,IAAI,SAAS,UAAU;EACzC,QAAQ;GACN,OAAO;EACT;CACF;CAEA,MAAMnB,uBAAuB,OAA8B;EACzD,IAAI;EACJ,MAAM,aAAa,IAAI,SAAe,SAAS,WAAW;GACxD,SAAS;IAAE;IAAS;GAAO;GAC3B,MAAM,UAAU,KAAKoB,kBAAkB,IAAI,KAAK,KAAK,CAAC;GACtD,QAAQ,KAAK,MAAM;GACnB,KAAKA,kBAAkB,IAAI,OAAO,OAAO;EAC3C,CAAC;EACD,IAAI;GACF,MAAM,KAAKtB,cAAc;IAAE,MAAM;IAAa;GAAM,CAAC;EACvD,SAAS,OAAO;GACd,KAAKuB,uBAAuB,OAAO,MAAM;GACzC,MAAM;EACR;EACA,MAAM;CACR;CAEA,uBAAuB,OAAe,QAAqC;EACzE,IAAI,CAAC,QAAQ;EACb,MAAM,UAAU,KAAKD,kBAAkB,IAAI,KAAK;EAChD,IAAI,CAAC,SAAS;EACd,MAAM,cAAc,QAAQ,QAAO,SAAQ,SAAS,MAAM;EAC1D,IAAI,YAAY,WAAW,GAAG;GAC5B,KAAKA,kBAAkB,OAAO,KAAK;GACnC;EACF;EACA,KAAKA,kBAAkB,IAAI,OAAO,WAAW;CAC/C;CAEA,wBAAwB,OAAe,OAAe;EACpD,MAAM,UAAU,KAAKA,kBAAkB,IAAI,KAAK;EAChD,KAAKA,kBAAkB,OAAO,KAAK;EACnC,IAAI,OAAO;GACT,SAAS,SAAQ,WAAU,OAAO,OAAO,KAAK,CAAC;GAC/C;EACF;EACA,SAAS,SAAQ,WAAU,OAAO,QAAQ,CAAC;CAC7C;CAEA,wBAAwB,OAAc;EACpC,KAAK,MAAM,SAAS,KAAKA,kBAAkB,KAAK,GAC9C,KAAKE,wBAAwB,OAAO,KAAK;CAE7C;CAEA,oBAAoB,QAAoB;EACtC,MAAM,SAAuB;GAC3B;GACA,+BAAe,IAAI,IAAI;GACvB,YAAY,QAAQ,QAAQ;GAC5B,aAAa;EACf;EACA,KAAK7B,eAAe,IAAI,QAAQ,MAAM;EACtC,WAAW,SAAQ,UAAS;GAC1B,MAAM,cAAc;GACpB,IAAI,YAAY,SAAS,aAAa;IACpC,OAAO,cAAc,IAAI,YAAY,KAAK;IAC1C,KAAK8B,0BAA0B,QAAQ;KAAE,MAAM;KAAc,OAAO,YAAY;IAAM,CAAC;GACzF,OAAO,IAAI,YAAY,SAAS,eAC9B,OAAO,cAAc,OAAO,YAAY,KAAK;QACxC,IAAI,YAAY,SAAS,WAC9B,KAAU3B,mBAAmB,YAAY,OAAO,YAAY,OAAO,QAAQ,YAAY,SAAS;EAEpG,CAAC;EACD,OAAO,GAAG,eAAe,KAAKQ,oBAAoB,MAAM,CAAC;EACzD,OAAO,GAAG,eAAe,KAAKA,oBAAoB,MAAM,CAAC;CAC3D;CAEA,0BAA0B,QAAsB,OAAoB;EAClE,IAAI,KAAKX,eAAe,IAAI,OAAO,MAAM,MAAM,UAAU,OAAO,OAAO,WAAW;EAElF,MAAM,kBAAkB,eAAe,KAAK;EAC5C,MAAM,cAAc,OAAO,WAAW,eAAe;EACrD,IAAI,OAAO,cAAc,cAAc,KAAKF,6BAA6B;GACvE,KAAKa,oBAAoB,MAAM;GAC/B;EACF;EAEA,OAAO,eAAe;EAEtB,MAAM,QAAQ,OAAO,WAClB,YAAY,CAAC,CAAC,CAAC,CACf,KAAK,YAAY;GAChB,IAAI,KAAKX,eAAe,IAAI,OAAO,MAAM,MAAM,UAAU,OAAO,OAAO,WAAW;GAClF,MAAM,qBAAqB,OAAO,QAAQ,eAAe;EAC3D,CAAC,CAAC,CACD,YAAY;GACX,KAAKW,oBAAoB,MAAM;EACjC,CAAC,CAAC,CACD,cAAc;GACb,OAAO,cAAc,KAAK,IAAI,GAAG,OAAO,cAAc,WAAW;EACnE,CAAC;EAEH,OAAO,aAAa;EACpB,KAAKH,eAAe,IAAI,KAAK;EAC7B,MAAW,cAAc,KAAKA,eAAe,OAAO,KAAK,CAAC;CAC5D;CAEA,oBAAoB,QAAsB;EACxC,IAAI,KAAKR,eAAe,IAAI,OAAO,MAAM,MAAM,QAAQ;EACvD,KAAKA,eAAe,OAAO,OAAO,MAAM;EACxC,OAAO,cAAc,MAAM;EAC3B,OAAO,cAAc;EACrB,OAAO,aAAa,QAAQ,QAAQ;EACpC,IAAI,CAAC,OAAO,OAAO,WACjB,OAAO,OAAO,QAAQ;CAE1B;CAEA,mBAAmB,OAAoB;EACrC,IAAI,MAAM,SAAS,cAAc;GAC/B,KAAK6B,wBAAwB,MAAM,KAAK;GACxC;EACF;EACA,IAAI,MAAM,SAAS,SAAS;EAG5B,KAAK3B,cAAc,MAAM,OAAO,MAAM,KAAK;CAC7C;CAEA,MAAMC,mBACJ,OACA,OACA,cACA,WACA;EACA,MAAM,cAAqB;GACzB,GAAG;GACH,KAAA,GAAA,OAAA,WAAA,CAAe;GACf,2BAAW,IAAI,KAAK;GACpB,iBAAiB;EACnB;EAEA,KAAKD,cAAc,OAAO,WAAW;EAGrC,IAAI,KAAKF,eAAe,SAAS,GAAG;EAQpC,IAAI,WAAW;GACb,IAAI,gBAAgB,aAAa,cAAc,IAAI,KAAK,KAAK,CAAC,aAAa,OAAO,WAChF,KAAK8B,0BAA0B,cAAc;IAAE,MAAM;IAAS;IAAO,OAAO;GAAY,CAAC;GAE3F;EACF;EAEA,IAAI;EACJ,KAAK,MAAM,UAAU,KAAK9B,eAAe,OAAO,GAAG;GACjD,IAAI,CAAC,OAAO,cAAc,IAAI,KAAK,KAAK,OAAO,OAAO,WAAW;GAEjE,UAAU;IAAE,MAAM;IAAS;IAAO,OAAO;GAAY;GACrD,KAAK8B,0BAA0B,QAAQ,KAAK;EAC9C;CACF;CAEA,cAAc,OAAe,OAAc;EACzC,MAAM,YAAY,KAAKxB,WAAW,IAAI,KAAK;EAC3C,IAAI,CAAC,WAAW;EAChB,KAAK,MAAM,MAAM,WACf,KAAKyB,qBAAqB,OAAO,OAAO,IAAI,CAAC;CAEjD;CAEA,qBAAqB,OAAe,OAAc,IAAmB,SAAiB;EACpF,IAAI,SAAS;EACb,MAAM,OAAO,YAAY;GACvB,IAAI,UAAU,KAAKtB,SAAS;GAC5B,SAAS;GACT,IAAI,WAAA,GAAmC;GAEvC,IAAI,CADoB,KAAKH,WAAW,IAAI,KAAK,CAAC,EAAE,IAAI,EAAE,GACpC;GAkBtB,iBAhBQ;IACJ,IAAI,KAAKG,SAAS;IAClB,IAAI,CAAC,KAAKH,WAAW,IAAI,KAAK,CAAC,EAAE,IAAI,EAAE,GAAG;IAC1C,MAAM,mBAA0B;KAC9B,GAAG;KACH,kBAAkB,MAAM,mBAAmB,KAAK;IAClD;IACA,KAAKyB,qBAAqB,OAAO,kBAAkB,IAAI,UAAU,CAAC;GACpE,GACA,uBAAuB,UAAU,EAO/B,CAAC,CAAC,QAAQ;EAChB;EACA,IAAI;GACF,MAAM,SAAU,GACd,OACA,YAAY,CAAC,GACb,IACF;GACA,IAAI,UAAU,OAAQ,OAAyB,UAAU,YACvD,OAA+B,YAAY,CAAC,CAAC;EAEjD,QAAQ,CAER;CACF;CAEA,MAAM1B,cAAc,OAAoB;EAKtC,MAAM,aAAa;EACnB,IAAI;EACJ,KAAK,IAAI,UAAU,GAAG,WAAW,YAAY,WAC3C,IAAI;GACF,IAAI,YAAY,GACd,MAAM,KAAK2B,oBAAoB,KAAK;QAC/B;IACL,IAAI,KAAKvB,SAAS,MAAM;IACxB,MAAM,eAAe,KAAKL;IAC1B,KAAKA,gBAAgB,KAAA;IACrB,cAAc,QAAQ;IACtB,MAAM,KAAKH,eAAe,IAAI;IAC9B,MAAM,KAAK+B,oBAAoB,KAAK;GACtC;GACA;EACF,SAAS,OAAO;GACd,YAAY;GACZ,IAAI,KAAKvB,SAAS,MAAM;GACxB,MAAM,OAAQ,OAAiC;GAkB/C,IAAI,EANF,SAAS,WACT,SAAS,gBACT,SAAS,cACR,OAAiB,SAAS,SAAS,sCAAsC,KACzE,OAAiB,SAAS,SAAS,0BAA0B,KAC7D,OAAiB,SAAS,SAAS,2BAA2B,MAC/C,YAAY,YAAY,MAAM;GAEhD,MAAM,IAAI,SAAQ,YAAW,WAAW,SAAS,MAAM,UAAU,EAAE,CAAC;EACtE;CAEJ;CAEA,MAAMuB,oBAAoB,OAAoB;EAC5C,MAAM,SAAS,KAAK5B;EACpB,IAAI,CAAC,UAAU,OAAO,WACpB,MAAM,KAAKH,eAAe,IAAI;EAEhC,IAAI,KAAKF,WAAW;GAClB,MAAM,KAAKkC,2BAA2B,KAAK;GAC3C;EACF;EACA,MAAM,eAAe,KAAK7B;EAC1B,IAAI,CAAC,gBAAgB,aAAa,WAGhC,MAAM,IAAI,MAAM,+CAA+C;EAEjE,MAAM,WAAW,cAAc,KAAK;CACtC;CAEA,MAAM6B,2BAA2B,OAAoB;EACnD,IAAI,MAAM,SAAS,aACjB,KAAKJ,wBAAwB,MAAM,KAAK;OACnC,IAAI,MAAM,SAAS,WACxB,MAAM,KAAK1B,mBAAmB,MAAM,OAAO,MAAM,KAAK;CAE1D;AACF"}