{"version":3,"file":"intelligence.cjs","names":["AgentRunner","Socket","Observable","AG_UI_CHANNEL_EVENT","EventType"],"sources":["../../../../src/v2/runtime/runner/intelligence.ts"],"sourcesContent":["import type {\n  AgentRunnerConnectRequest,\n  AgentRunnerIsRunningRequest,\n  AgentRunnerRunRequest,\n} from \"./agent-runner\";\nimport { AgentRunner } from \"./agent-runner\";\nimport type { AgentRunnerStopRequest } from \"./agent-runner\";\nimport { Observable } from \"rxjs\";\nimport type { AbstractAgent, BaseEvent, RunStartedEvent } from \"@ag-ui/client\";\nimport { EventType } from \"@ag-ui/client\";\nimport {\n  finalizeRunEvents,\n  AG_UI_CHANNEL_EVENT,\n  phoenixExponentialBackoff,\n} from \"@copilotkit/shared\";\nimport type { Channel } from \"phoenix\";\nimport { Socket } from \"phoenix\";\nimport { randomUUID } from \"node:crypto\";\n\nexport interface IntelligenceAgentRunnerOptions {\n  /** Phoenix runner websocket URL, e.g. \"ws://localhost:4000/runner\" */\n  url: string;\n  /** Optional Phoenix socket auth token used during websocket connect. */\n  authToken?: string;\n  /** Max delay (ms) for WebSocket reconnect backoff. @default 10_000 */\n  maxReconnectMs?: number;\n  /** Max delay (ms) for channel rejoin backoff. @default 30_000 */\n  maxRejoinMs?: number;\n}\n\nexport interface RunnerStartupBoundary {\n  events: Observable<BaseEvent>;\n  startup: Promise<void>;\n}\n\ninterface ThreadState {\n  threadId: string;\n  runId: string;\n  socket: Socket;\n  channel: Channel;\n  isRunning: boolean;\n  stopRequested: boolean;\n  agent: AbstractAgent | null;\n  currentEvents: BaseEvent[];\n  nextEventSeq: number;\n  hasRunStarted: boolean;\n  hasJoined: boolean;\n  supportsRunnerEventBatch: boolean;\n  producerFinished: boolean;\n  pendingEvents: Map<\n    string,\n    { payload: Record<string, unknown>; queuedAt: number }\n  >;\n  activeEventBatch: { eventIds: string[]; attempt: number } | null;\n  nextEventPushAttempt: number;\n  eventRetryTimer: ReturnType<typeof setTimeout> | null;\n  eventFlushTimer: ReturnType<typeof setTimeout> | null;\n  eventDeadlineTimer: ReturnType<typeof setTimeout> | null;\n  socketReconnectWatchdog: ReturnType<typeof setTimeout> | null;\n  eventRetryAttempt: number;\n  completeRun: () => void;\n  failRun: (error: Error) => void;\n}\n\nconst MAX_CONSECUTIVE_SOCKET_ERRORS = 5;\nconst EVENT_RETRY_BASE_MS = 100;\nconst EVENT_RETRY_MAX_MS = 2_000;\nconst RUNNER_EVENT_BATCH_CAPABILITY = \"runner_event_batch_v1\";\nconst MAX_RUNNER_EVENT_BATCH_SIZE = 32;\nconst RUNNER_EVENT_BATCH_FLUSH_MS = 5;\nconst EVENT_DURABILITY_DEADLINE_MS = 60_000;\n\nexport class IntelligenceAgentRunner extends AgentRunner {\n  private options: IntelligenceAgentRunnerOptions;\n  private threads = new Map<string, ThreadState>();\n\n  constructor(options: IntelligenceAgentRunnerOptions) {\n    super();\n    // Store config — sockets are created per-run, not eagerly.\n    this.options = options;\n  }\n\n  /**\n   * Create a new Phoenix socket with explicit exponential backoff.\n   *\n   * Each run/connect gets its own socket so that:\n   *  - A socket failure only affects a single thread, not all threads.\n   *  - Cleanup is simple: channel.leave() + socket.disconnect() tears\n   *    down everything for that run with no shared-state concerns.\n   *  - Each run gets its own independent retry budget.\n   *\n   * reconnectAfterMs — delay before Phoenix reconnects the WebSocket\n   *   after an unclean close. 100ms base, doubling up to maxReconnectMs (default 10s).\n   *\n   * rejoinAfterMs — delay before Phoenix re-joins a channel that\n   *   entered the \"errored\" state. 1s base, doubling up to maxRejoinMs (default 30s).\n   *\n   * These are set explicitly because Phoenix's default schedule is a\n   * fixed stepped array (not exponential), and any code that calls\n   * socket.disconnect() in an onError handler will set\n   * closeWasClean = true and reset the reconnect timer — permanently\n   * killing retries.\n   */\n  private createSocket(authToken = this.options.authToken): Socket {\n    const socket = new Socket(this.options.url, {\n      ...(authToken ? { authToken } : {}),\n      reconnectAfterMs: phoenixExponentialBackoff(\n        100,\n        this.options.maxReconnectMs ?? 10_000,\n      ),\n      rejoinAfterMs: phoenixExponentialBackoff(\n        1_000,\n        this.options.maxRejoinMs ?? 30_000,\n      ),\n    });\n    socket.connect();\n    return socket;\n  }\n\n  private createRunnerEventPayload(\n    event: BaseEvent,\n    request: AgentRunnerRunRequest,\n    state: ThreadState,\n  ): Record<string, unknown> {\n    const canonicalEvent = this.stampRunnerMetadata(\n      this.stampCanonicalRunOwnership(event, request),\n      state,\n    );\n    const payload = {\n      ...(canonicalEvent as Record<string, unknown>),\n    };\n\n    payload.threadId = request.threadId;\n    payload.runId = request.input.runId;\n    payload.thread_id = request.threadId;\n    payload.run_id = request.input.runId;\n\n    return payload;\n  }\n\n  private stampCanonicalRunOwnership(\n    event: BaseEvent,\n    request: AgentRunnerRunRequest,\n  ): BaseEvent {\n    const eventRecord = event as BaseEvent & Record<string, unknown>;\n    eventRecord.threadId = request.threadId;\n    eventRecord.runId = request.input.runId;\n    return eventRecord;\n  }\n\n  private stampRunnerMetadata(event: BaseEvent, state: ThreadState): BaseEvent {\n    const eventRecord = event as BaseEvent & {\n      metadata?: Record<string, unknown>;\n    };\n\n    const existingMetadata = eventRecord.metadata ?? {};\n    const hasEventId = typeof existingMetadata.cpki_event_id === \"string\";\n    const hasEventSeq = typeof existingMetadata.cpki_event_seq === \"number\";\n\n    if (hasEventId && hasEventSeq) {\n      const eventSeq = existingMetadata.cpki_event_seq as number;\n      state.nextEventSeq = Math.max(state.nextEventSeq, eventSeq + 1);\n      return eventRecord;\n    }\n\n    const eventSeq = state.nextEventSeq++;\n\n    eventRecord.metadata = {\n      ...existingMetadata,\n      cpki_event_id:\n        typeof existingMetadata.cpki_event_id === \"string\"\n          ? existingMetadata.cpki_event_id\n          : randomUUID(),\n      cpki_event_seq: eventSeq,\n    };\n    return eventRecord;\n  }\n\n  run(request: AgentRunnerRunRequest): Observable<BaseEvent> {\n    return this.createRunObservable(request);\n  }\n\n  runWithStartupBoundary(\n    request: AgentRunnerRunRequest,\n  ): RunnerStartupBoundary {\n    let resolveStartup: (() => void) | undefined;\n    let rejectStartup: ((reason: Error) => void) | undefined;\n    const startup = new Promise<void>((resolve, reject) => {\n      resolveStartup = resolve;\n      rejectStartup = reject;\n    });\n\n    return {\n      events: this.createRunObservable(request, {\n        resolveStartup: () => resolveStartup?.(),\n        rejectStartup: (error) => rejectStartup?.(error),\n      }),\n      startup,\n    };\n  }\n\n  private createRunObservable(\n    request: AgentRunnerRunRequest,\n    startupBoundary?: {\n      resolveStartup: () => void;\n      rejectStartup: (error: Error) => void;\n    },\n  ): Observable<BaseEvent> {\n    const { threadId, agent, input } = request;\n\n    const existing = this.threads.get(threadId);\n    if (existing?.isRunning) {\n      throw new Error(\"Thread already running\");\n    }\n\n    return new Observable((observer) => {\n      if (this.threads.get(threadId)?.isRunning) {\n        observer.error(new Error(\"Thread already running\"));\n        return;\n      }\n\n      const socket = this.createSocket(request.authToken);\n\n      const channel = socket.channel(`ingestion:${input.runId}`, {\n        thread_id: threadId,\n        run_id: input.runId,\n      });\n\n      const state: ThreadState = {\n        threadId,\n        runId: input.runId,\n        socket,\n        channel,\n        isRunning: true,\n        stopRequested: false,\n        agent,\n        currentEvents: [],\n        nextEventSeq: 1,\n        hasRunStarted: false,\n        hasJoined: false,\n        supportsRunnerEventBatch: false,\n        producerFinished: false,\n        pendingEvents: new Map(),\n        activeEventBatch: null,\n        nextEventPushAttempt: 0,\n        eventRetryTimer: null,\n        eventFlushTimer: null,\n        eventDeadlineTimer: null,\n        socketReconnectWatchdog: null,\n        eventRetryAttempt: 0,\n        completeRun: () => observer.complete(),\n        failRun: (error) => observer.error(error),\n      };\n      this.threads.set(threadId, state);\n\n      let consecutiveSocketErrors = 0;\n      let plannedRestart = false;\n\n      socket.onClose((event) => {\n        if (!this.isCurrentThreadState(threadId, state)) {\n          return;\n        }\n        plannedRestart = plannedRestart || event?.code === 1012;\n        if (event?.code !== 1000 && state.socketReconnectWatchdog === null) {\n          state.socketReconnectWatchdog = setTimeout(() => {\n            state.socketReconnectWatchdog = null;\n            if (!state.isRunning || socket.isConnected()) {\n              return;\n            }\n            socket.disconnect(() => {\n              if (state.isRunning && !socket.isConnected()) {\n                socket.connect();\n              }\n            });\n          }, 1_000);\n        }\n      });\n      socket.onOpen(() => {\n        if (!this.isCurrentThreadState(threadId, state)) {\n          return;\n        }\n        if (state.socketReconnectWatchdog !== null) {\n          clearTimeout(state.socketReconnectWatchdog);\n          state.socketReconnectWatchdog = null;\n        }\n        consecutiveSocketErrors = 0;\n        plannedRestart = false;\n      });\n      socket.onError(() => {\n        if (!this.isCurrentThreadState(threadId, state)) {\n          return;\n        }\n        // Once the gateway has accepted the run, transport recovery is not a\n        // terminal condition. The agent may still be producing the\n        // authoritative answer while Phoenix reconnects and replays durable\n        // events, so an arbitrary socket-error count must not abort it.\n        if (plannedRestart || state.hasJoined) {\n          return;\n        }\n\n        consecutiveSocketErrors += 1;\n        if (consecutiveSocketErrors >= MAX_CONSECUTIVE_SOCKET_ERRORS) {\n          state.agent?.abortRun();\n        }\n      });\n\n      // Listen for custom \"stop\" events pushed by the client over the\n      // channel. This must be registered before channel.join() so the\n      // handler is in place by the time the server starts relaying messages.\n      // The client sends the stop event before leaving the channel, so the\n      // runner is guaranteed to receive it while still joined.\n      channel.on(AG_UI_CHANNEL_EVENT, (payload: BaseEvent) => {\n        if (\n          this.isCurrentThreadState(threadId, state) &&\n          payload.type === EventType.CUSTOM &&\n          (payload as BaseEvent & { name?: string }).name === \"stop\"\n        ) {\n          this.stop({ threadId, runId: state.runId });\n        }\n      });\n\n      channel\n        .join()\n        .receive(\"ok\", (response) => {\n          if (!this.isCurrentThreadState(threadId, state)) {\n            return;\n          }\n          const supportsRunnerEventBatch =\n            this.supportsRunnerEventBatch(response);\n          const activeEventBatch =\n            state.supportsRunnerEventBatch === supportsRunnerEventBatch\n              ? state.activeEventBatch\n              : null;\n          state.supportsRunnerEventBatch = supportsRunnerEventBatch;\n          if (state.hasJoined) {\n            this.resetPendingEventRetry(state);\n            if (activeEventBatch !== null) {\n              state.activeEventBatch = activeEventBatch;\n              this.retryActiveEventBatch(state);\n            } else {\n              this.replayPendingEvents(state);\n            }\n            return;\n          }\n\n          state.hasJoined = true;\n          startupBoundary?.resolveStartup();\n          void this.executeAgentRun(request, state, threadId, (event) => {\n            observer.next(event);\n          });\n        })\n        .receive(\"error\", (resp) => {\n          if (!this.isCurrentThreadState(threadId, state)) {\n            return;\n          }\n          if (state.hasJoined) {\n            if (this.isPermanentEventFailure(resp)) {\n              const reason =\n                typeof (resp as { reason?: unknown }).reason === \"string\"\n                  ? `: ${(resp as { reason: string }).reason}`\n                  : \"\";\n              this.failThread(\n                threadId,\n                state,\n                new Error(\n                  `Gateway permanently rejected channel rejoin${reason}`,\n                ),\n              );\n            }\n            return;\n          }\n          if (this.isRetryableJoinError(resp)) {\n            return;\n          }\n\n          const error = new Error(\n            `Failed to join channel: ${JSON.stringify(resp)}`,\n          );\n          const errorEvent = {\n            type: EventType.RUN_ERROR,\n            message: error.message,\n            code: \"CHANNEL_JOIN_ERROR\",\n          } as BaseEvent;\n          observer.next(errorEvent);\n          state.currentEvents.push(errorEvent);\n          this.removeThread(threadId, state);\n          startupBoundary?.rejectStartup(error);\n          observer.complete();\n        })\n        .receive(\"timeout\", () => {\n          if (!this.isCurrentThreadState(threadId, state)) {\n            return;\n          }\n          if (state.hasJoined) {\n            return;\n          }\n\n          const error = new Error(\"Timed out joining channel\");\n          const errorEvent = {\n            type: EventType.RUN_ERROR,\n            message: error.message,\n            code: \"CHANNEL_JOIN_TIMEOUT\",\n          } as BaseEvent;\n          observer.next(errorEvent);\n          state.currentEvents.push(errorEvent);\n          this.removeThread(threadId, state);\n          startupBoundary?.rejectStartup(error);\n          observer.complete();\n        });\n\n      return () => {\n        this.removeThread(threadId, state);\n      };\n    });\n  }\n\n  connect(request: AgentRunnerConnectRequest): Observable<BaseEvent> {\n    const { threadId } = request;\n\n    return new Observable((observer) => {\n      const socket = this.createSocket();\n\n      const channel = socket.channel(`thread:${threadId}`);\n\n      channel.on(\"ag_ui_event\", (payload: BaseEvent) => {\n        observer.next(payload);\n\n        if (\n          payload.type === EventType.RUN_FINISHED ||\n          payload.type === EventType.RUN_ERROR\n        ) {\n          observer.complete();\n        }\n      });\n\n      const cleanup = () => {\n        channel.leave();\n        socket.disconnect();\n      };\n\n      channel\n        .join()\n        .receive(\"ok\", () => undefined)\n        .receive(\"error\", (resp) => {\n          observer.error(\n            new Error(`Failed to join channel: ${JSON.stringify(resp)}`),\n          );\n          cleanup();\n        })\n        .receive(\"timeout\", () => {\n          observer.error(new Error(\"Timed out joining channel\"));\n          cleanup();\n        });\n\n      return () => {\n        cleanup();\n      };\n    });\n  }\n\n  isRunning(request: AgentRunnerIsRunningRequest): Promise<boolean> {\n    const state = this.threads.get(request.threadId);\n    return Promise.resolve(state?.isRunning ?? false);\n  }\n\n  stop(request: AgentRunnerStopRequest): Promise<boolean | undefined> {\n    const state = this.threads.get(request.threadId);\n    if (!state || !state.isRunning || state.stopRequested) {\n      return Promise.resolve(false);\n    }\n    if (request.runId !== undefined && state.runId !== request.runId) {\n      return Promise.resolve(false);\n    }\n\n    state.stopRequested = true;\n\n    // Direct local abort — the runtime is the authority.\n    if (state.agent) {\n      try {\n        state.agent.abortRun();\n      } catch {\n        // Ignore abort errors.\n      }\n    }\n\n    return Promise.resolve(true);\n  }\n\n  private async executeAgentRun(\n    request: AgentRunnerRunRequest,\n    state: ThreadState,\n    threadId: string,\n    onRunError: (event: BaseEvent) => void,\n  ): Promise<void> {\n    const { currentEvents } = state;\n    const pushCanonicalEvent = (event: BaseEvent): void => {\n      if (!this.isCurrentThreadState(threadId, state)) {\n        return;\n      }\n      const canonicalEvent = this.stampRunnerMetadata(\n        this.stampCanonicalRunOwnership(event, request),\n        state,\n      );\n      currentEvents.push(canonicalEvent);\n\n      if (canonicalEvent.type === EventType.RUN_STARTED) {\n        state.hasRunStarted = true;\n      }\n\n      this.queueRunnerEvent(\n        this.createRunnerEventPayload(canonicalEvent, request, state),\n        state,\n      );\n    };\n\n    const getPersistedInputMessages = () =>\n      request.persistedInputMessages ?? request.input.messages;\n\n    const buildRunStartedEvent = (\n      source?: RunStartedEvent,\n    ): RunStartedEvent => {\n      const baseInput = source?.input ?? request.input;\n      const persistedInputMessages = getPersistedInputMessages();\n      const event =\n        source ??\n        ({\n          type: EventType.RUN_STARTED,\n          threadId: request.threadId,\n          runId: request.input.runId,\n        } as RunStartedEvent);\n      event.threadId = request.threadId;\n      event.runId = request.input.runId;\n      event.input = {\n        ...baseInput,\n        threadId: request.threadId,\n        runId: request.input.runId,\n        ...(persistedInputMessages !== undefined\n          ? { messages: persistedInputMessages }\n          : {}),\n      };\n      return event;\n    };\n\n    const ensureRunStarted = (): void => {\n      if (!state.hasRunStarted) {\n        state.hasRunStarted = true;\n        pushCanonicalEvent(buildRunStartedEvent());\n      }\n    };\n\n    try {\n      await request.agent.runAgent(request.input, {\n        onEvent: ({ event }: { event: BaseEvent }) => {\n          if (event.type === EventType.RUN_STARTED) {\n            pushCanonicalEvent(buildRunStartedEvent(event as RunStartedEvent));\n            return;\n          }\n\n          ensureRunStarted();\n          pushCanonicalEvent(event);\n        },\n      });\n    } catch (error) {\n      if (!this.isCurrentThreadState(threadId, state)) {\n        return;\n      }\n      ensureRunStarted();\n      const existingError = currentEvents.find(\n        (event) => event.type === EventType.RUN_ERROR,\n      );\n      if (existingError) {\n        onRunError(existingError);\n      } else {\n        const errorEvent = {\n          type: EventType.RUN_ERROR,\n          message: error instanceof Error ? error.message : String(error),\n        } as BaseEvent;\n        pushCanonicalEvent(errorEvent);\n        onRunError(errorEvent);\n      }\n    } finally {\n      if (!this.isCurrentThreadState(threadId, state)) {\n        return;\n      }\n      ensureRunStarted();\n      const appended = finalizeRunEvents(currentEvents, {\n        stopRequested: state.stopRequested,\n      });\n      for (const event of appended) {\n        const canonicalEvent = this.stampRunnerMetadata(\n          this.stampCanonicalRunOwnership(event, request),\n          state,\n        );\n        this.queueRunnerEvent(\n          this.createRunnerEventPayload(canonicalEvent, request, state),\n          state,\n        );\n      }\n      state.producerFinished = true;\n      this.completeWhenDurable(threadId, state);\n    }\n  }\n\n  /** Queue one immutable event payload until Redis-backed gateway acknowledgement. */\n  private queueRunnerEvent(\n    payload: Record<string, unknown>,\n    state: ThreadState,\n  ): void {\n    if (!this.isCurrentThreadState(state.threadId, state)) {\n      return;\n    }\n    const eventId = this.runnerEventId(payload);\n    if (!state.pendingEvents.has(eventId)) {\n      state.pendingEvents.set(eventId, {\n        payload: structuredClone(payload),\n        queuedAt: Date.now(),\n      });\n    }\n    this.scheduleEventDeadline(state);\n    this.replayPendingEvents(state);\n  }\n\n  private pushPendingEventBatch(\n    eventIds: string[],\n    events: Array<{ payload: Record<string, unknown>; queuedAt: number }>,\n    state: ThreadState,\n  ): void {\n    if (\n      !this.isCurrentThreadState(state.threadId, state) ||\n      this.failIfEventDeadlineExceeded(state) ||\n      state.channel.state !== \"joined\" ||\n      !state.socket.isConnected()\n    ) {\n      return;\n    }\n\n    const attempt = ++state.nextEventPushAttempt;\n    state.activeEventBatch = { eventIds, attempt };\n    const payloads = events.map((event) => event.payload);\n    const isBatch = state.supportsRunnerEventBatch;\n\n    state.channel\n      .push(\n        isBatch ? \"events\" : \"event\",\n        isBatch ? { events: payloads } : payloads[0],\n      )\n      .receive(\"ok\", () => {\n        if (\n          !this.isCurrentThreadState(state.threadId, state) ||\n          state.activeEventBatch?.attempt !== attempt\n        ) {\n          return;\n        }\n        if (this.failIfEventDeadlineExceeded(state)) {\n          return;\n        }\n        for (const eventId of eventIds) {\n          state.pendingEvents.delete(eventId);\n        }\n        state.activeEventBatch = null;\n        state.eventRetryAttempt = 0;\n        if (state.pendingEvents.size === 0) {\n          this.clearPendingEventRetry(state);\n          this.clearPendingEventFlush(state);\n        }\n        this.scheduleEventDeadline(state);\n        this.completeWhenDurable(state.threadId, state);\n        this.replayPendingEvents(state);\n      })\n      .receive(\"error\", (response) =>\n        this.handlePendingEventFailure(state, attempt, response),\n      )\n      .receive(\"timeout\", () => this.handlePendingEventFailure(state, attempt));\n  }\n\n  private handlePendingEventFailure(\n    state: ThreadState,\n    attempt: number,\n    response?: unknown,\n  ): void {\n    if (\n      !this.isCurrentThreadState(state.threadId, state) ||\n      state.activeEventBatch?.attempt !== attempt\n    ) {\n      return;\n    }\n\n    if (this.failIfEventDeadlineExceeded(state)) {\n      return;\n    }\n    if (this.isPermanentEventFailure(response)) {\n      const reason =\n        typeof response === \"object\" &&\n        response !== null &&\n        typeof (response as { reason?: unknown }).reason === \"string\"\n          ? (response as { reason: string }).reason\n          : \"permanent_gateway_rejection\";\n      this.failThread(\n        state.threadId,\n        state,\n        new Error(`Runner event durability failed: ${reason}`),\n      );\n      return;\n    }\n    this.schedulePendingEventRetry(state);\n  }\n\n  private schedulePendingEventRetry(state: ThreadState): void {\n    if (\n      !this.isCurrentThreadState(state.threadId, state) ||\n      state.pendingEvents.size === 0 ||\n      state.eventRetryTimer !== null\n    ) {\n      return;\n    }\n\n    const delay = Math.min(\n      EVENT_RETRY_BASE_MS * 2 ** state.eventRetryAttempt,\n      EVENT_RETRY_MAX_MS,\n    );\n    state.eventRetryAttempt += 1;\n    state.eventRetryTimer = setTimeout(() => {\n      state.eventRetryTimer = null;\n      if (\n        !this.isCurrentThreadState(state.threadId, state) ||\n        state.pendingEvents.size === 0\n      ) {\n        return;\n      }\n      this.retryActiveEventBatch(state);\n    }, delay);\n  }\n\n  private clearPendingEventRetry(state: ThreadState): void {\n    if (state.eventRetryTimer !== null) {\n      clearTimeout(state.eventRetryTimer);\n      state.eventRetryTimer = null;\n    }\n  }\n\n  private schedulePendingEventFlush(state: ThreadState): void {\n    if (\n      !this.isCurrentThreadState(state.threadId, state) ||\n      state.pendingEvents.size === 0 ||\n      state.activeEventBatch !== null ||\n      state.eventFlushTimer !== null\n    ) {\n      return;\n    }\n\n    const oldestQueuedAt = Math.min(\n      ...[...state.pendingEvents.values()].map((event) => event.queuedAt),\n    );\n    const delay = Math.max(\n      0,\n      oldestQueuedAt + RUNNER_EVENT_BATCH_FLUSH_MS - Date.now(),\n    );\n    state.eventFlushTimer = setTimeout(() => {\n      state.eventFlushTimer = null;\n      if (!this.isCurrentThreadState(state.threadId, state)) {\n        return;\n      }\n      this.flushPendingEventBatch(state);\n    }, delay);\n  }\n\n  private clearPendingEventFlush(state: ThreadState): void {\n    if (state.eventFlushTimer !== null) {\n      clearTimeout(state.eventFlushTimer);\n      state.eventFlushTimer = null;\n    }\n  }\n\n  private resetPendingEventRetry(state: ThreadState): void {\n    this.clearPendingEventRetry(state);\n    this.clearPendingEventFlush(state);\n    state.eventRetryAttempt = 0;\n    state.activeEventBatch = null;\n  }\n\n  private replayPendingEvents(state: ThreadState): void {\n    if (\n      !this.isCurrentThreadState(state.threadId, state) ||\n      state.activeEventBatch !== null\n    ) {\n      return;\n    }\n\n    const pendingEvents = [...state.pendingEvents.entries()].sort(\n      ([, left], [, right]) =>\n        this.runnerEventSeq(left.payload) - this.runnerEventSeq(right.payload),\n    );\n    if (pendingEvents.length === 0) {\n      return;\n    }\n\n    const batchSize = state.supportsRunnerEventBatch\n      ? MAX_RUNNER_EVENT_BATCH_SIZE\n      : 1;\n    if (!state.supportsRunnerEventBatch || pendingEvents.length >= batchSize) {\n      this.clearPendingEventFlush(state);\n      const batch = pendingEvents.slice(0, batchSize);\n      this.pushPendingEventBatch(\n        batch.map(([eventId]) => eventId),\n        batch.map(([, event]) => event),\n        state,\n      );\n      return;\n    }\n\n    this.schedulePendingEventFlush(state);\n  }\n\n  private flushPendingEventBatch(state: ThreadState): void {\n    if (\n      !this.isCurrentThreadState(state.threadId, state) ||\n      state.activeEventBatch !== null ||\n      state.pendingEvents.size === 0\n    ) {\n      return;\n    }\n\n    const batch = [...state.pendingEvents.entries()]\n      .sort(\n        ([, left], [, right]) =>\n          this.runnerEventSeq(left.payload) -\n          this.runnerEventSeq(right.payload),\n      )\n      .slice(0, MAX_RUNNER_EVENT_BATCH_SIZE);\n    this.pushPendingEventBatch(\n      batch.map(([eventId]) => eventId),\n      batch.map(([, event]) => event),\n      state,\n    );\n  }\n\n  private retryActiveEventBatch(state: ThreadState): void {\n    if (\n      !this.isCurrentThreadState(state.threadId, state) ||\n      this.failIfEventDeadlineExceeded(state)\n    ) {\n      return;\n    }\n    const activeBatch = state.activeEventBatch;\n    if (activeBatch === null) {\n      this.replayPendingEvents(state);\n      return;\n    }\n\n    const events = activeBatch.eventIds\n      .map((eventId) => state.pendingEvents.get(eventId))\n      .filter(\n        (\n          event,\n        ): event is { payload: Record<string, unknown>; queuedAt: number } =>\n          event !== undefined,\n      );\n    if (events.length !== activeBatch.eventIds.length) {\n      state.activeEventBatch = null;\n      this.replayPendingEvents(state);\n      return;\n    }\n\n    this.pushPendingEventBatch(activeBatch.eventIds, events, state);\n  }\n\n  private completeWhenDurable(threadId: string, state: ThreadState): void {\n    if (\n      !this.isCurrentThreadState(threadId, state) ||\n      !state.producerFinished ||\n      state.pendingEvents.size !== 0\n    ) {\n      return;\n    }\n    if (this.threads.get(threadId) !== state) {\n      return;\n    }\n\n    this.removeThread(threadId, state);\n    state.completeRun();\n  }\n\n  private runnerEventId(payload: Record<string, unknown>): string {\n    const metadata = payload.metadata as Record<string, unknown>;\n    return metadata.cpki_event_id as string;\n  }\n\n  private runnerEventSeq(payload: Record<string, unknown>): number {\n    const metadata = payload.metadata as Record<string, unknown>;\n    return metadata.cpki_event_seq as number;\n  }\n\n  private supportsRunnerEventBatch(response: unknown): boolean {\n    if (typeof response !== \"object\" || response === null) {\n      return false;\n    }\n    const capabilities = (response as { capabilities?: unknown }).capabilities;\n    return (\n      Array.isArray(capabilities) &&\n      capabilities.includes(RUNNER_EVENT_BATCH_CAPABILITY)\n    );\n  }\n\n  private isRetryableJoinError(response: unknown): boolean {\n    if (typeof response !== \"object\" || response === null) {\n      return false;\n    }\n    const value = response as { reason?: unknown; retryable?: unknown };\n    if (value.retryable === false) {\n      return false;\n    }\n    return value.retryable === true || value.reason === \"gateway_draining\";\n  }\n\n  private isPermanentEventFailure(response: unknown): boolean {\n    return (\n      typeof response === \"object\" &&\n      response !== null &&\n      (response as { retryable?: unknown }).retryable === false\n    );\n  }\n\n  private scheduleEventDeadline(state: ThreadState): void {\n    if (state.eventDeadlineTimer !== null) {\n      clearTimeout(state.eventDeadlineTimer);\n      state.eventDeadlineTimer = null;\n    }\n    if (\n      !this.isCurrentThreadState(state.threadId, state) ||\n      state.pendingEvents.size === 0\n    ) {\n      return;\n    }\n\n    const oldestQueuedAt = Math.min(\n      ...[...state.pendingEvents.values()].map((event) => event.queuedAt),\n    );\n    const deadline = oldestQueuedAt + EVENT_DURABILITY_DEADLINE_MS;\n    state.eventDeadlineTimer = setTimeout(\n      () => {\n        state.eventDeadlineTimer = null;\n        if (!this.isCurrentThreadState(state.threadId, state)) {\n          return;\n        }\n\n        if (!this.failIfEventDeadlineExceeded(state)) {\n          this.scheduleEventDeadline(state);\n        }\n      },\n      Math.max(0, deadline - Date.now()),\n    );\n  }\n\n  private failThread(threadId: string, state: ThreadState, error: Error): void {\n    if (!this.isCurrentThreadState(threadId, state)) {\n      return;\n    }\n    this.removeThread(threadId, state);\n    try {\n      state.agent?.abortRun();\n    } catch {\n      // The terminal durability error must still reach the subscriber.\n    }\n    state.failRun(error);\n  }\n\n  private failIfEventDeadlineExceeded(state: ThreadState): boolean {\n    if (state.pendingEvents.size === 0) {\n      return false;\n    }\n    const oldestQueuedAt = Math.min(\n      ...[...state.pendingEvents.values()].map((event) => event.queuedAt),\n    );\n    if (Date.now() < oldestQueuedAt + EVENT_DURABILITY_DEADLINE_MS) {\n      return false;\n    }\n    this.failThread(\n      state.threadId,\n      state,\n      new Error(\"Timed out trying to durably deliver runner events\"),\n    );\n    return true;\n  }\n\n  private isCurrentThreadState(threadId: string, state: ThreadState): boolean {\n    return state.isRunning && this.threads.get(threadId) === state;\n  }\n\n  /**\n   * Tear down all resources for a thread: leave the channel,\n   * disconnect the per-run socket, and remove the thread state.\n   *\n   * Idempotent — safe to call multiple times for the same threadId\n   * (e.g. from join error handlers, finalize, and Observable teardown).\n   */\n  private removeThread(threadId: string, state: ThreadState): void {\n    if (this.threads.get(threadId) !== state) {\n      return;\n    }\n\n    // Delete first so concurrent calls see the entry as already removed.\n    this.threads.delete(threadId);\n    state.isRunning = false;\n    this.clearPendingEventRetry(state);\n    this.clearPendingEventFlush(state);\n    if (state.eventDeadlineTimer !== null) {\n      clearTimeout(state.eventDeadlineTimer);\n      state.eventDeadlineTimer = null;\n    }\n    state.activeEventBatch = null;\n    if (state.socketReconnectWatchdog !== null) {\n      clearTimeout(state.socketReconnectWatchdog);\n      state.socketReconnectWatchdog = null;\n    }\n\n    try {\n      state.channel.leave();\n    } catch {\n      // Channel may already be closed/left.\n    }\n    try {\n      state.socket.disconnect();\n    } catch {\n      // Socket may already be disconnected.\n    }\n  }\n}\n"],"mappings":";;;;;;;;;;AAgEA,MAAM,gCAAgC;AACtC,MAAM,sBAAsB;AAC5B,MAAM,qBAAqB;AAC3B,MAAM,gCAAgC;AACtC,MAAM,8BAA8B;AACpC,MAAM,8BAA8B;AACpC,MAAM,+BAA+B;AAErC,IAAa,0BAAb,cAA6CA,iCAAY;CAIvD,YAAY,SAAyC;AACnD,SAAO;iCAHS,IAAI,KAA0B;AAK9C,OAAK,UAAU;;;;;;;;;;;;;;;;;;;;;;;CAwBjB,AAAQ,aAAa,YAAY,KAAK,QAAQ,WAAmB;EAC/D,MAAM,SAAS,IAAIC,eAAO,KAAK,QAAQ,KAAK;GAC1C,GAAI,YAAY,EAAE,WAAW,GAAG,EAAE;GAClC,oEACE,KACA,KAAK,QAAQ,kBAAkB,IAChC;GACD,iEACE,KACA,KAAK,QAAQ,eAAe,IAC7B;GACF,CAAC;AACF,SAAO,SAAS;AAChB,SAAO;;CAGT,AAAQ,yBACN,OACA,SACA,OACyB;EAKzB,MAAM,UAAU,EACd,GALqB,KAAK,oBAC1B,KAAK,2BAA2B,OAAO,QAAQ,EAC/C,MACD,EAGA;AAED,UAAQ,WAAW,QAAQ;AAC3B,UAAQ,QAAQ,QAAQ,MAAM;AAC9B,UAAQ,YAAY,QAAQ;AAC5B,UAAQ,SAAS,QAAQ,MAAM;AAE/B,SAAO;;CAGT,AAAQ,2BACN,OACA,SACW;EACX,MAAM,cAAc;AACpB,cAAY,WAAW,QAAQ;AAC/B,cAAY,QAAQ,QAAQ,MAAM;AAClC,SAAO;;CAGT,AAAQ,oBAAoB,OAAkB,OAA+B;EAC3E,MAAM,cAAc;EAIpB,MAAM,mBAAmB,YAAY,YAAY,EAAE;EACnD,MAAM,aAAa,OAAO,iBAAiB,kBAAkB;EAC7D,MAAM,cAAc,OAAO,iBAAiB,mBAAmB;AAE/D,MAAI,cAAc,aAAa;GAC7B,MAAM,WAAW,iBAAiB;AAClC,SAAM,eAAe,KAAK,IAAI,MAAM,cAAc,WAAW,EAAE;AAC/D,UAAO;;EAGT,MAAM,WAAW,MAAM;AAEvB,cAAY,WAAW;GACrB,GAAG;GACH,eACE,OAAO,iBAAiB,kBAAkB,WACtC,iBAAiB,6CACL;GAClB,gBAAgB;GACjB;AACD,SAAO;;CAGT,IAAI,SAAuD;AACzD,SAAO,KAAK,oBAAoB,QAAQ;;CAG1C,uBACE,SACuB;EACvB,IAAI;EACJ,IAAI;EACJ,MAAM,UAAU,IAAI,SAAe,SAAS,WAAW;AACrD,oBAAiB;AACjB,mBAAgB;IAChB;AAEF,SAAO;GACL,QAAQ,KAAK,oBAAoB,SAAS;IACxC,sBAAsB,kBAAkB;IACxC,gBAAgB,UAAU,gBAAgB,MAAM;IACjD,CAAC;GACF;GACD;;CAGH,AAAQ,oBACN,SACA,iBAIuB;EACvB,MAAM,EAAE,UAAU,OAAO,UAAU;AAGnC,MADiB,KAAK,QAAQ,IAAI,SAAS,EAC7B,UACZ,OAAM,IAAI,MAAM,yBAAyB;AAG3C,SAAO,IAAIC,iBAAY,aAAa;AAClC,OAAI,KAAK,QAAQ,IAAI,SAAS,EAAE,WAAW;AACzC,aAAS,sBAAM,IAAI,MAAM,yBAAyB,CAAC;AACnD;;GAGF,MAAM,SAAS,KAAK,aAAa,QAAQ,UAAU;GAEnD,MAAM,UAAU,OAAO,QAAQ,aAAa,MAAM,SAAS;IACzD,WAAW;IACX,QAAQ,MAAM;IACf,CAAC;GAEF,MAAM,QAAqB;IACzB;IACA,OAAO,MAAM;IACb;IACA;IACA,WAAW;IACX,eAAe;IACf;IACA,eAAe,EAAE;IACjB,cAAc;IACd,eAAe;IACf,WAAW;IACX,0BAA0B;IAC1B,kBAAkB;IAClB,+BAAe,IAAI,KAAK;IACxB,kBAAkB;IAClB,sBAAsB;IACtB,iBAAiB;IACjB,iBAAiB;IACjB,oBAAoB;IACpB,yBAAyB;IACzB,mBAAmB;IACnB,mBAAmB,SAAS,UAAU;IACtC,UAAU,UAAU,SAAS,MAAM,MAAM;IAC1C;AACD,QAAK,QAAQ,IAAI,UAAU,MAAM;GAEjC,IAAI,0BAA0B;GAC9B,IAAI,iBAAiB;AAErB,UAAO,SAAS,UAAU;AACxB,QAAI,CAAC,KAAK,qBAAqB,UAAU,MAAM,CAC7C;AAEF,qBAAiB,kBAAkB,OAAO,SAAS;AACnD,QAAI,OAAO,SAAS,OAAQ,MAAM,4BAA4B,KAC5D,OAAM,0BAA0B,iBAAiB;AAC/C,WAAM,0BAA0B;AAChC,SAAI,CAAC,MAAM,aAAa,OAAO,aAAa,CAC1C;AAEF,YAAO,iBAAiB;AACtB,UAAI,MAAM,aAAa,CAAC,OAAO,aAAa,CAC1C,QAAO,SAAS;OAElB;OACD,IAAM;KAEX;AACF,UAAO,aAAa;AAClB,QAAI,CAAC,KAAK,qBAAqB,UAAU,MAAM,CAC7C;AAEF,QAAI,MAAM,4BAA4B,MAAM;AAC1C,kBAAa,MAAM,wBAAwB;AAC3C,WAAM,0BAA0B;;AAElC,8BAA0B;AAC1B,qBAAiB;KACjB;AACF,UAAO,cAAc;AACnB,QAAI,CAAC,KAAK,qBAAqB,UAAU,MAAM,CAC7C;AAMF,QAAI,kBAAkB,MAAM,UAC1B;AAGF,+BAA2B;AAC3B,QAAI,2BAA2B,8BAC7B,OAAM,OAAO,UAAU;KAEzB;AAOF,WAAQ,GAAGC,yCAAsB,YAAuB;AACtD,QACE,KAAK,qBAAqB,UAAU,MAAM,IAC1C,QAAQ,SAASC,wBAAU,UAC1B,QAA0C,SAAS,OAEpD,MAAK,KAAK;KAAE;KAAU,OAAO,MAAM;KAAO,CAAC;KAE7C;AAEF,WACG,MAAM,CACN,QAAQ,OAAO,aAAa;AAC3B,QAAI,CAAC,KAAK,qBAAqB,UAAU,MAAM,CAC7C;IAEF,MAAM,2BACJ,KAAK,yBAAyB,SAAS;IACzC,MAAM,mBACJ,MAAM,6BAA6B,2BAC/B,MAAM,mBACN;AACN,UAAM,2BAA2B;AACjC,QAAI,MAAM,WAAW;AACnB,UAAK,uBAAuB,MAAM;AAClC,SAAI,qBAAqB,MAAM;AAC7B,YAAM,mBAAmB;AACzB,WAAK,sBAAsB,MAAM;WAEjC,MAAK,oBAAoB,MAAM;AAEjC;;AAGF,UAAM,YAAY;AAClB,qBAAiB,gBAAgB;AACjC,IAAK,KAAK,gBAAgB,SAAS,OAAO,WAAW,UAAU;AAC7D,cAAS,KAAK,MAAM;MACpB;KACF,CACD,QAAQ,UAAU,SAAS;AAC1B,QAAI,CAAC,KAAK,qBAAqB,UAAU,MAAM,CAC7C;AAEF,QAAI,MAAM,WAAW;AACnB,SAAI,KAAK,wBAAwB,KAAK,EAAE;MACtC,MAAM,SACJ,OAAQ,KAA8B,WAAW,WAC7C,KAAM,KAA4B,WAClC;AACN,WAAK,WACH,UACA,uBACA,IAAI,MACF,8CAA8C,SAC/C,CACF;;AAEH;;AAEF,QAAI,KAAK,qBAAqB,KAAK,CACjC;IAGF,MAAM,wBAAQ,IAAI,MAChB,2BAA2B,KAAK,UAAU,KAAK,GAChD;IACD,MAAM,aAAa;KACjB,MAAMA,wBAAU;KAChB,SAAS,MAAM;KACf,MAAM;KACP;AACD,aAAS,KAAK,WAAW;AACzB,UAAM,cAAc,KAAK,WAAW;AACpC,SAAK,aAAa,UAAU,MAAM;AAClC,qBAAiB,cAAc,MAAM;AACrC,aAAS,UAAU;KACnB,CACD,QAAQ,iBAAiB;AACxB,QAAI,CAAC,KAAK,qBAAqB,UAAU,MAAM,CAC7C;AAEF,QAAI,MAAM,UACR;IAGF,MAAM,wBAAQ,IAAI,MAAM,4BAA4B;IACpD,MAAM,aAAa;KACjB,MAAMA,wBAAU;KAChB,SAAS,MAAM;KACf,MAAM;KACP;AACD,aAAS,KAAK,WAAW;AACzB,UAAM,cAAc,KAAK,WAAW;AACpC,SAAK,aAAa,UAAU,MAAM;AAClC,qBAAiB,cAAc,MAAM;AACrC,aAAS,UAAU;KACnB;AAEJ,gBAAa;AACX,SAAK,aAAa,UAAU,MAAM;;IAEpC;;CAGJ,QAAQ,SAA2D;EACjE,MAAM,EAAE,aAAa;AAErB,SAAO,IAAIF,iBAAY,aAAa;GAClC,MAAM,SAAS,KAAK,cAAc;GAElC,MAAM,UAAU,OAAO,QAAQ,UAAU,WAAW;AAEpD,WAAQ,GAAG,gBAAgB,YAAuB;AAChD,aAAS,KAAK,QAAQ;AAEtB,QACE,QAAQ,SAASE,wBAAU,gBAC3B,QAAQ,SAASA,wBAAU,UAE3B,UAAS,UAAU;KAErB;GAEF,MAAM,gBAAgB;AACpB,YAAQ,OAAO;AACf,WAAO,YAAY;;AAGrB,WACG,MAAM,CACN,QAAQ,YAAY,OAAU,CAC9B,QAAQ,UAAU,SAAS;AAC1B,aAAS,sBACP,IAAI,MAAM,2BAA2B,KAAK,UAAU,KAAK,GAAG,CAC7D;AACD,aAAS;KACT,CACD,QAAQ,iBAAiB;AACxB,aAAS,sBAAM,IAAI,MAAM,4BAA4B,CAAC;AACtD,aAAS;KACT;AAEJ,gBAAa;AACX,aAAS;;IAEX;;CAGJ,UAAU,SAAwD;EAChE,MAAM,QAAQ,KAAK,QAAQ,IAAI,QAAQ,SAAS;AAChD,SAAO,QAAQ,QAAQ,OAAO,aAAa,MAAM;;CAGnD,KAAK,SAA+D;EAClE,MAAM,QAAQ,KAAK,QAAQ,IAAI,QAAQ,SAAS;AAChD,MAAI,CAAC,SAAS,CAAC,MAAM,aAAa,MAAM,cACtC,QAAO,QAAQ,QAAQ,MAAM;AAE/B,MAAI,QAAQ,UAAU,UAAa,MAAM,UAAU,QAAQ,MACzD,QAAO,QAAQ,QAAQ,MAAM;AAG/B,QAAM,gBAAgB;AAGtB,MAAI,MAAM,MACR,KAAI;AACF,SAAM,MAAM,UAAU;UAChB;AAKV,SAAO,QAAQ,QAAQ,KAAK;;CAG9B,MAAc,gBACZ,SACA,OACA,UACA,YACe;EACf,MAAM,EAAE,kBAAkB;EAC1B,MAAM,sBAAsB,UAA2B;AACrD,OAAI,CAAC,KAAK,qBAAqB,UAAU,MAAM,CAC7C;GAEF,MAAM,iBAAiB,KAAK,oBAC1B,KAAK,2BAA2B,OAAO,QAAQ,EAC/C,MACD;AACD,iBAAc,KAAK,eAAe;AAElC,OAAI,eAAe,SAASA,wBAAU,YACpC,OAAM,gBAAgB;AAGxB,QAAK,iBACH,KAAK,yBAAyB,gBAAgB,SAAS,MAAM,EAC7D,MACD;;EAGH,MAAM,kCACJ,QAAQ,0BAA0B,QAAQ,MAAM;EAElD,MAAM,wBACJ,WACoB;GACpB,MAAM,YAAY,QAAQ,SAAS,QAAQ;GAC3C,MAAM,yBAAyB,2BAA2B;GAC1D,MAAM,QACJ,UACC;IACC,MAAMA,wBAAU;IAChB,UAAU,QAAQ;IAClB,OAAO,QAAQ,MAAM;IACtB;AACH,SAAM,WAAW,QAAQ;AACzB,SAAM,QAAQ,QAAQ,MAAM;AAC5B,SAAM,QAAQ;IACZ,GAAG;IACH,UAAU,QAAQ;IAClB,OAAO,QAAQ,MAAM;IACrB,GAAI,2BAA2B,SAC3B,EAAE,UAAU,wBAAwB,GACpC,EAAE;IACP;AACD,UAAO;;EAGT,MAAM,yBAA+B;AACnC,OAAI,CAAC,MAAM,eAAe;AACxB,UAAM,gBAAgB;AACtB,uBAAmB,sBAAsB,CAAC;;;AAI9C,MAAI;AACF,SAAM,QAAQ,MAAM,SAAS,QAAQ,OAAO,EAC1C,UAAU,EAAE,YAAkC;AAC5C,QAAI,MAAM,SAASA,wBAAU,aAAa;AACxC,wBAAmB,qBAAqB,MAAyB,CAAC;AAClE;;AAGF,sBAAkB;AAClB,uBAAmB,MAAM;MAE5B,CAAC;WACK,OAAO;AACd,OAAI,CAAC,KAAK,qBAAqB,UAAU,MAAM,CAC7C;AAEF,qBAAkB;GAClB,MAAM,gBAAgB,cAAc,MACjC,UAAU,MAAM,SAASA,wBAAU,UACrC;AACD,OAAI,cACF,YAAW,cAAc;QACpB;IACL,MAAM,aAAa;KACjB,MAAMA,wBAAU;KAChB,SAAS,iBAAiB,QAAQ,MAAM,UAAU,OAAO,MAAM;KAChE;AACD,uBAAmB,WAAW;AAC9B,eAAW,WAAW;;YAEhB;AACR,OAAI,CAAC,KAAK,qBAAqB,UAAU,MAAM,CAC7C;AAEF,qBAAkB;GAClB,MAAM,qDAA6B,eAAe,EAChD,eAAe,MAAM,eACtB,CAAC;AACF,QAAK,MAAM,SAAS,UAAU;IAC5B,MAAM,iBAAiB,KAAK,oBAC1B,KAAK,2BAA2B,OAAO,QAAQ,EAC/C,MACD;AACD,SAAK,iBACH,KAAK,yBAAyB,gBAAgB,SAAS,MAAM,EAC7D,MACD;;AAEH,SAAM,mBAAmB;AACzB,QAAK,oBAAoB,UAAU,MAAM;;;;CAK7C,AAAQ,iBACN,SACA,OACM;AACN,MAAI,CAAC,KAAK,qBAAqB,MAAM,UAAU,MAAM,CACnD;EAEF,MAAM,UAAU,KAAK,cAAc,QAAQ;AAC3C,MAAI,CAAC,MAAM,cAAc,IAAI,QAAQ,CACnC,OAAM,cAAc,IAAI,SAAS;GAC/B,SAAS,gBAAgB,QAAQ;GACjC,UAAU,KAAK,KAAK;GACrB,CAAC;AAEJ,OAAK,sBAAsB,MAAM;AACjC,OAAK,oBAAoB,MAAM;;CAGjC,AAAQ,sBACN,UACA,QACA,OACM;AACN,MACE,CAAC,KAAK,qBAAqB,MAAM,UAAU,MAAM,IACjD,KAAK,4BAA4B,MAAM,IACvC,MAAM,QAAQ,UAAU,YACxB,CAAC,MAAM,OAAO,aAAa,CAE3B;EAGF,MAAM,UAAU,EAAE,MAAM;AACxB,QAAM,mBAAmB;GAAE;GAAU;GAAS;EAC9C,MAAM,WAAW,OAAO,KAAK,UAAU,MAAM,QAAQ;EACrD,MAAM,UAAU,MAAM;AAEtB,QAAM,QACH,KACC,UAAU,WAAW,SACrB,UAAU,EAAE,QAAQ,UAAU,GAAG,SAAS,GAC3C,CACA,QAAQ,YAAY;AACnB,OACE,CAAC,KAAK,qBAAqB,MAAM,UAAU,MAAM,IACjD,MAAM,kBAAkB,YAAY,QAEpC;AAEF,OAAI,KAAK,4BAA4B,MAAM,CACzC;AAEF,QAAK,MAAM,WAAW,SACpB,OAAM,cAAc,OAAO,QAAQ;AAErC,SAAM,mBAAmB;AACzB,SAAM,oBAAoB;AAC1B,OAAI,MAAM,cAAc,SAAS,GAAG;AAClC,SAAK,uBAAuB,MAAM;AAClC,SAAK,uBAAuB,MAAM;;AAEpC,QAAK,sBAAsB,MAAM;AACjC,QAAK,oBAAoB,MAAM,UAAU,MAAM;AAC/C,QAAK,oBAAoB,MAAM;IAC/B,CACD,QAAQ,UAAU,aACjB,KAAK,0BAA0B,OAAO,SAAS,SAAS,CACzD,CACA,QAAQ,iBAAiB,KAAK,0BAA0B,OAAO,QAAQ,CAAC;;CAG7E,AAAQ,0BACN,OACA,SACA,UACM;AACN,MACE,CAAC,KAAK,qBAAqB,MAAM,UAAU,MAAM,IACjD,MAAM,kBAAkB,YAAY,QAEpC;AAGF,MAAI,KAAK,4BAA4B,MAAM,CACzC;AAEF,MAAI,KAAK,wBAAwB,SAAS,EAAE;GAC1C,MAAM,SACJ,OAAO,aAAa,YACpB,aAAa,QACb,OAAQ,SAAkC,WAAW,WAChD,SAAgC,SACjC;AACN,QAAK,WACH,MAAM,UACN,uBACA,IAAI,MAAM,mCAAmC,SAAS,CACvD;AACD;;AAEF,OAAK,0BAA0B,MAAM;;CAGvC,AAAQ,0BAA0B,OAA0B;AAC1D,MACE,CAAC,KAAK,qBAAqB,MAAM,UAAU,MAAM,IACjD,MAAM,cAAc,SAAS,KAC7B,MAAM,oBAAoB,KAE1B;EAGF,MAAM,QAAQ,KAAK,IACjB,sBAAsB,KAAK,MAAM,mBACjC,mBACD;AACD,QAAM,qBAAqB;AAC3B,QAAM,kBAAkB,iBAAiB;AACvC,SAAM,kBAAkB;AACxB,OACE,CAAC,KAAK,qBAAqB,MAAM,UAAU,MAAM,IACjD,MAAM,cAAc,SAAS,EAE7B;AAEF,QAAK,sBAAsB,MAAM;KAChC,MAAM;;CAGX,AAAQ,uBAAuB,OAA0B;AACvD,MAAI,MAAM,oBAAoB,MAAM;AAClC,gBAAa,MAAM,gBAAgB;AACnC,SAAM,kBAAkB;;;CAI5B,AAAQ,0BAA0B,OAA0B;AAC1D,MACE,CAAC,KAAK,qBAAqB,MAAM,UAAU,MAAM,IACjD,MAAM,cAAc,SAAS,KAC7B,MAAM,qBAAqB,QAC3B,MAAM,oBAAoB,KAE1B;EAGF,MAAM,iBAAiB,KAAK,IAC1B,GAAG,CAAC,GAAG,MAAM,cAAc,QAAQ,CAAC,CAAC,KAAK,UAAU,MAAM,SAAS,CACpE;EACD,MAAM,QAAQ,KAAK,IACjB,GACA,iBAAiB,8BAA8B,KAAK,KAAK,CAC1D;AACD,QAAM,kBAAkB,iBAAiB;AACvC,SAAM,kBAAkB;AACxB,OAAI,CAAC,KAAK,qBAAqB,MAAM,UAAU,MAAM,CACnD;AAEF,QAAK,uBAAuB,MAAM;KACjC,MAAM;;CAGX,AAAQ,uBAAuB,OAA0B;AACvD,MAAI,MAAM,oBAAoB,MAAM;AAClC,gBAAa,MAAM,gBAAgB;AACnC,SAAM,kBAAkB;;;CAI5B,AAAQ,uBAAuB,OAA0B;AACvD,OAAK,uBAAuB,MAAM;AAClC,OAAK,uBAAuB,MAAM;AAClC,QAAM,oBAAoB;AAC1B,QAAM,mBAAmB;;CAG3B,AAAQ,oBAAoB,OAA0B;AACpD,MACE,CAAC,KAAK,qBAAqB,MAAM,UAAU,MAAM,IACjD,MAAM,qBAAqB,KAE3B;EAGF,MAAM,gBAAgB,CAAC,GAAG,MAAM,cAAc,SAAS,CAAC,CAAC,MACtD,GAAG,OAAO,GAAG,WACZ,KAAK,eAAe,KAAK,QAAQ,GAAG,KAAK,eAAe,MAAM,QAAQ,CACzE;AACD,MAAI,cAAc,WAAW,EAC3B;EAGF,MAAM,YAAY,MAAM,2BACpB,8BACA;AACJ,MAAI,CAAC,MAAM,4BAA4B,cAAc,UAAU,WAAW;AACxE,QAAK,uBAAuB,MAAM;GAClC,MAAM,QAAQ,cAAc,MAAM,GAAG,UAAU;AAC/C,QAAK,sBACH,MAAM,KAAK,CAAC,aAAa,QAAQ,EACjC,MAAM,KAAK,GAAG,WAAW,MAAM,EAC/B,MACD;AACD;;AAGF,OAAK,0BAA0B,MAAM;;CAGvC,AAAQ,uBAAuB,OAA0B;AACvD,MACE,CAAC,KAAK,qBAAqB,MAAM,UAAU,MAAM,IACjD,MAAM,qBAAqB,QAC3B,MAAM,cAAc,SAAS,EAE7B;EAGF,MAAM,QAAQ,CAAC,GAAG,MAAM,cAAc,SAAS,CAAC,CAC7C,MACE,GAAG,OAAO,GAAG,WACZ,KAAK,eAAe,KAAK,QAAQ,GACjC,KAAK,eAAe,MAAM,QAAQ,CACrC,CACA,MAAM,GAAG,4BAA4B;AACxC,OAAK,sBACH,MAAM,KAAK,CAAC,aAAa,QAAQ,EACjC,MAAM,KAAK,GAAG,WAAW,MAAM,EAC/B,MACD;;CAGH,AAAQ,sBAAsB,OAA0B;AACtD,MACE,CAAC,KAAK,qBAAqB,MAAM,UAAU,MAAM,IACjD,KAAK,4BAA4B,MAAM,CAEvC;EAEF,MAAM,cAAc,MAAM;AAC1B,MAAI,gBAAgB,MAAM;AACxB,QAAK,oBAAoB,MAAM;AAC/B;;EAGF,MAAM,SAAS,YAAY,SACxB,KAAK,YAAY,MAAM,cAAc,IAAI,QAAQ,CAAC,CAClD,QAEG,UAEA,UAAU,OACb;AACH,MAAI,OAAO,WAAW,YAAY,SAAS,QAAQ;AACjD,SAAM,mBAAmB;AACzB,QAAK,oBAAoB,MAAM;AAC/B;;AAGF,OAAK,sBAAsB,YAAY,UAAU,QAAQ,MAAM;;CAGjE,AAAQ,oBAAoB,UAAkB,OAA0B;AACtE,MACE,CAAC,KAAK,qBAAqB,UAAU,MAAM,IAC3C,CAAC,MAAM,oBACP,MAAM,cAAc,SAAS,EAE7B;AAEF,MAAI,KAAK,QAAQ,IAAI,SAAS,KAAK,MACjC;AAGF,OAAK,aAAa,UAAU,MAAM;AAClC,QAAM,aAAa;;CAGrB,AAAQ,cAAc,SAA0C;AAE9D,SADiB,QAAQ,SACT;;CAGlB,AAAQ,eAAe,SAA0C;AAE/D,SADiB,QAAQ,SACT;;CAGlB,AAAQ,yBAAyB,UAA4B;AAC3D,MAAI,OAAO,aAAa,YAAY,aAAa,KAC/C,QAAO;EAET,MAAM,eAAgB,SAAwC;AAC9D,SACE,MAAM,QAAQ,aAAa,IAC3B,aAAa,SAAS,8BAA8B;;CAIxD,AAAQ,qBAAqB,UAA4B;AACvD,MAAI,OAAO,aAAa,YAAY,aAAa,KAC/C,QAAO;EAET,MAAM,QAAQ;AACd,MAAI,MAAM,cAAc,MACtB,QAAO;AAET,SAAO,MAAM,cAAc,QAAQ,MAAM,WAAW;;CAGtD,AAAQ,wBAAwB,UAA4B;AAC1D,SACE,OAAO,aAAa,YACpB,aAAa,QACZ,SAAqC,cAAc;;CAIxD,AAAQ,sBAAsB,OAA0B;AACtD,MAAI,MAAM,uBAAuB,MAAM;AACrC,gBAAa,MAAM,mBAAmB;AACtC,SAAM,qBAAqB;;AAE7B,MACE,CAAC,KAAK,qBAAqB,MAAM,UAAU,MAAM,IACjD,MAAM,cAAc,SAAS,EAE7B;EAMF,MAAM,WAHiB,KAAK,IAC1B,GAAG,CAAC,GAAG,MAAM,cAAc,QAAQ,CAAC,CAAC,KAAK,UAAU,MAAM,SAAS,CACpE,GACiC;AAClC,QAAM,qBAAqB,iBACnB;AACJ,SAAM,qBAAqB;AAC3B,OAAI,CAAC,KAAK,qBAAqB,MAAM,UAAU,MAAM,CACnD;AAGF,OAAI,CAAC,KAAK,4BAA4B,MAAM,CAC1C,MAAK,sBAAsB,MAAM;KAGrC,KAAK,IAAI,GAAG,WAAW,KAAK,KAAK,CAAC,CACnC;;CAGH,AAAQ,WAAW,UAAkB,OAAoB,OAAoB;AAC3E,MAAI,CAAC,KAAK,qBAAqB,UAAU,MAAM,CAC7C;AAEF,OAAK,aAAa,UAAU,MAAM;AAClC,MAAI;AACF,SAAM,OAAO,UAAU;UACjB;AAGR,QAAM,QAAQ,MAAM;;CAGtB,AAAQ,4BAA4B,OAA6B;AAC/D,MAAI,MAAM,cAAc,SAAS,EAC/B,QAAO;EAET,MAAM,iBAAiB,KAAK,IAC1B,GAAG,CAAC,GAAG,MAAM,cAAc,QAAQ,CAAC,CAAC,KAAK,UAAU,MAAM,SAAS,CACpE;AACD,MAAI,KAAK,KAAK,GAAG,iBAAiB,6BAChC,QAAO;AAET,OAAK,WACH,MAAM,UACN,uBACA,IAAI,MAAM,oDAAoD,CAC/D;AACD,SAAO;;CAGT,AAAQ,qBAAqB,UAAkB,OAA6B;AAC1E,SAAO,MAAM,aAAa,KAAK,QAAQ,IAAI,SAAS,KAAK;;;;;;;;;CAU3D,AAAQ,aAAa,UAAkB,OAA0B;AAC/D,MAAI,KAAK,QAAQ,IAAI,SAAS,KAAK,MACjC;AAIF,OAAK,QAAQ,OAAO,SAAS;AAC7B,QAAM,YAAY;AAClB,OAAK,uBAAuB,MAAM;AAClC,OAAK,uBAAuB,MAAM;AAClC,MAAI,MAAM,uBAAuB,MAAM;AACrC,gBAAa,MAAM,mBAAmB;AACtC,SAAM,qBAAqB;;AAE7B,QAAM,mBAAmB;AACzB,MAAI,MAAM,4BAA4B,MAAM;AAC1C,gBAAa,MAAM,wBAAwB;AAC3C,SAAM,0BAA0B;;AAGlC,MAAI;AACF,SAAM,QAAQ,OAAO;UACf;AAGR,MAAI;AACF,SAAM,OAAO,YAAY;UACnB"}