{"version":3,"file":"connection.cjs","names":["WaitGroup","ConnectionState","ReconnectError","AuthError","ConnectionLimitError","expBackoff","waitWithCancel","createStartRequest","headers: Record<string, string>","headerKeys","parseStartResponse","resolveWebsocketConnected:\n      | ((value: void | PromiseLike<void>) => void)\n      | undefined","rejectWebsocketConnected: ((reason?: unknown) => void) | undefined","onConnectionError: (error: unknown) => void | Promise<void>","workerDisconnectReasonToJSON","WorkerDisconnectReason","heartbeatIntervalMs: number | undefined","extendLeaseIntervalMs: number | undefined","parseConnectMessage","gatewayMessageTypeToJSON","GatewayMessageType","WorkerConnectRequestData","getPlatformName","allProcessEnv","version","retrieveSystemAttributes","getHostname","ensureUnsharedArrayBuffer","ConnectMessage","GatewayConnectionReadyData","conn: Connection","parseGatewayExecutorRequest","WorkerRequestAckData","extendLeaseInterval: NodeJS.Timeout | undefined","WorkerRequestExtendLeaseData","parseWorkerReplyAck","WorkerRequestExtendLeaseAckData","heartbeatInterval: NodeJS.Timeout | undefined","resolveApiBaseUrl"],"sources":["../../../../../src/components/connect/strategies/core/connection.ts"],"sourcesContent":["/**\n * Shared connection core logic used by both SameThreadStrategy and\n * WorkerThreadStrategy.\n *\n * This module extracts the common WebSocket connection management, handshake,\n * heartbeat, lease extension, and reconnection logic.\n */\n\nimport { WaitGroup } from \"@jpwilliams/waitgroup\";\nimport ms from \"ms\";\nimport { headerKeys } from \"../../../../helpers/consts.ts\";\nimport { allProcessEnv, getPlatformName } from \"../../../../helpers/env.ts\";\nimport { resolveApiBaseUrl } from \"../../../../helpers/url.ts\";\nimport type { Logger } from \"../../../../middleware/logger.ts\";\nimport {\n  ConnectMessage,\n  GatewayConnectionReadyData,\n  type GatewayExecutorRequestData,\n  GatewayMessageType,\n  gatewayMessageTypeToJSON,\n  WorkerConnectRequestData,\n  WorkerDisconnectReason,\n  WorkerRequestAckData,\n  WorkerRequestExtendLeaseAckData,\n  WorkerRequestExtendLeaseData,\n  workerDisconnectReasonToJSON,\n} from \"../../../../proto/src/components/connect/protobuf/connect.ts\";\nimport { version } from \"../../../../version.ts\";\nimport { ensureUnsharedArrayBuffer } from \"../../buffer.ts\";\nimport {\n  createStartRequest,\n  parseConnectMessage,\n  parseGatewayExecutorRequest,\n  parseStartResponse,\n  parseWorkerReplyAck,\n} from \"../../messages.ts\";\nimport { getHostname, retrieveSystemAttributes } from \"../../os.ts\";\nimport { ConnectionState } from \"../../types.ts\";\nimport {\n  AuthError,\n  ConnectionLimitError,\n  expBackoff,\n  ReconnectError,\n  waitWithCancel,\n} from \"../../util.ts\";\nimport type { BaseConnectionConfig } from \"./types.ts\";\n\nconst ConnectWebSocketProtocol = \"v0.connect.inngest.com\";\n\nfunction toError(value: unknown): Error {\n  if (value instanceof Error) {\n    return value;\n  }\n  return new Error(String(value));\n}\n\n/**\n * Connection object representing an active WebSocket connection.\n */\nexport interface Connection {\n  id: string;\n  ws: WebSocket;\n  cleanup: () => void | Promise<void>;\n  pendingHeartbeats: number;\n}\n\n/**\n * Configuration for the connection core.\n * Extends BaseConnectionConfig with connection-specific options.\n */\nexport interface ConnectionCoreConfig extends BaseConnectionConfig {\n  instanceId?: string;\n  maxWorkerConcurrency?: number;\n  gatewayUrl?: string;\n  appIds: string[];\n}\n\n/**\n * Callbacks for connection core events.\n */\nexport interface ConnectionCoreCallbacks {\n  logger: Logger;\n  onStateChange: (state: ConnectionState) => void;\n  getState: () => ConnectionState;\n  handleExecutionRequest: (\n    request: GatewayExecutorRequestData,\n  ) => Promise<Uint8Array>;\n  onReplyAck?: (requestId: string) => void;\n  onBufferResponse?: (requestId: string, responseBytes: Uint8Array) => void;\n  beforeConnect?: (signingKey: string | undefined) => Promise<void>;\n}\n\n/**\n * Core connection manager that handles WebSocket connection lifecycle,\n * handshake, heartbeat, lease extension, and reconnection.\n */\nexport class ConnectionCore {\n  private config: ConnectionCoreConfig;\n  private callbacks: ConnectionCoreCallbacks;\n\n  private currentConnection: Connection | undefined;\n  private excludeGateways: Set<string> = new Set();\n\n  private inProgressRequests: {\n    wg: WaitGroup;\n    requestLeases: Record<string, string>;\n  } = {\n    wg: new WaitGroup(),\n    requestLeases: {},\n  };\n\n  constructor(\n    config: ConnectionCoreConfig,\n    callbacks: ConnectionCoreCallbacks,\n  ) {\n    this.config = config;\n    this.callbacks = callbacks;\n  }\n\n  get connection(): Connection | undefined {\n    return this.currentConnection;\n  }\n\n  get connectionId(): string | undefined {\n    return this.currentConnection?.id;\n  }\n\n  /**\n   * Wait for all in-progress requests to complete.\n   */\n  async waitForInProgress(): Promise<void> {\n    await this.inProgressRequests.wg.wait();\n  }\n\n  /**\n   * Main connection loop with reconnection logic.\n   */\n  async connect(attempt = 0, path: string[] = []): Promise<void> {\n    if (typeof WebSocket === \"undefined\") {\n      throw new Error(\"WebSockets not supported in current environment\");\n    }\n\n    const state = this.callbacks.getState();\n    if (state === ConnectionState.CLOSED) {\n      throw new Error(\"Connection already closed\");\n    }\n\n    this.callbacks.logger.debug({ attempt }, \"Establishing connection\");\n\n    let useSigningKey = this.config.hashedSigningKey;\n\n    while (true) {\n      const currentState = this.callbacks.getState();\n      if (currentState === ConnectionState.CLOSED) {\n        break;\n      }\n\n      // NOTE: We can get here when the state is CLOSING, therefore it's\n      // possible to reconnect while the CLOSING. This is intentional so that\n      // the worker can reconnect while waiting for pending requests during a\n      // graceful shutdown. If we didn't allow reconnect in that case,\n      // heartbeats and lease extensions would stop and the Inngest Server would\n      // think the worker died.\n      //\n      // However, the state can be CLOSING during a shutdown without pending\n      // requests. The window of that happening is very small, but it's\n      // technically possible that we could mistakenly reconnect during a\n      // shutdown if the Inngest Server send a drain message.\n\n      // Flush any pending messages before attempting connection\n      if (this.callbacks.beforeConnect) {\n        await this.callbacks.beforeConnect(useSigningKey);\n      }\n\n      try {\n        await this.prepareConnection(useSigningKey, attempt, [...path]);\n        return;\n      } catch (err) {\n        this.callbacks.logger.warn({ err: toError(err) }, \"Failed to connect\");\n\n        if (!(err instanceof ReconnectError)) {\n          throw err;\n        }\n\n        attempt = err.attempt;\n\n        if (err instanceof AuthError) {\n          const switchToFallback =\n            useSigningKey === this.config.hashedSigningKey;\n          if (switchToFallback) {\n            this.callbacks.logger.debug(\"Switching to fallback signing key\");\n          }\n          useSigningKey = switchToFallback\n            ? this.config.hashedFallbackKey\n            : this.config.hashedSigningKey;\n        }\n\n        if (err instanceof ConnectionLimitError) {\n          this.callbacks.logger.error(\n            \"You have reached the maximum number of concurrent connections. Please disconnect other active workers to continue.\",\n          );\n        }\n\n        const delay = expBackoff(attempt);\n        this.callbacks.logger.debug({ delay }, \"Reconnecting\");\n\n        const cancelled = await waitWithCancel(delay, () => {\n          return this.callbacks.getState() === ConnectionState.CLOSED;\n        });\n        if (cancelled) {\n          this.callbacks.logger.debug(\"Reconnect backoff cancelled\");\n          break;\n        }\n\n        attempt++;\n      }\n    }\n\n    this.callbacks.logger.debug(\"Exiting connect loop\");\n  }\n\n  /**\n   * Clean up the current connection.\n   */\n  async cleanup(): Promise<void> {\n    if (this.currentConnection) {\n      await this.currentConnection.cleanup();\n      this.currentConnection = undefined;\n    }\n  }\n\n  private async sendStartRequest(\n    hashedSigningKey: string | undefined,\n    attempt: number,\n  ) {\n    const msg = createStartRequest(Array.from(this.excludeGateways));\n\n    const headers: Record<string, string> = {\n      \"Content-Type\": \"application/protobuf\",\n      ...(hashedSigningKey\n        ? { Authorization: `Bearer ${hashedSigningKey}` }\n        : {}),\n    };\n\n    if (this.config.envName) {\n      headers[headerKeys.Environment] = this.config.envName;\n    }\n\n    const targetUrl = new URL(\"/v0/connect/start\", await this.getApiBaseUrl());\n\n    let resp;\n    try {\n      resp = await fetch(targetUrl, {\n        method: \"POST\",\n        body: new Uint8Array(msg),\n        headers: headers,\n      });\n    } catch (err) {\n      const errMsg = err instanceof Error ? err.message : \"Unknown error\";\n      throw new ReconnectError(\n        `Failed initial API handshake request to ${targetUrl.toString()}, ${errMsg}`,\n        attempt,\n      );\n    }\n\n    if (!resp.ok) {\n      if (resp.status === 401) {\n        throw new AuthError(\n          `Failed initial API handshake request to ${targetUrl.toString()}${\n            this.config.envName ? ` (env: ${this.config.envName})` : \"\"\n          }, ${await resp.text()}`,\n          attempt,\n        );\n      }\n\n      if (resp.status === 429) {\n        throw new ConnectionLimitError(attempt);\n      }\n\n      throw new ReconnectError(\n        `Failed initial API handshake request to ${targetUrl.toString()}, ${await resp.text()}`,\n        attempt,\n      );\n    }\n\n    const startResp = await parseStartResponse(resp);\n    return startResp;\n  }\n\n  private async prepareConnection(\n    hashedSigningKey: string | undefined,\n    attempt: number,\n    path: string[] = [],\n  ): Promise<Connection> {\n    let closed = false;\n\n    this.callbacks.logger.debug({ attempt, path }, \"Preparing connection\");\n\n    const startedAt = new Date();\n    const startResp = await this.sendStartRequest(hashedSigningKey, attempt);\n\n    const connectionId = startResp.connectionId;\n    path.push(connectionId);\n\n    let resolveWebsocketConnected:\n      | ((value: void | PromiseLike<void>) => void)\n      | undefined;\n    let rejectWebsocketConnected: ((reason?: unknown) => void) | undefined;\n    const websocketConnectedPromise = new Promise<void>((resolve, reject) => {\n      resolveWebsocketConnected = resolve;\n      rejectWebsocketConnected = reject;\n    });\n\n    const connectTimeout = setTimeout(() => {\n      this.excludeGateways.add(startResp.gatewayGroup);\n      rejectWebsocketConnected?.(\n        new ReconnectError(`Connection ${connectionId} timed out`, attempt),\n      );\n    }, 10_000);\n\n    const finalEndpoint = this.config.gatewayUrl || startResp.gatewayEndpoint;\n    if (finalEndpoint !== startResp.gatewayEndpoint) {\n      this.callbacks.logger.debug(\n        { original: startResp.gatewayEndpoint, override: finalEndpoint },\n        \"Overriding gateway endpoint\",\n      );\n    }\n\n    this.callbacks.logger.debug(\n      {\n        endpoint: finalEndpoint,\n        gatewayGroup: startResp.gatewayGroup,\n        connectionId,\n      },\n      \"Connecting to gateway\",\n    );\n\n    const ws = new WebSocket(finalEndpoint, [ConnectWebSocketProtocol]);\n    ws.binaryType = \"arraybuffer\";\n\n    let onConnectionError: (error: unknown) => void | Promise<void>;\n    {\n      onConnectionError = (error: unknown) => {\n        if (closed) {\n          this.callbacks.logger.debug(\n            { connectionId },\n            \"Connection error while initializing but already in closed state, skipping\",\n          );\n          return;\n        }\n        closed = true;\n\n        this.callbacks.logger.debug(\n          { connectionId },\n          \"Connection error in connecting state, rejecting promise\",\n        );\n\n        this.excludeGateways.add(startResp.gatewayGroup);\n        clearTimeout(connectTimeout);\n\n        ws.onerror = () => {};\n        ws.onclose = () => {};\n        ws.close(\n          4001,\n          workerDisconnectReasonToJSON(WorkerDisconnectReason.UNEXPECTED),\n        );\n\n        rejectWebsocketConnected?.(\n          new ReconnectError(\n            `Error while connecting (${connectionId}): ${\n              error instanceof Error ? error.message : \"Unknown error\"\n            }`,\n            attempt,\n          ),\n        );\n      };\n\n      ws.onerror = (err) => onConnectionError(err);\n      ws.onclose = (ev) => {\n        void onConnectionError(\n          new ReconnectError(\n            `Connection ${connectionId} closed: ${ev.reason}`,\n            attempt,\n          ),\n        );\n      };\n    }\n\n    const setupState = {\n      receivedGatewayHello: false,\n      sentWorkerConnect: false,\n      receivedConnectionReady: false,\n    };\n\n    let heartbeatIntervalMs: number | undefined;\n    let extendLeaseIntervalMs: number | undefined;\n\n    ws.onmessage = async (event) => {\n      const messageBytes = new Uint8Array(event.data as ArrayBuffer);\n      const connectMessage = parseConnectMessage(messageBytes);\n\n      this.callbacks.logger.debug(\n        { kind: gatewayMessageTypeToJSON(connectMessage.kind), connectionId },\n        \"Received message\",\n      );\n\n      if (!setupState.receivedGatewayHello) {\n        if (connectMessage.kind !== GatewayMessageType.GATEWAY_HELLO) {\n          void onConnectionError(\n            new ReconnectError(\n              `Expected hello message, got ${gatewayMessageTypeToJSON(\n                connectMessage.kind,\n              )}`,\n              attempt,\n            ),\n          );\n          return;\n        }\n        setupState.receivedGatewayHello = true;\n      }\n\n      if (!setupState.sentWorkerConnect) {\n        const workerConnectRequestMsg = WorkerConnectRequestData.create({\n          connectionId: startResp.connectionId,\n          environment: this.config.envName,\n          platform: getPlatformName({ ...allProcessEnv() }),\n          sdkVersion: `v${version}`,\n          sdkLanguage: \"typescript\",\n          framework: \"connect\",\n          workerManualReadinessAck:\n            this.config.connectionData.manualReadinessAck,\n          systemAttributes: await retrieveSystemAttributes(),\n          authData: {\n            sessionToken: startResp.sessionToken,\n            syncToken: startResp.syncToken,\n          },\n          apps: this.config.connectionData.apps,\n          capabilities: new TextEncoder().encode(\n            this.config.connectionData.marshaledCapabilities,\n          ),\n          startedAt: startedAt,\n          instanceId: this.config.instanceId || (await getHostname()),\n          maxWorkerConcurrency: this.config.maxWorkerConcurrency,\n        });\n\n        const workerConnectRequestMsgBytes = WorkerConnectRequestData.encode(\n          workerConnectRequestMsg,\n        ).finish();\n\n        ws.send(\n          ensureUnsharedArrayBuffer(\n            ConnectMessage.encode(\n              ConnectMessage.create({\n                kind: GatewayMessageType.WORKER_CONNECT,\n                payload: workerConnectRequestMsgBytes,\n              }),\n            ).finish(),\n          ),\n        );\n\n        setupState.sentWorkerConnect = true;\n        return;\n      }\n\n      if (!setupState.receivedConnectionReady) {\n        if (\n          connectMessage.kind !== GatewayMessageType.GATEWAY_CONNECTION_READY\n        ) {\n          void onConnectionError(\n            new ReconnectError(\n              `Expected ready message, got ${gatewayMessageTypeToJSON(\n                connectMessage.kind,\n              )}`,\n              attempt,\n            ),\n          );\n          return;\n        }\n\n        const readyPayload = GatewayConnectionReadyData.decode(\n          connectMessage.payload,\n        );\n\n        setupState.receivedConnectionReady = true;\n\n        heartbeatIntervalMs =\n          readyPayload.heartbeatInterval.length > 0\n            ? ms(readyPayload.heartbeatInterval as ms.StringValue)\n            : 10_000;\n        extendLeaseIntervalMs =\n          readyPayload.extendLeaseInterval.length > 0\n            ? ms(readyPayload.extendLeaseInterval as ms.StringValue)\n            : 5_000;\n\n        resolveWebsocketConnected?.();\n        return;\n      }\n\n      this.callbacks.logger.warn(\n        {\n          kind: gatewayMessageTypeToJSON(connectMessage.kind),\n          rawKind: connectMessage.kind,\n          attempt,\n          setupState,\n          state: this.callbacks.getState(),\n          connectionId,\n        },\n        \"Unexpected message type during setup\",\n      );\n    };\n\n    await websocketConnectedPromise;\n\n    clearTimeout(connectTimeout);\n\n    this.excludeGateways.delete(startResp.gatewayGroup);\n\n    attempt = 0;\n\n    const conn: Connection = {\n      id: connectionId,\n      ws,\n      cleanup: () => {\n        if (closed) {\n          return;\n        }\n        closed = true;\n        ws.onerror = () => {};\n        ws.onclose = () => {};\n        ws.close();\n      },\n      pendingHeartbeats: 0,\n    };\n    this.currentConnection = conn;\n\n    // Set state to ACTIVE after currentConnection is set, so that\n    // connectionId is available in the onStateChange callback.\n    this.callbacks.onStateChange(ConnectionState.ACTIVE);\n\n    let isDraining = false;\n    {\n      onConnectionError = async (error: unknown) => {\n        if (closed) {\n          this.callbacks.logger.debug(\n            { connectionId },\n            \"Connection error but already in closed state, skipping\",\n          );\n          return;\n        }\n        closed = true;\n\n        await conn.cleanup();\n\n        const currentState = this.callbacks.getState();\n        if (\n          currentState === ConnectionState.CLOSING ||\n          currentState === ConnectionState.CLOSED\n        ) {\n          this.callbacks.logger.debug(\n            { connectionId },\n            \"Connection error but already closing or closed, skipping\",\n          );\n          return;\n        }\n\n        this.callbacks.onStateChange(ConnectionState.RECONNECTING);\n        this.excludeGateways.add(startResp.gatewayGroup);\n\n        if (isDraining) {\n          this.callbacks.logger.debug(\n            { connectionId },\n            \"Connection error but already draining, skipping\",\n          );\n          return;\n        }\n\n        this.callbacks.logger.warn(\n          {\n            connectionId,\n            err: toError(error),\n          },\n          \"Connection error\",\n        );\n        this.connect(attempt + 1, [...path, \"onConnectionError\"]);\n      };\n\n      ws.onerror = (err) => onConnectionError(err);\n      ws.onclose = (ev) => {\n        void onConnectionError(\n          new ReconnectError(`Connection closed: ${ev.reason}`, attempt),\n        );\n      };\n    }\n\n    ws.onmessage = async (event) => {\n      const messageBytes = new Uint8Array(event.data as ArrayBuffer);\n      const connectMessage = parseConnectMessage(messageBytes);\n\n      if (connectMessage.kind === GatewayMessageType.GATEWAY_CLOSING) {\n        isDraining = true;\n        this.callbacks.logger.info(\n          { connectionId },\n          \"Received draining message\",\n        );\n        try {\n          this.callbacks.logger.debug(\n            { connectionId },\n            \"Setting up new connection while keeping previous connection open\",\n          );\n\n          await this.connect(0, [...path]);\n          await conn.cleanup();\n        } catch (err) {\n          this.callbacks.logger.warn(\n            {\n              connectionId,\n              err: toError(err),\n            },\n            \"Failed to reconnect after receiving draining message\",\n          );\n\n          await conn.cleanup();\n\n          void onConnectionError(\n            new ReconnectError(\n              `Failed to reconnect after receiving draining message (${connectionId})`,\n              attempt,\n            ),\n          );\n        }\n        return;\n      }\n\n      if (connectMessage.kind === GatewayMessageType.GATEWAY_HEARTBEAT) {\n        conn.pendingHeartbeats = 0;\n        this.callbacks.logger.debug(\n          { connectionId },\n          \"Handled gateway heartbeat\",\n        );\n        return;\n      }\n\n      if (connectMessage.kind === GatewayMessageType.GATEWAY_EXECUTOR_REQUEST) {\n        const currentState = this.callbacks.getState();\n        if (currentState !== ConnectionState.ACTIVE) {\n          this.callbacks.logger.warn(\n            { connectionId },\n            \"Received request while not active, skipping\",\n          );\n          return;\n        }\n\n        const gatewayExecutorRequest = parseGatewayExecutorRequest(\n          connectMessage.payload,\n        );\n\n        this.callbacks.logger.debug(\n          {\n            requestId: gatewayExecutorRequest.requestId,\n            appId: gatewayExecutorRequest.appId,\n            appName: gatewayExecutorRequest.appName,\n            functionSlug: gatewayExecutorRequest.functionSlug,\n            stepId: gatewayExecutorRequest.stepId,\n            connectionId,\n          },\n          \"Received gateway executor request\",\n        );\n\n        if (\n          typeof gatewayExecutorRequest.appName !== \"string\" ||\n          gatewayExecutorRequest.appName.length === 0\n        ) {\n          this.callbacks.logger.warn(\n            {\n              requestId: gatewayExecutorRequest.requestId,\n              appId: gatewayExecutorRequest.appId,\n              functionSlug: gatewayExecutorRequest.functionSlug,\n              stepId: gatewayExecutorRequest.stepId,\n              connectionId,\n            },\n            \"No app name in request, skipping\",\n          );\n          return;\n        }\n\n        if (!this.config.appIds.includes(gatewayExecutorRequest.appName)) {\n          this.callbacks.logger.warn(\n            {\n              requestId: gatewayExecutorRequest.requestId,\n              appId: gatewayExecutorRequest.appId,\n              appName: gatewayExecutorRequest.appName,\n              functionSlug: gatewayExecutorRequest.functionSlug,\n              stepId: gatewayExecutorRequest.stepId,\n              connectionId,\n            },\n            \"No request handler found for app, skipping\",\n          );\n          return;\n        }\n\n        // Send ACK\n        ws.send(\n          ensureUnsharedArrayBuffer(\n            ConnectMessage.encode(\n              ConnectMessage.create({\n                kind: GatewayMessageType.WORKER_REQUEST_ACK,\n                payload: WorkerRequestAckData.encode(\n                  WorkerRequestAckData.create({\n                    accountId: gatewayExecutorRequest.accountId,\n                    envId: gatewayExecutorRequest.envId,\n                    appId: gatewayExecutorRequest.appId,\n                    functionSlug: gatewayExecutorRequest.functionSlug,\n                    requestId: gatewayExecutorRequest.requestId,\n                    stepId: gatewayExecutorRequest.stepId,\n                    userTraceCtx: gatewayExecutorRequest.userTraceCtx,\n                    systemTraceCtx: gatewayExecutorRequest.systemTraceCtx,\n                    runId: gatewayExecutorRequest.runId,\n                  }),\n                ).finish(),\n              }),\n            ).finish(),\n          ),\n        );\n\n        this.inProgressRequests.wg.add(1);\n        this.inProgressRequests.requestLeases[\n          gatewayExecutorRequest.requestId\n        ] = gatewayExecutorRequest.leaseId;\n\n        // Start lease extension interval\n        let extendLeaseInterval: NodeJS.Timeout | undefined;\n        extendLeaseInterval = setInterval(() => {\n          if (extendLeaseIntervalMs === undefined) {\n            return;\n          }\n\n          const currentLeaseId =\n            this.inProgressRequests.requestLeases[\n              gatewayExecutorRequest.requestId\n            ];\n          if (!currentLeaseId) {\n            clearInterval(extendLeaseInterval);\n            return;\n          }\n\n          // Use the current live connection's WebSocket for lease\n          // extensions. During a drain, the original WebSocket may be\n          // closed by the gateway while the request is still in flight,\n          // causing lease extension messages to silently fail and the\n          // gateway to time out the request.\n          const latestConn = {\n            ws: this.currentConnection?.ws ?? ws,\n            id: this.currentConnection?.id ?? connectionId,\n          };\n\n          this.callbacks.logger.debug(\n            { connectionId: latestConn.id, leaseId: currentLeaseId },\n            \"Extending lease\",\n          );\n\n          if (latestConn.ws.readyState !== WebSocket.OPEN) {\n            this.callbacks.logger.warn(\n              {\n                connectionId: latestConn.id,\n                requestId: gatewayExecutorRequest.requestId,\n              },\n              \"Cannot extend lease, no open WebSocket available\",\n            );\n            return;\n          }\n\n          latestConn.ws.send(\n            ensureUnsharedArrayBuffer(\n              ConnectMessage.encode(\n                ConnectMessage.create({\n                  kind: GatewayMessageType.WORKER_REQUEST_EXTEND_LEASE,\n                  payload: WorkerRequestExtendLeaseData.encode(\n                    WorkerRequestExtendLeaseData.create({\n                      accountId: gatewayExecutorRequest.accountId,\n                      envId: gatewayExecutorRequest.envId,\n                      appId: gatewayExecutorRequest.appId,\n                      functionSlug: gatewayExecutorRequest.functionSlug,\n                      requestId: gatewayExecutorRequest.requestId,\n                      stepId: gatewayExecutorRequest.stepId,\n                      runId: gatewayExecutorRequest.runId,\n                      userTraceCtx: gatewayExecutorRequest.userTraceCtx,\n                      systemTraceCtx: gatewayExecutorRequest.systemTraceCtx,\n                      leaseId: currentLeaseId,\n                    }),\n                  ).finish(),\n                }),\n              ).finish(),\n            ),\n          );\n        }, extendLeaseIntervalMs);\n\n        try {\n          // Handle execution via callback\n          const responseBytes = await this.callbacks.handleExecutionRequest(\n            gatewayExecutorRequest,\n          );\n\n          if (!this.currentConnection) {\n            this.callbacks.logger.warn(\n              { requestId: gatewayExecutorRequest.requestId },\n              \"No current WebSocket, buffering response\",\n            );\n            if (this.callbacks.onBufferResponse) {\n              this.callbacks.onBufferResponse(\n                gatewayExecutorRequest.requestId,\n                responseBytes,\n              );\n            }\n            return;\n          }\n\n          this.callbacks.logger.debug(\n            {\n              connectionId: this.currentConnection.id,\n              requestId: gatewayExecutorRequest.requestId,\n            },\n            \"Sending worker reply\",\n          );\n\n          this.currentConnection.ws.send(\n            ensureUnsharedArrayBuffer(\n              ConnectMessage.encode(\n                ConnectMessage.create({\n                  kind: GatewayMessageType.WORKER_REPLY,\n                  payload: responseBytes,\n                }),\n              ).finish(),\n            ),\n          );\n        } catch (err) {\n          this.callbacks.logger.debug(\n            {\n              requestId: gatewayExecutorRequest.requestId,\n              err: toError(err),\n            },\n            \"Execution error\",\n          );\n        } finally {\n          this.inProgressRequests.wg.done();\n          delete this.inProgressRequests.requestLeases[\n            gatewayExecutorRequest.requestId\n          ];\n          clearInterval(extendLeaseInterval);\n        }\n\n        return;\n      }\n\n      if (connectMessage.kind === GatewayMessageType.WORKER_REPLY_ACK) {\n        const replyAck = parseWorkerReplyAck(connectMessage.payload);\n\n        this.callbacks.logger.debug(\n          { connectionId, requestId: replyAck.requestId },\n          \"Acknowledging reply ack\",\n        );\n\n        this.callbacks.onReplyAck?.(replyAck.requestId);\n        return;\n      }\n\n      if (\n        connectMessage.kind ===\n        GatewayMessageType.WORKER_REQUEST_EXTEND_LEASE_ACK\n      ) {\n        const extendLeaseAck = WorkerRequestExtendLeaseAckData.decode(\n          connectMessage.payload,\n        );\n\n        this.callbacks.logger.debug(\n          { connectionId, newLeaseId: extendLeaseAck.newLeaseId },\n          \"Received extend lease ack\",\n        );\n\n        if (extendLeaseAck.newLeaseId) {\n          this.inProgressRequests.requestLeases[extendLeaseAck.requestId] =\n            extendLeaseAck.newLeaseId;\n        } else {\n          this.callbacks.logger.warn(\n            { connectionId, requestId: extendLeaseAck.requestId },\n            \"Unable to extend lease\",\n          );\n          delete this.inProgressRequests.requestLeases[\n            extendLeaseAck.requestId\n          ];\n        }\n\n        return;\n      }\n\n      this.callbacks.logger.warn(\n        {\n          kind: gatewayMessageTypeToJSON(connectMessage.kind),\n          rawKind: connectMessage.kind,\n          attempt,\n          setupState,\n          state: this.callbacks.getState(),\n          connectionId,\n        },\n        \"Unexpected message type\",\n      );\n    };\n\n    // Heartbeat interval\n    let heartbeatInterval: NodeJS.Timeout | undefined;\n    if (heartbeatIntervalMs !== undefined) {\n      heartbeatInterval = setInterval(() => {\n        if (heartbeatIntervalMs === undefined) {\n          return;\n        }\n\n        // Skip heartbeat ticks when the WebSocket is no longer open. During\n        // drain the Gateway may close the old WS while in-flight requests are\n        // still running. Sending on a closed socket is a no-op and we must not\n        // treat the missing response as a failure.\n        //\n        // This is safe because each connection gets its own heartbeat\n        // interval. Once we reconnect, we can safely skip the old WS\n        // heartbeats because the new WS is heartbeating.\n        //\n        // TODO: We need a better way to handle this. This isn't a horrible\n        // hack, but it isn't ideal. Over the life of a worker, it'll have N-1\n        // noop heartbeat intervals, where N is the number of times it\n        // reconnected.\n        if (ws.readyState !== WebSocket.OPEN) {\n          return;\n        }\n\n        if (conn.pendingHeartbeats >= 2) {\n          this.callbacks.logger.warn(\n            { connectionId },\n            \"Gateway heartbeat missed\",\n          );\n          void onConnectionError(\n            new ReconnectError(\n              `Consecutive gateway heartbeats missed (${connectionId})`,\n              attempt,\n            ),\n          );\n          return;\n        }\n\n        this.callbacks.logger.debug(\n          { connectionId },\n          \"Sending worker heartbeat\",\n        );\n\n        conn.pendingHeartbeats++;\n        ws.send(\n          ensureUnsharedArrayBuffer(\n            ConnectMessage.encode(\n              ConnectMessage.create({\n                kind: GatewayMessageType.WORKER_HEARTBEAT,\n              }),\n            ).finish(),\n          ),\n        );\n      }, heartbeatIntervalMs);\n    }\n\n    conn.cleanup = async () => {\n      if (closed) {\n        return;\n      }\n      closed = true;\n\n      this.callbacks.logger.debug({ connectionId }, \"Cleaning up connection\");\n      if (ws.readyState === WebSocket.OPEN) {\n        this.callbacks.logger.debug({ connectionId }, \"Sending pause message\");\n        ws.send(\n          ensureUnsharedArrayBuffer(\n            ConnectMessage.encode(\n              ConnectMessage.create({\n                kind: GatewayMessageType.WORKER_PAUSE,\n              }),\n            ).finish(),\n          ),\n        );\n      }\n\n      this.callbacks.logger.debug({ connectionId }, \"Closing connection\");\n      ws.onerror = () => {};\n      ws.onclose = () => {};\n\n      await this.inProgressRequests.wg.wait();\n\n      ws.close(\n        1000,\n        workerDisconnectReasonToJSON(WorkerDisconnectReason.WORKER_SHUTDOWN),\n      );\n\n      if (this.currentConnection?.id === connectionId) {\n        this.currentConnection = undefined;\n      }\n\n      this.callbacks.logger.debug(\n        { connectionId },\n        \"Cleaning up worker heartbeat\",\n      );\n      clearInterval(heartbeatInterval);\n    };\n\n    return conn;\n  }\n\n  async getApiBaseUrl(): Promise<string> {\n    return resolveApiBaseUrl({\n      apiBaseUrl: this.config.apiBaseUrl,\n      mode: this.config.mode,\n    });\n  }\n}\n"],"mappings":";;;;;;;;;;;;;;;;;;;;;;;AA+CA,MAAM,2BAA2B;AAEjC,SAAS,QAAQ,OAAuB;AACtC,KAAI,iBAAiB,MACnB,QAAO;AAET,QAAO,IAAI,MAAM,OAAO,MAAM,CAAC;;;;;;AA2CjC,IAAa,iBAAb,MAA4B;CAC1B,AAAQ;CACR,AAAQ;CAER,AAAQ;CACR,AAAQ,kCAA+B,IAAI,KAAK;CAEhD,AAAQ,qBAGJ;EACF,IAAI,IAAIA,kCAAW;EACnB,eAAe,EAAE;EAClB;CAED,YACE,QACA,WACA;AACA,OAAK,SAAS;AACd,OAAK,YAAY;;CAGnB,IAAI,aAAqC;AACvC,SAAO,KAAK;;CAGd,IAAI,eAAmC;AACrC,SAAO,KAAK,mBAAmB;;;;;CAMjC,MAAM,oBAAmC;AACvC,QAAM,KAAK,mBAAmB,GAAG,MAAM;;;;;CAMzC,MAAM,QAAQ,UAAU,GAAG,OAAiB,EAAE,EAAiB;AAC7D,MAAI,OAAO,cAAc,YACvB,OAAM,IAAI,MAAM,kDAAkD;AAIpE,MADc,KAAK,UAAU,UAAU,KACzBC,8BAAgB,OAC5B,OAAM,IAAI,MAAM,4BAA4B;AAG9C,OAAK,UAAU,OAAO,MAAM,EAAE,SAAS,EAAE,0BAA0B;EAEnE,IAAI,gBAAgB,KAAK,OAAO;AAEhC,SAAO,MAAM;AAEX,OADqB,KAAK,UAAU,UAAU,KACzBA,8BAAgB,OACnC;AAgBF,OAAI,KAAK,UAAU,cACjB,OAAM,KAAK,UAAU,cAAc,cAAc;AAGnD,OAAI;AACF,UAAM,KAAK,kBAAkB,eAAe,SAAS,CAAC,GAAG,KAAK,CAAC;AAC/D;YACO,KAAK;AACZ,SAAK,UAAU,OAAO,KAAK,EAAE,KAAK,QAAQ,IAAI,EAAE,EAAE,oBAAoB;AAEtE,QAAI,EAAE,eAAeC,6BACnB,OAAM;AAGR,cAAU,IAAI;AAEd,QAAI,eAAeC,wBAAW;KAC5B,MAAM,mBACJ,kBAAkB,KAAK,OAAO;AAChC,SAAI,iBACF,MAAK,UAAU,OAAO,MAAM,oCAAoC;AAElE,qBAAgB,mBACZ,KAAK,OAAO,oBACZ,KAAK,OAAO;;AAGlB,QAAI,eAAeC,kCACjB,MAAK,UAAU,OAAO,MACpB,qHACD;IAGH,MAAM,QAAQC,wBAAW,QAAQ;AACjC,SAAK,UAAU,OAAO,MAAM,EAAE,OAAO,EAAE,eAAe;AAKtD,QAHkB,MAAMC,4BAAe,aAAa;AAClD,YAAO,KAAK,UAAU,UAAU,KAAKL,8BAAgB;MACrD,EACa;AACb,UAAK,UAAU,OAAO,MAAM,8BAA8B;AAC1D;;AAGF;;;AAIJ,OAAK,UAAU,OAAO,MAAM,uBAAuB;;;;;CAMrD,MAAM,UAAyB;AAC7B,MAAI,KAAK,mBAAmB;AAC1B,SAAM,KAAK,kBAAkB,SAAS;AACtC,QAAK,oBAAoB;;;CAI7B,MAAc,iBACZ,kBACA,SACA;EACA,MAAM,MAAMM,oCAAmB,MAAM,KAAK,KAAK,gBAAgB,CAAC;EAEhE,MAAMC,UAAkC;GACtC,gBAAgB;GAChB,GAAI,mBACA,EAAE,eAAe,UAAU,oBAAoB,GAC/C,EAAE;GACP;AAED,MAAI,KAAK,OAAO,QACd,SAAQC,0BAAW,eAAe,KAAK,OAAO;EAGhD,MAAM,YAAY,IAAI,IAAI,qBAAqB,MAAM,KAAK,eAAe,CAAC;EAE1E,IAAI;AACJ,MAAI;AACF,UAAO,MAAM,MAAM,WAAW;IAC5B,QAAQ;IACR,MAAM,IAAI,WAAW,IAAI;IAChB;IACV,CAAC;WACK,KAAK;GACZ,MAAM,SAAS,eAAe,QAAQ,IAAI,UAAU;AACpD,SAAM,IAAIP,4BACR,2CAA2C,UAAU,UAAU,CAAC,IAAI,UACpE,QACD;;AAGH,MAAI,CAAC,KAAK,IAAI;AACZ,OAAI,KAAK,WAAW,IAClB,OAAM,IAAIC,uBACR,2CAA2C,UAAU,UAAU,GAC7D,KAAK,OAAO,UAAU,UAAU,KAAK,OAAO,QAAQ,KAAK,GAC1D,IAAI,MAAM,KAAK,MAAM,IACtB,QACD;AAGH,OAAI,KAAK,WAAW,IAClB,OAAM,IAAIC,kCAAqB,QAAQ;AAGzC,SAAM,IAAIF,4BACR,2CAA2C,UAAU,UAAU,CAAC,IAAI,MAAM,KAAK,MAAM,IACrF,QACD;;AAIH,SADkB,MAAMQ,oCAAmB,KAAK;;CAIlD,MAAc,kBACZ,kBACA,SACA,OAAiB,EAAE,EACE;EACrB,IAAI,SAAS;AAEb,OAAK,UAAU,OAAO,MAAM;GAAE;GAAS;GAAM,EAAE,uBAAuB;EAEtE,MAAM,4BAAY,IAAI,MAAM;EAC5B,MAAM,YAAY,MAAM,KAAK,iBAAiB,kBAAkB,QAAQ;EAExE,MAAM,eAAe,UAAU;AAC/B,OAAK,KAAK,aAAa;EAEvB,IAAIC;EAGJ,IAAIC;EACJ,MAAM,4BAA4B,IAAI,SAAe,SAAS,WAAW;AACvE,+BAA4B;AAC5B,8BAA2B;IAC3B;EAEF,MAAM,iBAAiB,iBAAiB;AACtC,QAAK,gBAAgB,IAAI,UAAU,aAAa;AAChD,8BACE,IAAIV,4BAAe,cAAc,aAAa,aAAa,QAAQ,CACpE;KACA,IAAO;EAEV,MAAM,gBAAgB,KAAK,OAAO,cAAc,UAAU;AAC1D,MAAI,kBAAkB,UAAU,gBAC9B,MAAK,UAAU,OAAO,MACpB;GAAE,UAAU,UAAU;GAAiB,UAAU;GAAe,EAChE,8BACD;AAGH,OAAK,UAAU,OAAO,MACpB;GACE,UAAU;GACV,cAAc,UAAU;GACxB;GACD,EACD,wBACD;EAED,MAAM,KAAK,IAAI,UAAU,eAAe,CAAC,yBAAyB,CAAC;AACnE,KAAG,aAAa;EAEhB,IAAIW,qBAEmB,UAAmB;AACtC,OAAI,QAAQ;AACV,SAAK,UAAU,OAAO,MACpB,EAAE,cAAc,EAChB,4EACD;AACD;;AAEF,YAAS;AAET,QAAK,UAAU,OAAO,MACpB,EAAE,cAAc,EAChB,0DACD;AAED,QAAK,gBAAgB,IAAI,UAAU,aAAa;AAChD,gBAAa,eAAe;AAE5B,MAAG,gBAAgB;AACnB,MAAG,gBAAgB;AACnB,MAAG,MACD,MACAC,6CAA6BC,uCAAuB,WAAW,CAChE;AAED,8BACE,IAAIb,4BACF,2BAA2B,aAAa,KACtC,iBAAiB,QAAQ,MAAM,UAAU,mBAE3C,QACD,CACF;;AAGH,KAAG,WAAW,QAAQ,kBAAkB,IAAI;AAC5C,KAAG,WAAW,OAAO;AACnB,GAAK,kBACH,IAAIA,4BACF,cAAc,aAAa,WAAW,GAAG,UACzC,QACD,CACF;;EAIL,MAAM,aAAa;GACjB,sBAAsB;GACtB,mBAAmB;GACnB,yBAAyB;GAC1B;EAED,IAAIc;EACJ,IAAIC;AAEJ,KAAG,YAAY,OAAO,UAAU;GAE9B,MAAM,iBAAiBC,qCADF,IAAI,WAAW,MAAM,KAAoB,CACN;AAExD,QAAK,UAAU,OAAO,MACpB;IAAE,MAAMC,yCAAyB,eAAe,KAAK;IAAE;IAAc,EACrE,mBACD;AAED,OAAI,CAAC,WAAW,sBAAsB;AACpC,QAAI,eAAe,SAASC,mCAAmB,eAAe;AAC5D,KAAK,kBACH,IAAIlB,4BACF,+BAA+BiB,yCAC7B,eAAe,KAChB,IACD,QACD,CACF;AACD;;AAEF,eAAW,uBAAuB;;AAGpC,OAAI,CAAC,WAAW,mBAAmB;IACjC,MAAM,0BAA0BE,yCAAyB,OAAO;KAC9D,cAAc,UAAU;KACxB,aAAa,KAAK,OAAO;KACzB,UAAUC,4BAAgB,EAAE,GAAGC,2BAAe,EAAE,CAAC;KACjD,YAAY,IAAIC;KAChB,aAAa;KACb,WAAW;KACX,0BACE,KAAK,OAAO,eAAe;KAC7B,kBAAkB,MAAMC,qCAA0B;KAClD,UAAU;MACR,cAAc,UAAU;MACxB,WAAW,UAAU;MACtB;KACD,MAAM,KAAK,OAAO,eAAe;KACjC,cAAc,IAAI,aAAa,CAAC,OAC9B,KAAK,OAAO,eAAe,sBAC5B;KACU;KACX,YAAY,KAAK,OAAO,cAAe,MAAMC,wBAAa;KAC1D,sBAAsB,KAAK,OAAO;KACnC,CAAC;IAEF,MAAM,+BAA+BL,yCAAyB,OAC5D,wBACD,CAAC,QAAQ;AAEV,OAAG,KACDM,yCACEC,+BAAe,OACbA,+BAAe,OAAO;KACpB,MAAMR,mCAAmB;KACzB,SAAS;KACV,CAAC,CACH,CAAC,QAAQ,CACX,CACF;AAED,eAAW,oBAAoB;AAC/B;;AAGF,OAAI,CAAC,WAAW,yBAAyB;AACvC,QACE,eAAe,SAASA,mCAAmB,0BAC3C;AACA,KAAK,kBACH,IAAIlB,4BACF,+BAA+BiB,yCAC7B,eAAe,KAChB,IACD,QACD,CACF;AACD;;IAGF,MAAM,eAAeU,2CAA2B,OAC9C,eAAe,QAChB;AAED,eAAW,0BAA0B;AAErC,0BACE,aAAa,kBAAkB,SAAS,oBACjC,aAAa,kBAAoC,GACpD;AACN,4BACE,aAAa,oBAAoB,SAAS,oBACnC,aAAa,oBAAsC,GACtD;AAEN,iCAA6B;AAC7B;;AAGF,QAAK,UAAU,OAAO,KACpB;IACE,MAAMV,yCAAyB,eAAe,KAAK;IACnD,SAAS,eAAe;IACxB;IACA;IACA,OAAO,KAAK,UAAU,UAAU;IAChC;IACD,EACD,uCACD;;AAGH,QAAM;AAEN,eAAa,eAAe;AAE5B,OAAK,gBAAgB,OAAO,UAAU,aAAa;AAEnD,YAAU;EAEV,MAAMW,OAAmB;GACvB,IAAI;GACJ;GACA,eAAe;AACb,QAAI,OACF;AAEF,aAAS;AACT,OAAG,gBAAgB;AACnB,OAAG,gBAAgB;AACnB,OAAG,OAAO;;GAEZ,mBAAmB;GACpB;AACD,OAAK,oBAAoB;AAIzB,OAAK,UAAU,cAAc7B,8BAAgB,OAAO;EAEpD,IAAI,aAAa;AAEf,sBAAoB,OAAO,UAAmB;AAC5C,OAAI,QAAQ;AACV,SAAK,UAAU,OAAO,MACpB,EAAE,cAAc,EAChB,yDACD;AACD;;AAEF,YAAS;AAET,SAAM,KAAK,SAAS;GAEpB,MAAM,eAAe,KAAK,UAAU,UAAU;AAC9C,OACE,iBAAiBA,8BAAgB,WACjC,iBAAiBA,8BAAgB,QACjC;AACA,SAAK,UAAU,OAAO,MACpB,EAAE,cAAc,EAChB,2DACD;AACD;;AAGF,QAAK,UAAU,cAAcA,8BAAgB,aAAa;AAC1D,QAAK,gBAAgB,IAAI,UAAU,aAAa;AAEhD,OAAI,YAAY;AACd,SAAK,UAAU,OAAO,MACpB,EAAE,cAAc,EAChB,kDACD;AACD;;AAGF,QAAK,UAAU,OAAO,KACpB;IACE;IACA,KAAK,QAAQ,MAAM;IACpB,EACD,mBACD;AACD,QAAK,QAAQ,UAAU,GAAG,CAAC,GAAG,MAAM,oBAAoB,CAAC;;AAG3D,KAAG,WAAW,QAAQ,kBAAkB,IAAI;AAC5C,KAAG,WAAW,OAAO;AACnB,GAAK,kBACH,IAAIC,4BAAe,sBAAsB,GAAG,UAAU,QAAQ,CAC/D;;AAIL,KAAG,YAAY,OAAO,UAAU;GAE9B,MAAM,iBAAiBgB,qCADF,IAAI,WAAW,MAAM,KAAoB,CACN;AAExD,OAAI,eAAe,SAASE,mCAAmB,iBAAiB;AAC9D,iBAAa;AACb,SAAK,UAAU,OAAO,KACpB,EAAE,cAAc,EAChB,4BACD;AACD,QAAI;AACF,UAAK,UAAU,OAAO,MACpB,EAAE,cAAc,EAChB,mEACD;AAED,WAAM,KAAK,QAAQ,GAAG,CAAC,GAAG,KAAK,CAAC;AAChC,WAAM,KAAK,SAAS;aACb,KAAK;AACZ,UAAK,UAAU,OAAO,KACpB;MACE;MACA,KAAK,QAAQ,IAAI;MAClB,EACD,uDACD;AAED,WAAM,KAAK,SAAS;AAEpB,KAAK,kBACH,IAAIlB,4BACF,yDAAyD,aAAa,IACtE,QACD,CACF;;AAEH;;AAGF,OAAI,eAAe,SAASkB,mCAAmB,mBAAmB;AAChE,SAAK,oBAAoB;AACzB,SAAK,UAAU,OAAO,MACpB,EAAE,cAAc,EAChB,4BACD;AACD;;AAGF,OAAI,eAAe,SAASA,mCAAmB,0BAA0B;AAEvE,QADqB,KAAK,UAAU,UAAU,KACzBnB,8BAAgB,QAAQ;AAC3C,UAAK,UAAU,OAAO,KACpB,EAAE,cAAc,EAChB,8CACD;AACD;;IAGF,MAAM,yBAAyB8B,6CAC7B,eAAe,QAChB;AAED,SAAK,UAAU,OAAO,MACpB;KACE,WAAW,uBAAuB;KAClC,OAAO,uBAAuB;KAC9B,SAAS,uBAAuB;KAChC,cAAc,uBAAuB;KACrC,QAAQ,uBAAuB;KAC/B;KACD,EACD,oCACD;AAED,QACE,OAAO,uBAAuB,YAAY,YAC1C,uBAAuB,QAAQ,WAAW,GAC1C;AACA,UAAK,UAAU,OAAO,KACpB;MACE,WAAW,uBAAuB;MAClC,OAAO,uBAAuB;MAC9B,cAAc,uBAAuB;MACrC,QAAQ,uBAAuB;MAC/B;MACD,EACD,mCACD;AACD;;AAGF,QAAI,CAAC,KAAK,OAAO,OAAO,SAAS,uBAAuB,QAAQ,EAAE;AAChE,UAAK,UAAU,OAAO,KACpB;MACE,WAAW,uBAAuB;MAClC,OAAO,uBAAuB;MAC9B,SAAS,uBAAuB;MAChC,cAAc,uBAAuB;MACrC,QAAQ,uBAAuB;MAC/B;MACD,EACD,6CACD;AACD;;AAIF,OAAG,KACDJ,yCACEC,+BAAe,OACbA,+BAAe,OAAO;KACpB,MAAMR,mCAAmB;KACzB,SAASY,qCAAqB,OAC5BA,qCAAqB,OAAO;MAC1B,WAAW,uBAAuB;MAClC,OAAO,uBAAuB;MAC9B,OAAO,uBAAuB;MAC9B,cAAc,uBAAuB;MACrC,WAAW,uBAAuB;MAClC,QAAQ,uBAAuB;MAC/B,cAAc,uBAAuB;MACrC,gBAAgB,uBAAuB;MACvC,OAAO,uBAAuB;MAC/B,CAAC,CACH,CAAC,QAAQ;KACX,CAAC,CACH,CAAC,QAAQ,CACX,CACF;AAED,SAAK,mBAAmB,GAAG,IAAI,EAAE;AACjC,SAAK,mBAAmB,cACtB,uBAAuB,aACrB,uBAAuB;IAG3B,IAAIC;AACJ,0BAAsB,kBAAkB;AACtC,SAAI,0BAA0B,OAC5B;KAGF,MAAM,iBACJ,KAAK,mBAAmB,cACtB,uBAAuB;AAE3B,SAAI,CAAC,gBAAgB;AACnB,oBAAc,oBAAoB;AAClC;;KAQF,MAAM,aAAa;MACjB,IAAI,KAAK,mBAAmB,MAAM;MAClC,IAAI,KAAK,mBAAmB,MAAM;MACnC;AAED,UAAK,UAAU,OAAO,MACpB;MAAE,cAAc,WAAW;MAAI,SAAS;MAAgB,EACxD,kBACD;AAED,SAAI,WAAW,GAAG,eAAe,UAAU,MAAM;AAC/C,WAAK,UAAU,OAAO,KACpB;OACE,cAAc,WAAW;OACzB,WAAW,uBAAuB;OACnC,EACD,mDACD;AACD;;AAGF,gBAAW,GAAG,KACZN,yCACEC,+BAAe,OACbA,+BAAe,OAAO;MACpB,MAAMR,mCAAmB;MACzB,SAASc,6CAA6B,OACpCA,6CAA6B,OAAO;OAClC,WAAW,uBAAuB;OAClC,OAAO,uBAAuB;OAC9B,OAAO,uBAAuB;OAC9B,cAAc,uBAAuB;OACrC,WAAW,uBAAuB;OAClC,QAAQ,uBAAuB;OAC/B,OAAO,uBAAuB;OAC9B,cAAc,uBAAuB;OACrC,gBAAgB,uBAAuB;OACvC,SAAS;OACV,CAAC,CACH,CAAC,QAAQ;MACX,CAAC,CACH,CAAC,QAAQ,CACX,CACF;OACA,sBAAsB;AAEzB,QAAI;KAEF,MAAM,gBAAgB,MAAM,KAAK,UAAU,uBACzC,uBACD;AAED,SAAI,CAAC,KAAK,mBAAmB;AAC3B,WAAK,UAAU,OAAO,KACpB,EAAE,WAAW,uBAAuB,WAAW,EAC/C,2CACD;AACD,UAAI,KAAK,UAAU,iBACjB,MAAK,UAAU,iBACb,uBAAuB,WACvB,cACD;AAEH;;AAGF,UAAK,UAAU,OAAO,MACpB;MACE,cAAc,KAAK,kBAAkB;MACrC,WAAW,uBAAuB;MACnC,EACD,uBACD;AAED,UAAK,kBAAkB,GAAG,KACxBP,yCACEC,+BAAe,OACbA,+BAAe,OAAO;MACpB,MAAMR,mCAAmB;MACzB,SAAS;MACV,CAAC,CACH,CAAC,QAAQ,CACX,CACF;aACM,KAAK;AACZ,UAAK,UAAU,OAAO,MACpB;MACE,WAAW,uBAAuB;MAClC,KAAK,QAAQ,IAAI;MAClB,EACD,kBACD;cACO;AACR,UAAK,mBAAmB,GAAG,MAAM;AACjC,YAAO,KAAK,mBAAmB,cAC7B,uBAAuB;AAEzB,mBAAc,oBAAoB;;AAGpC;;AAGF,OAAI,eAAe,SAASA,mCAAmB,kBAAkB;IAC/D,MAAM,WAAWe,qCAAoB,eAAe,QAAQ;AAE5D,SAAK,UAAU,OAAO,MACpB;KAAE;KAAc,WAAW,SAAS;KAAW,EAC/C,0BACD;AAED,SAAK,UAAU,aAAa,SAAS,UAAU;AAC/C;;AAGF,OACE,eAAe,SACff,mCAAmB,iCACnB;IACA,MAAM,iBAAiBgB,gDAAgC,OACrD,eAAe,QAChB;AAED,SAAK,UAAU,OAAO,MACpB;KAAE;KAAc,YAAY,eAAe;KAAY,EACvD,4BACD;AAED,QAAI,eAAe,WACjB,MAAK,mBAAmB,cAAc,eAAe,aACnD,eAAe;SACZ;AACL,UAAK,UAAU,OAAO,KACpB;MAAE;MAAc,WAAW,eAAe;MAAW,EACrD,yBACD;AACD,YAAO,KAAK,mBAAmB,cAC7B,eAAe;;AAInB;;AAGF,QAAK,UAAU,OAAO,KACpB;IACE,MAAMjB,yCAAyB,eAAe,KAAK;IACnD,SAAS,eAAe;IACxB;IACA;IACA,OAAO,KAAK,UAAU,UAAU;IAChC;IACD,EACD,0BACD;;EAIH,IAAIkB;AACJ,MAAI,wBAAwB,OAC1B,qBAAoB,kBAAkB;AACpC,OAAI,wBAAwB,OAC1B;AAgBF,OAAI,GAAG,eAAe,UAAU,KAC9B;AAGF,OAAI,KAAK,qBAAqB,GAAG;AAC/B,SAAK,UAAU,OAAO,KACpB,EAAE,cAAc,EAChB,2BACD;AACD,IAAK,kBACH,IAAInC,4BACF,0CAA0C,aAAa,IACvD,QACD,CACF;AACD;;AAGF,QAAK,UAAU,OAAO,MACpB,EAAE,cAAc,EAChB,2BACD;AAED,QAAK;AACL,MAAG,KACDyB,yCACEC,+BAAe,OACbA,+BAAe,OAAO,EACpB,MAAMR,mCAAmB,kBAC1B,CAAC,CACH,CAAC,QAAQ,CACX,CACF;KACA,oBAAoB;AAGzB,OAAK,UAAU,YAAY;AACzB,OAAI,OACF;AAEF,YAAS;AAET,QAAK,UAAU,OAAO,MAAM,EAAE,cAAc,EAAE,yBAAyB;AACvE,OAAI,GAAG,eAAe,UAAU,MAAM;AACpC,SAAK,UAAU,OAAO,MAAM,EAAE,cAAc,EAAE,wBAAwB;AACtE,OAAG,KACDO,yCACEC,+BAAe,OACbA,+BAAe,OAAO,EACpB,MAAMR,mCAAmB,cAC1B,CAAC,CACH,CAAC,QAAQ,CACX,CACF;;AAGH,QAAK,UAAU,OAAO,MAAM,EAAE,cAAc,EAAE,qBAAqB;AACnE,MAAG,gBAAgB;AACnB,MAAG,gBAAgB;AAEnB,SAAM,KAAK,mBAAmB,GAAG,MAAM;AAEvC,MAAG,MACD,KACAN,6CAA6BC,uCAAuB,gBAAgB,CACrE;AAED,OAAI,KAAK,mBAAmB,OAAO,aACjC,MAAK,oBAAoB;AAG3B,QAAK,UAAU,OAAO,MACpB,EAAE,cAAc,EAChB,+BACD;AACD,iBAAc,kBAAkB;;AAGlC,SAAO;;CAGT,MAAM,gBAAiC;AACrC,SAAOuB,8BAAkB;GACvB,YAAY,KAAK,OAAO;GACxB,MAAM,KAAK,OAAO;GACnB,CAAC"}