{"version":3,"file":"storage-C4FD5U8Z.cjs","names":["readPositiveIntEnv","EventEmitterPubSub","#getPubSub","NoopLeaseProvider","isLeaseProvider","#getLeaseProvider","#hasLiveThreadLease","#resolveLeaseProvider","#id","#stopLeaseRenewal","#getState","#transferThreadLease","#startLeaseRenewal","#acquireOrTransferThreadLease","#statesByPubSub","#isApprovalSuspendedRun","#isSuspendedRun","#threadKey","#isThreadBlockingRun","#publishAndWait","#threadTopic","#markRunSuspending","#getSourceId","#getThreadTarget","#releaseThreadLease","#publish","#resetState","#persistSignal","#nextStreamIdentity","#withBroadcastStream","#persistAndBroadcastIdleSignal","#broadcastPersistedSignal","#clearSuspendedRun","#sweepStaleSuspendedRecords","#watchThreadRunCompletion","#cleanupPreparedRun","#hasPendingThreadWork","#drainPendingSignals","#serializeSignal","#drainPendingContinuations","#drainPendingIdleSignals","#startContinuation","getErrorFromUnknown","#waitForRemoteRunToFinish","createSignal","#createMessageSignalInput","createMessageSignal","#generateSignalMessageId","parseMemoryRequestContext","applyStateSignal","resolveDeliveryAttributes","createSignal","StorageDomain","#notifications"],"sources":["../src/agent/thread-stream-runtime.ts","../src/notifications/delivery-policy.ts","../src/notifications/signals.ts","../src/notifications/dispatcher.ts","../src/notifications/storage.ts"],"sourcesContent":["import { randomUUID } from 'node:crypto';\n\nimport { getErrorFromUnknown } from '../error';\nimport { EventEmitterPubSub } from '../events/event-emitter';\nimport { isLeaseProvider, NoopLeaseProvider } from '../events/pubsub';\nimport type { LeaseProvider, PubSub } from '../events/pubsub';\nimport type { EventCallback } from '../events/types';\nimport { parseMemoryRequestContext } from '../memory/types';\nimport type { RequestContext } from '../request-context';\nimport { MASTRA_RESOURCE_ID_KEY, MASTRA_THREAD_ID_KEY } from '../request-context';\nimport type { MastraModelOutput } from '../stream/base/output';\nimport { readPositiveIntEnv } from '../utils';\nimport type { Agent } from './agent';\nimport type { AgentExecutionOptions } from './agent.types';\nimport type { MessageListInput } from './message-list';\nimport { createMessageSignal, createSignal, resolveDeliveryAttributes } from './signals';\nimport type { AgentMessageInput, AgentStateSignalInput, CreatedAgentSignal } from './signals';\nimport { applyStateSignal } from './state-signals';\nimport type {\n  AgentSignal,\n  AgentSubscribeToThreadOptions,\n  AgentThreadSubscription,\n  QueueAgentMessageOptions,\n  QueueAgentMessageResult,\n  SendAgentMessageOptions,\n  SendAgentMessageResult,\n  SendAgentSignalOptions,\n  SendAgentSignalAccepted,\n  SendAgentSignalResult,\n  SendAgentStateSignalOptions,\n  SendAgentStateSignalResult,\n} from './types';\n\nconst AGENT_THREAD_KEY_SEPARATOR = '\\u0000';\nconst AGENT_THREAD_STREAM_TOPIC_PREFIX = 'agent.thread-stream';\n/**\n * Lease TTL for the cross-process thread lease acquired in the idle-wake\n * path. Kept short so a crashed owner process frees the thread quickly; a\n * background timer renews it while the run is still running. Overridable via\n * `MASTRA_AGENT_THREAD_LEASE_TTL_MS` (production keeps the 15s default).\n */\nconst AGENT_THREAD_LEASE_TTL_MS = readPositiveIntEnv('MASTRA_AGENT_THREAD_LEASE_TTL_MS', 15_000);\n/**\n * Interval at which the owner process renews its lease. Defaults to TTL/3,\n * leaving room for two missed renewals (network blip, GC pause) before the\n * lease expires. Overridable via `MASTRA_AGENT_THREAD_LEASE_RENEW_INTERVAL_MS`.\n */\nconst AGENT_THREAD_LEASE_RENEW_INTERVAL_MS = readPositiveIntEnv(\n  'MASTRA_AGENT_THREAD_LEASE_RENEW_INTERVAL_MS',\n  Math.floor(AGENT_THREAD_LEASE_TTL_MS / 3),\n);\n/**\n * TTL for a suspended run's warm in-memory state — the parked thread-run record\n * (swept by #sweepStaleSuspendedRecords). The Mastra internal-workflow registry\n * reads the same `MASTRA_SUSPENDED_RUN_TTL_MS` so both expire on one bound. A\n * suspended run is kept warm so a same-instance resume can reattach and the thread\n * stays blocked; once it lapses the state is evicted and resume falls back to the\n * durable snapshot. Multi-instance deployments (resume rarely lands on the origin)\n * can shed it sooner; 30 minute default.\n */\nconst AGENT_SUSPENDED_RUN_TTL_MS = readPositiveIntEnv('MASTRA_SUSPENDED_RUN_TTL_MS', 30 * 60 * 1000);\n\nexport let defaultAgentThreadPubSub: PubSub = new EventEmitterPubSub();\n\nfunction withThreadMemory(memory: unknown, resourceId: string, threadId: string) {\n  return {\n    ...((memory && typeof memory === 'object' ? memory : {}) as Record<string, unknown>),\n    resource: (memory as { resource?: string } | undefined)?.resource ?? resourceId,\n    thread: (memory as { thread?: string } | undefined)?.thread ?? threadId,\n  };\n}\n\ntype AgentThreadRunLifecycle = 'running' | 'suspending' | 'suspended' | 'completed' | 'failed' | 'aborted';\n\ntype AgentThreadRunSuspension = {\n  toolCallId?: string;\n  toolName?: string;\n  kind: 'approval' | 'generic-tool';\n};\n\ntype AgentThreadRunRecord<OUTPUT = unknown> = {\n  agent: Agent<any, any, any, any>;\n  output: MastraModelOutput<OUTPUT>;\n  runId: string;\n  streamId: string;\n  streamSeq: number;\n  lifecycle: AgentThreadRunLifecycle;\n  suspension?: AgentThreadRunSuspension;\n  /** When the record was parked as suspended (ms epoch); drives the TTL sweep. */\n  suspendedAt?: number;\n  threadId: string;\n  resourceId?: string;\n  streamOptions: AgentExecutionOptions<OUTPUT>;\n  createSubscriberStream?: () => ReadableStream<unknown>;\n};\n\ntype PreparedThreadRun = {\n  abortController: AbortController;\n  cleanup: () => void;\n};\n\ntype PendingIdleSignal<OUTPUT = unknown> = {\n  agent: Agent<any, any, any, any>;\n  signal: CreatedAgentSignal;\n  runId: string;\n  resourceId: string;\n  threadId: string;\n  streamOptions?: AgentExecutionOptions<OUTPUT>;\n};\n\ntype PendingContinuation<OUTPUT = unknown> = {\n  agent: Agent<any, any, any, any>;\n  messages: MessageListInput;\n  runId: string;\n  resourceId: string;\n  threadId: string;\n  streamOptions?: AgentExecutionOptions<OUTPUT>;\n};\n\ntype AgentThreadRuntimeState = {\n  threadRunsById: Map<string, AgentThreadRunRecord<any>>;\n  threadRunsByStreamId: Map<string, AgentThreadRunRecord<any>>;\n  threadKeysByRunId: Map<string, string>;\n  remoteThreadKeysByRunId: Map<string, string>;\n  activeThreadRunIds: Map<string, string>;\n  activeThreadStreamIds: Map<string, string>;\n  streamSeqByRunId: Map<string, number>;\n  approvalSuspendedRunIds: Set<string>;\n  suspendedRunIds: Set<string>;\n  suspensionMetadataByRunId: Map<string, AgentThreadRunSuspension>;\n  pendingSignalsByThread: Map<string, CreatedAgentSignal[]>;\n  // Signals queued for a run that is starting but has not made its first model\n  // request yet. The first LLM step drains these and folds them into that\n  // request; `pendingSignalsByThread` follow-ups instead become their own turn.\n  preRunSignalsByThread: Map<string, CreatedAgentSignal[]>;\n  pendingIdleSignalsByThread: Map<string, PendingIdleSignal<any>[]>;\n  pendingContinuationsByThread: Map<string, PendingContinuation<any>[]>;\n  watchedThreadStreamIds: Set<string>;\n  preparedRunsById: Map<string, PreparedThreadRun>;\n  abortedRunIds: Set<string>;\n  /**\n   * Active lease-renewal timers keyed by runId. Set when the owner\n   * process wins the cross-process lease, cleared on release. Stored\n   * here (not on a Map<key,timer>) so a run's renewal timer survives even\n   * if `activeThreadRunIds` is rotated by a follow-up signal.\n   */\n  leaseRenewalTimers: Map<string, ReturnType<typeof setInterval>>;\n};\n\nexport type AgentThreadState = 'active' | 'idle';\n\ntype SerializableAgentSignal = AgentSignal & Pick<CreatedAgentSignal, 'id' | 'createdAt'>;\n\ntype AgentThreadStreamRuntimeEvent =\n  | { type: 'run-registered'; runId: string; streamId: string; streamSeq: number }\n  | { type: 'stream-part'; runId: string; streamId: string; part: unknown; sourceId: string }\n  | { type: 'run-completed'; runId: string; streamId?: string }\n  | { type: 'run-suspended'; runId: string; streamId?: string }\n  | { type: 'run-abort-requested'; runId: string; streamId: string }\n  | { type: 'run-aborted'; runId: string; streamId?: string }\n  | { type: 'run-failed'; runId: string; streamId?: string; error: string }\n  | { type: 'signal-enqueued'; runId: string; signal: SerializableAgentSignal; sourceId: string; preRun?: boolean };\n\nfunction createRuntimeState(): AgentThreadRuntimeState {\n  return {\n    threadRunsById: new Map(),\n    threadRunsByStreamId: new Map(),\n    threadKeysByRunId: new Map(),\n    remoteThreadKeysByRunId: new Map(),\n    activeThreadRunIds: new Map(),\n    activeThreadStreamIds: new Map(),\n    streamSeqByRunId: new Map(),\n    approvalSuspendedRunIds: new Set(),\n    suspendedRunIds: new Set(),\n    suspensionMetadataByRunId: new Map(),\n    pendingSignalsByThread: new Map(),\n    preRunSignalsByThread: new Map(),\n    pendingIdleSignalsByThread: new Map(),\n    pendingContinuationsByThread: new Map(),\n    watchedThreadStreamIds: new Set(),\n    preparedRunsById: new Map(),\n    abortedRunIds: new Set(),\n    leaseRenewalTimers: new Map(),\n  };\n}\n\nexport class AgentThreadStreamRuntime {\n  #id?: string;\n  #statesByPubSub = new WeakMap<PubSub, AgentThreadRuntimeState>();\n\n  #getPubSub(pubsub?: PubSub): PubSub {\n    return pubsub ?? defaultAgentThreadPubSub;\n  }\n\n  /**\n   * Resolve the {@link LeaseProvider} for the configured pubsub. Leasing is\n   * a separate capability from event delivery: a backend only implements it\n   * when it can genuinely coordinate a distributed lock (Redis via SET-NX,\n   * in-memory for single-process). We feature-detect once here so all lease\n   * call sites can use the resolved provider unconditionally.\n   *\n   * `CachingPubSub` exposes its inner's lease provider via `getLeaseProvider`\n   * (caching is transparent to leasing). Otherwise we duck-type the pubsub\n   * directly. Backends that cannot lease fall back to {@link NoopLeaseProvider}\n   * (always-win / no-op), preserving single-process behavior.\n   */\n  #getLeaseProvider(pubsub?: PubSub): LeaseProvider {\n    const resolved = this.#getPubSub(pubsub);\n    const unwrap = (resolved as { getLeaseProvider?: () => LeaseProvider | undefined }).getLeaseProvider;\n    if (typeof unwrap === 'function') {\n      const inner = unwrap.call(resolved);\n      return inner ?? NoopLeaseProvider;\n    }\n    return isLeaseProvider(resolved) ? resolved : NoopLeaseProvider;\n  }\n\n  #resolveLeaseProvider(pubsub?: PubSub): { provider: LeaseProvider; isFallback: boolean } {\n    const provider = this.#getLeaseProvider(pubsub);\n    return { provider, isFallback: provider === NoopLeaseProvider };\n  }\n\n  async #hasLiveThreadLease(pubsub: PubSub, key: string, runId: string): Promise<boolean> {\n    const { provider, isFallback } = this.#resolveLeaseProvider(pubsub);\n    if (isFallback) return true;\n    return provider\n      .getLeaseOwner(key)\n      .then(owner => owner === runId)\n      .catch(() => false);\n  }\n\n  #getSourceId(): string {\n    this.#id ??= randomUUID();\n    return this.#id;\n  }\n\n  /**\n   * Fire-and-forget release of the cross-process thread lease held by\n   * this owner. Safe to call when no lease was ever acquired — the\n   * pubsub's `releaseLease` is a no-op for non-owners (Lua-guarded\n   * GET+DEL on Redis), and the default in-memory implementation is\n   * identical. Also stops the renewal timer if one is running for\n   * this run.\n   */\n  #releaseThreadLease(pubsub: PubSub | undefined, key: string, runId: string): void {\n    const resolved = this.#getPubSub(pubsub);\n    this.#stopLeaseRenewal(resolved, runId);\n    void this.#getLeaseProvider(resolved)\n      .releaseLease(key, runId)\n      .catch(() => {});\n  }\n\n  /**\n   * Start a background timer that renews the cross-process lease at\n   * TTL/3 intervals while the run is still going. If the lease is lost\n   * (e.g. expired due to clock skew or pubsub outage) the renewal\n   * stops itself — there's nothing useful we can do from the runner\n   * side beyond log; the original owner will keep running until the run\n   * itself errors or completes.\n   */\n  #startLeaseRenewal(pubsub: PubSub, key: string, runId: string): void {\n    const state = this.#getState(pubsub);\n    if (state.leaseRenewalTimers.has(runId)) return;\n    const leaseProvider = this.#getLeaseProvider(pubsub);\n    const timer = setInterval(() => {\n      void leaseProvider\n        .renewLease(key, runId, AGENT_THREAD_LEASE_TTL_MS)\n        .then(renewed => {\n          if (!renewed) {\n            // If renewLease reports the lease is gone, stop renewing; the current stream may still finish,\n            // but another process can now claim the thread until this run completes or errors.\n            this.#stopLeaseRenewal(pubsub, runId);\n          }\n        })\n        .catch(() => {});\n    }, AGENT_THREAD_LEASE_RENEW_INTERVAL_MS);\n    // Don't keep the process alive solely to renew a lease.\n    if (typeof timer === 'object' && timer && typeof (timer as any).unref === 'function') {\n      (timer as any).unref();\n    }\n    state.leaseRenewalTimers.set(runId, timer);\n  }\n\n  #stopLeaseRenewal(pubsub: PubSub, runId: string): void {\n    const state = this.#getState(pubsub);\n    const timer = state.leaseRenewalTimers.get(runId);\n    if (!timer) return;\n    clearInterval(timer);\n    state.leaseRenewalTimers.delete(runId);\n  }\n\n  /**\n   * Hand the cross-process thread lease from a finishing run (`fromRunId`)\n   * to the run that will drain queued follow-up work next (`toRunId`),\n   * without the lease key ever going empty.\n   *\n   * The previous owner releases its renewal timer and the new owner starts\n   * its own; the lease key is re-stamped by `transferLease` (with a full fresh\n   * TTL). On atomic backends (Redis, in-memory) a racing process cannot win a\n   * freed key between a release and a re-acquire. Backends that can't transfer\n   * atomically implement `transferLease` as release+acquire internally and own\n   * that race cost. Returns `true` if the new owner now holds the lease.\n   */\n  async #transferThreadLease(\n    pubsub: PubSub | undefined,\n    key: string,\n    fromRunId: string,\n    toRunId: string,\n  ): Promise<boolean> {\n    const resolved = this.#getPubSub(pubsub);\n    const leaseProvider = this.#getLeaseProvider(resolved);\n    // `transferLease` is a required `LeaseProvider` method. Atomic backends\n    // (Redis, in-memory) swap the key gap-free; backends that can't be atomic\n    // implement it as release+acquire internally and own that race cost.\n    const held = await leaseProvider\n      .transferLease(key, fromRunId, toRunId, AGENT_THREAD_LEASE_TTL_MS)\n      .catch(() => false);\n    // Move the renewal timer to the new owner regardless: the old timer is\n    // owner-guarded and would only no-op now, and the new owner needs its\n    // own keep-alive for long drains.\n    this.#stopLeaseRenewal(resolved, fromRunId);\n    if (held) {\n      this.#startLeaseRenewal(resolved, key, toRunId);\n    }\n    return held;\n  }\n\n  /**\n   * Ensure this process owns the cross-process lease for `toRunId` before it\n   * starts a run, regardless of whether it already held the lease.\n   *\n   * - When `fromRunId` is provided (draining after a run this process owned),\n   *   atomically transfer the held lease to `toRunId` — gap-free, no empty key.\n   * - When `fromRunId` is absent, or the transfer reports the old owner no\n   *   longer holds the lease, fall back to a fresh `acquireLease`. This covers\n   *   a *different* process that observed the owner finish via pub/sub and now\n   *   wants to wake the thread: it never held the lease, so it must win one.\n   *\n   * On success the renewal timer is started for `toRunId`. On failure the\n   * returned `owner` is the current holder so the caller can forward work to it.\n   */\n  async #acquireOrTransferThreadLease(\n    pubsub: PubSub | undefined,\n    key: string,\n    toRunId: string,\n    fromRunId?: string,\n  ): Promise<{ acquired: boolean; owner?: string }> {\n    const resolved = this.#getPubSub(pubsub);\n    if (fromRunId) {\n      const transferred = await this.#transferThreadLease(pubsub, key, fromRunId, toRunId);\n      if (transferred) return { acquired: true, owner: toRunId };\n      // Old owner lost the lease before the handoff — fall through to acquire.\n    }\n    const leaseProvider = this.#getLeaseProvider(resolved);\n    const result = await leaseProvider\n      .acquireLease(key, toRunId, AGENT_THREAD_LEASE_TTL_MS)\n      .catch(() => ({ acquired: false as boolean, owner: undefined as string | undefined }));\n    if (result.acquired) {\n      this.#startLeaseRenewal(resolved, key, toRunId);\n      return { acquired: true, owner: toRunId };\n    }\n    return { acquired: false, owner: result.owner };\n  }\n\n  /**\n   * Whether the thread has any queued follow-up work that a finishing run's\n   * completion handler would drain next: pending follow-up signals (including\n   * any pre-run leftover that will be folded in), queued continuations, or\n   * queued idle signals.\n   */\n  #hasPendingThreadWork(state: AgentThreadRuntimeState, key: string): boolean {\n    return (\n      (state.pendingSignalsByThread.get(key)?.length ?? 0) > 0 ||\n      (state.preRunSignalsByThread.get(key)?.length ?? 0) > 0 ||\n      (state.pendingContinuationsByThread.get(key)?.length ?? 0) > 0 ||\n      (state.pendingIdleSignalsByThread.get(key)?.length ?? 0) > 0\n    );\n  }\n\n  #getState(pubsub?: PubSub): AgentThreadRuntimeState {\n    const resolvedPubSub = this.#getPubSub(pubsub);\n    let state = this.#statesByPubSub.get(resolvedPubSub);\n    if (!state) {\n      state = createRuntimeState();\n      this.#statesByPubSub.set(resolvedPubSub, state);\n    }\n    return state;\n  }\n\n  #threadKey(resourceId: string | undefined, threadId: string): string {\n    return [resourceId ?? '', threadId].join(AGENT_THREAD_KEY_SEPARATOR);\n  }\n\n  #threadTopic(key: string): string {\n    return `${AGENT_THREAD_STREAM_TOPIC_PREFIX}.${encodeURIComponent(key)}`;\n  }\n\n  #isApprovalSuspendedRun(state: AgentThreadRuntimeState, runId: string) {\n    return state.approvalSuspendedRunIds.has(runId);\n  }\n\n  #isSuspendedRun(state: AgentThreadRuntimeState, runId: string) {\n    return state.suspendedRunIds.has(runId) || this.#isApprovalSuspendedRun(state, runId);\n  }\n\n  #isThreadBlockingRun(state: AgentThreadRuntimeState, record: AgentThreadRunRecord<any>) {\n    return (\n      record.output.status === 'running' ||\n      record.output.status === 'suspended' ||\n      record.lifecycle === 'suspending' ||\n      record.lifecycle === 'suspended' ||\n      !!record.suspension ||\n      this.#isSuspendedRun(state, record.runId)\n    );\n  }\n\n  #serializeSignal(signal: CreatedAgentSignal): SerializableAgentSignal {\n    return signal;\n  }\n\n  #nextStreamIdentity(state: AgentThreadRuntimeState, runId: string) {\n    const streamSeq = (state.streamSeqByRunId.get(runId) ?? 0) + 1;\n    state.streamSeqByRunId.set(runId, streamSeq);\n    return { streamId: randomUUID(), streamSeq };\n  }\n\n  #markRunSuspending(\n    state: AgentThreadRuntimeState,\n    runId: string,\n    streamId: string,\n    suspension: AgentThreadRunSuspension,\n  ) {\n    state.suspendedRunIds.add(runId);\n    state.suspensionMetadataByRunId.set(runId, suspension);\n    const record = state.threadRunsByStreamId.get(streamId) ?? state.threadRunsById.get(runId);\n    if (record) {\n      record.lifecycle = 'suspending';\n      record.suspension = suspension;\n    }\n    if (suspension.kind === 'approval') {\n      state.approvalSuspendedRunIds.add(runId);\n    }\n  }\n\n  #clearSuspendedRun(state: AgentThreadRuntimeState, runId: string) {\n    state.suspendedRunIds.delete(runId);\n    state.suspensionMetadataByRunId.delete(runId);\n    state.approvalSuspendedRunIds.delete(runId);\n  }\n\n  #generateSignalMessageId(\n    agent: Agent<any, any, any, any>,\n    target: { threadId?: string; resourceId?: string },\n  ): string {\n    return (\n      agent.getMastraInstance?.()?.generateId({\n        idType: 'message',\n        source: 'agent',\n        entityId: agent.id,\n        threadId: target.threadId,\n        resourceId: target.resourceId,\n      }) ?? randomUUID()\n    );\n  }\n\n  #createMessageSignalInput(message: AgentMessageInput): AgentSignal {\n    const normalizedMessage = typeof message === 'string' || Array.isArray(message) ? { contents: message } : message;\n    return {\n      ...normalizedMessage,\n      type: 'user',\n      tagName: 'user',\n    };\n  }\n\n  getThreadState(options: { resourceId?: string; threadId: string }, pubsub?: PubSub): AgentThreadState {\n    const state = this.#getState(pubsub);\n    const key = this.#threadKey(options.resourceId, options.threadId);\n    const activeRunId = state.activeThreadRunIds.get(key);\n    if (!activeRunId) return 'idle';\n\n    const activeRecord = state.threadRunsById.get(activeRunId);\n    if (activeRecord && !this.#isThreadBlockingRun(state, activeRecord)) {\n      state.activeThreadRunIds.delete(key);\n      return 'idle';\n    }\n\n    return 'active';\n  }\n\n  #publish(pubsub: PubSub | undefined, key: string, event: AgentThreadStreamRuntimeEvent) {\n    void this.#publishAndWait(pubsub, key, event).catch(() => {});\n  }\n\n  async #publishAndWait(pubsub: PubSub | undefined, key: string, event: AgentThreadStreamRuntimeEvent) {\n    await this.#getPubSub(pubsub).publish(this.#threadTopic(key), {\n      type: event.type,\n      runId: event.runId,\n      data: event,\n    });\n  }\n\n  #withBroadcastStream<OUTPUT>(\n    output: MastraModelOutput<OUTPUT>,\n    pubsub: PubSub | undefined,\n    key: string,\n    streamId: string,\n  ) {\n    const runtime = this;\n\n    const parts: unknown[] = [];\n    const waiters = new Set<() => void>();\n    let started = false;\n    let done = false;\n    let error: unknown;\n\n    const wake = () => {\n      const pending = [...waiters];\n      waiters.clear();\n      for (const waiter of pending) waiter();\n    };\n\n    const emitPart = async (part: unknown) => {\n      if (part && typeof part === 'object' && 'type' in part) {\n        const typedPart = part as { type?: string; payload?: { toolCallId?: string; toolName?: string } };\n        if (typedPart.type === 'tool-call-approval' || typedPart.type === 'tool-call-suspended') {\n          runtime.#markRunSuspending(runtime.#getState(pubsub), output.runId, streamId, {\n            toolCallId: typedPart.payload?.toolCallId,\n            toolName: typedPart.payload?.toolName,\n            kind: typedPart.type === 'tool-call-approval' ? 'approval' : 'generic-tool',\n          });\n        }\n      }\n      parts.push(part);\n      await runtime.#publishAndWait(pubsub, key, {\n        type: 'stream-part',\n        runId: output.runId,\n        streamId,\n        part,\n        sourceId: runtime.#getSourceId(),\n      });\n      wake();\n    };\n\n    const start = () => {\n      if (started) return;\n      started = true;\n      void (async () => {\n        try {\n          const source = output.fullStream as ReadableStream<unknown> | undefined;\n          if (!source) return;\n\n          if (typeof source.getReader === 'function') {\n            const reader = source.getReader();\n            try {\n              while (true) {\n                const { value: part, done: streamDone } = await reader.read();\n                if (streamDone) break;\n                await emitPart(part);\n              }\n            } finally {\n              reader.releaseLock();\n            }\n          } else {\n            for await (const part of source as any) {\n              await emitPart(part);\n            }\n          }\n        } catch (caught) {\n          error = caught;\n        } finally {\n          done = true;\n          wake();\n        }\n      })();\n    };\n\n    const createStream = () => {\n      let index = 0;\n      let closed = false;\n      let waiter: (() => void) | undefined;\n      return new ReadableStream({\n        async pull(controller) {\n          start();\n          while (!closed) {\n            if (index < parts.length) {\n              controller.enqueue(parts[index++]);\n              return;\n            }\n            if (error) {\n              controller.error(error);\n              return;\n            }\n            if (done) {\n              controller.close();\n              return;\n            }\n            await new Promise<void>(resolve => {\n              waiter = resolve;\n              waiters.add(resolve);\n            });\n            if (waiter) {\n              waiters.delete(waiter);\n              waiter = undefined;\n            }\n          }\n        },\n        cancel() {\n          closed = true;\n          if (waiter) {\n            waiters.delete(waiter);\n            waiter();\n            waiter = undefined;\n          }\n        },\n      });\n    };\n\n    return { output, createSubscriberStream: createStream, startBroadcast: start };\n  }\n\n  #getThreadTarget(options?: { memory?: AgentExecutionOptions<any>['memory']; requestContext?: RequestContext }) {\n    const thread = options?.memory?.thread;\n    const threadId =\n      (options?.requestContext?.get(MASTRA_THREAD_ID_KEY) as string | undefined) ||\n      (typeof thread === 'string' ? thread : thread?.id);\n    const resourceId =\n      (options?.requestContext?.get(MASTRA_RESOURCE_ID_KEY) as string | undefined) || options?.memory?.resource;\n\n    return { threadId, resourceId };\n  }\n\n  prepareRunOptions<OUTPUT>(options: AgentExecutionOptions<OUTPUT>, pubsub?: PubSub): AgentExecutionOptions<OUTPUT> {\n    const { threadId } = this.#getThreadTarget(options);\n    if (!threadId || !options.runId) return options;\n\n    const state = this.#getState(pubsub);\n    const abortController = new AbortController();\n    const upstreamAbortSignal = options.abortSignal;\n    const abort = () => abortController.abort();\n    if (upstreamAbortSignal?.aborted) {\n      abort();\n    } else {\n      upstreamAbortSignal?.addEventListener('abort', abort, { once: true });\n    }\n\n    state.preparedRunsById.set(options.runId, {\n      abortController,\n      cleanup: () => upstreamAbortSignal?.removeEventListener('abort', abort),\n    });\n\n    if (state.abortedRunIds.has(options.runId)) {\n      abort();\n    }\n\n    return {\n      ...options,\n      abortSignal: abortController.signal,\n    };\n  }\n\n  abortRun(runId: string, pubsub?: PubSub): boolean {\n    const state = this.#getState(pubsub);\n    const preparedRun = state.preparedRunsById.get(runId);\n    if (!preparedRun) {\n      state.abortedRunIds.add(runId);\n      return false;\n    }\n\n    preparedRun.abortController.abort();\n    state.abortedRunIds.add(runId);\n\n    const key = state.threadKeysByRunId.get(runId);\n    if (key) {\n      const streamId = state.activeThreadRunIds.get(key) === runId ? state.activeThreadStreamIds.get(key) : undefined;\n      this.#releaseThreadLease(pubsub, key, runId);\n      this.#publish(pubsub, key, { type: 'run-aborted', runId, streamId });\n    }\n\n    return true;\n  }\n\n  getActiveThreadRunId(options: AgentSubscribeToThreadOptions, pubsub?: PubSub): string | undefined {\n    const state = this.#getState(pubsub);\n    const key = this.#threadKey(options.resourceId, options.threadId);\n    const activeRunId = state.activeThreadRunIds.get(key);\n    if (!activeRunId) return undefined;\n\n    const record = state.threadRunsById.get(activeRunId);\n    if (record && !this.#isThreadBlockingRun(state, record)) return undefined;\n\n    return activeRunId;\n  }\n\n  getResumableThreadRun(\n    options: AgentSubscribeToThreadOptions & { runId: string; toolCallId?: string },\n    pubsub?: PubSub,\n  ): { runId: string; toolCallId?: string } | undefined {\n    const state = this.#getState(pubsub);\n    const key = this.#threadKey(options.resourceId, options.threadId);\n    const record = state.threadRunsById.get(options.runId);\n    const isSuspended = this.#isSuspendedRun(state, options.runId);\n    if (!record || state.threadKeysByRunId.get(options.runId) !== key || !isSuspended) {\n      return undefined;\n    }\n\n    const suspension = record.suspension ?? state.suspensionMetadataByRunId.get(options.runId);\n    if (options.toolCallId && suspension?.toolCallId && suspension.toolCallId !== options.toolCallId) {\n      return undefined;\n    }\n\n    return { runId: options.runId, toolCallId: options.toolCallId ?? suspension?.toolCallId };\n  }\n\n  abortThread(options: AgentSubscribeToThreadOptions, pubsub?: PubSub): boolean {\n    const resolvedPubSub = this.#getPubSub(pubsub);\n    const state = this.#getState(resolvedPubSub);\n    const key = this.#threadKey(options.resourceId, options.threadId);\n    const runId = this.getActiveThreadRunId(options, resolvedPubSub);\n    if (!runId) return false;\n    if (state.preparedRunsById.has(runId)) return this.abortRun(runId, resolvedPubSub);\n    if (state.threadKeysByRunId.get(runId) === key) {\n      // Reserved locally (a sendSignal wake that has not prepared its run yet):\n      // record the abort intent in abortedRunIds so prepareRunOptions aborts the\n      // run the moment it starts, instead of letting it run to completion.\n      this.abortRun(runId, resolvedPubSub);\n      return true;\n    }\n    if (state.remoteThreadKeysByRunId.get(runId) !== key) return false;\n    const streamId = state.activeThreadStreamIds.get(key);\n    if (!streamId) return false;\n    this.#publish(resolvedPubSub, key, { type: 'run-abort-requested', runId, streamId });\n    return true;\n  }\n\n  /** @internal */\n  resetForTests() {\n    for (const pubsub of [defaultAgentThreadPubSub]) {\n      this.#resetState(pubsub);\n      void (pubsub as { close?: () => Promise<void> }).close?.();\n    }\n    defaultAgentThreadPubSub = new EventEmitterPubSub();\n  }\n\n  #resetState(pubsub: PubSub) {\n    const state = this.#statesByPubSub.get(pubsub);\n    if (!state) return;\n\n    state.preparedRunsById.forEach(preparedRun => {\n      preparedRun.abortController.abort();\n      preparedRun.cleanup();\n    });\n    state.leaseRenewalTimers.forEach(timer => clearInterval(timer));\n    state.leaseRenewalTimers.clear();\n    state.threadRunsById.clear();\n    state.threadRunsByStreamId.clear();\n    state.threadKeysByRunId.clear();\n    state.remoteThreadKeysByRunId.clear();\n    state.activeThreadRunIds.clear();\n    state.approvalSuspendedRunIds.clear();\n    state.suspendedRunIds.clear();\n    state.suspensionMetadataByRunId.clear();\n    state.pendingSignalsByThread.clear();\n    state.preRunSignalsByThread.clear();\n    state.pendingIdleSignalsByThread.clear();\n    state.pendingContinuationsByThread.clear();\n    state.activeThreadStreamIds.clear();\n    state.streamSeqByRunId.clear();\n    state.watchedThreadStreamIds.clear();\n    state.preparedRunsById.clear();\n    state.abortedRunIds.clear();\n  }\n\n  #cleanupPreparedRun(state: AgentThreadRuntimeState, runId: string) {\n    state.preparedRunsById.get(runId)?.cleanup();\n    state.preparedRunsById.delete(runId);\n    state.abortedRunIds.delete(runId);\n  }\n\n  async #persistSignal(\n    agent: Agent<any, any, any, any>,\n    signal: CreatedAgentSignal,\n    resourceId: string,\n    threadId: string,\n    requestContext?: RequestContext,\n  ) {\n    // Transient signals are delivery-only: never write them to storage, even when the\n    // active-behavior asked to persist. Honored here (not just in the memory layer) so it holds\n    // for any memory implementation, including ones without a signal-aware save filter.\n    if (signal.transient) return;\n    const memory = await agent.getMemory({ requestContext });\n    if (!memory) return;\n    await memory.saveMessages({\n      messages: [signal.toDBMessage({ resourceId, threadId })],\n    });\n  }\n\n  #broadcastPersistedSignal(\n    state: AgentThreadRuntimeState,\n    pubsub: PubSub | undefined,\n    key: string,\n    runId: string,\n    signal: CreatedAgentSignal,\n    resourceId: string,\n    threadId: string,\n  ) {\n    let finish!: () => void;\n    const finished = new Promise<void>(resolve => {\n      finish = resolve;\n    });\n    const parts: any[] = [\n      { type: 'start', runId },\n      { ...signal.toDataPart(), runId },\n      {\n        type: 'finish',\n        runId,\n        payload: {\n          stepResult: { reason: 'stop' },\n          output: {\n            usage: { inputTokens: 0, outputTokens: 0, totalTokens: 0 },\n          },\n        },\n      },\n    ];\n    const output = {\n      runId,\n      status: 'running',\n      fullStream: new ReadableStream({\n        start(controller) {\n          for (const part of parts) controller.enqueue(part);\n          controller.close();\n          finish();\n        },\n      }),\n      _waitUntilFinished: () => finished,\n    } as MastraModelOutput<any>;\n    const { streamId, streamSeq } = this.#nextStreamIdentity(state, runId);\n    const {\n      output: outputForSubscribers,\n      createSubscriberStream,\n      startBroadcast,\n    } = this.#withBroadcastStream(output, pubsub, key, streamId);\n    const record: AgentThreadRunRecord<any> = {\n      agent: { id: `persisted-signal:${signal.id}` } as Agent<any, any, any, any>,\n      output: outputForSubscribers,\n      runId,\n      streamId,\n      streamSeq,\n      lifecycle: 'running',\n      threadId,\n      resourceId,\n      streamOptions: {},\n      createSubscriberStream,\n    };\n\n    state.threadRunsById.set(runId, record);\n    state.threadRunsByStreamId.set(streamId, record);\n    state.threadKeysByRunId.set(runId, key);\n    state.activeThreadStreamIds.set(key, streamId);\n    const registered = this.#publishAndWait(pubsub, key, { type: 'run-registered', runId, streamId, streamSeq });\n    void registered.then(startBroadcast, startBroadcast);\n    void outputForSubscribers._waitUntilFinished().finally(() => {\n      setTimeout(() => {\n        state.threadRunsByStreamId.delete(streamId);\n        if (state.threadRunsById.get(runId) === record) {\n          state.threadRunsById.delete(runId);\n          state.threadKeysByRunId.delete(runId);\n        }\n        if (state.activeThreadRunIds.get(key) === runId && state.activeThreadStreamIds.get(key) === streamId) {\n          state.activeThreadRunIds.delete(key);\n          state.activeThreadStreamIds.delete(key);\n        }\n        this.#releaseThreadLease(pubsub, key, runId);\n        this.#publish(pubsub, key, { type: 'run-completed', runId, streamId });\n      }, 0);\n    });\n  }\n\n  async #persistAndBroadcastIdleSignal(\n    state: AgentThreadRuntimeState,\n    pubsub: PubSub | undefined,\n    key: string,\n    runId: string,\n    agent: Agent<any, any, any, any>,\n    signal: CreatedAgentSignal,\n    resourceId: string,\n    threadId: string,\n    requestContext?: RequestContext,\n  ) {\n    if (signal.transient) return;\n\n    await this.#persistSignal(agent, signal, resourceId, threadId, requestContext);\n    this.#broadcastPersistedSignal(state, pubsub, key, runId, signal, resourceId, threadId);\n  }\n\n  /**\n   * Evict SUSPENDED records parked longer than {@link AGENT_SUSPENDED_RUN_TTL_MS}.\n   * Called lazily on each registration so cleanup is proportional to activity and\n   * zero-cost when idle — mirrors the internal-workflow registry sweep. Bounds the\n   * records left behind by abandoned suspends and by resumes that land on a\n   * different instance (which never clean the origin instance's record).\n   *\n   * When the expiring record is still the run's current record — an abandoned\n   * suspend, not one superseded by a same-instance resume — the teardown mirrors\n   * #watchThreadRunCompletion's terminal path: it clears run-level state, releases\n   * the cross-process lease, and publishes `run-completed` so remote subscribers\n   * stop treating the thread as blocked and drain any queued follow-up work. A\n   * superseded older stream just has its stream entry dropped; the resumed run\n   * keeps its lease, suspended marker, and active slot.\n   */\n  #sweepStaleSuspendedRecords(state: AgentThreadRuntimeState, pubsub: PubSub | undefined) {\n    const now = Date.now();\n    for (const [streamId, record] of state.threadRunsByStreamId) {\n      if (record.lifecycle !== 'suspended' || record.suspendedAt === undefined) continue;\n      if (now - record.suspendedAt <= AGENT_SUSPENDED_RUN_TTL_MS) continue;\n      state.threadRunsByStreamId.delete(streamId);\n      state.watchedThreadStreamIds.delete(streamId);\n      // A same-instance resume re-registers the run under a newer streamId, so a\n      // record that is no longer the run's current record is just the superseded\n      // older stream: dropping its stream entry above is enough. Only the current\n      // record (an abandoned suspend) gets the full run-level teardown below.\n      if (state.threadRunsById.get(record.runId) !== record) continue;\n      const staleKey = this.#threadKey(record.resourceId, record.threadId);\n      state.threadRunsById.delete(record.runId);\n      state.threadKeysByRunId.delete(record.runId);\n      this.#clearSuspendedRun(state, record.runId);\n      // Stop renewing and release the cross-process lease, otherwise the run's\n      // lease-renewal timer keeps the thread owned forever on other instances.\n      this.#releaseThreadLease(pubsub, staleKey, record.runId);\n      if (\n        state.activeThreadRunIds.get(staleKey) === record.runId &&\n        state.activeThreadStreamIds.get(staleKey) === streamId\n      ) {\n        state.activeThreadRunIds.delete(staleKey);\n        state.activeThreadStreamIds.delete(staleKey);\n      }\n      this.#publish(pubsub, staleKey, { type: 'run-completed', runId: record.runId, streamId });\n    }\n  }\n\n  registerRun<OUTPUT>(\n    agent: Agent<any, any, any, any>,\n    output: MastraModelOutput<OUTPUT>,\n    streamOptions: AgentExecutionOptions<OUTPUT>,\n    pubsub?: PubSub,\n  ): Promise<void> | undefined {\n    const { threadId, resourceId } = this.#getThreadTarget(streamOptions);\n    if (!threadId) return;\n\n    const state = this.#getState(pubsub);\n    this.#sweepStaleSuspendedRecords(state, pubsub);\n    const key = this.#threadKey(resourceId, threadId);\n    const { streamId, streamSeq } = this.#nextStreamIdentity(state, output.runId);\n    const {\n      output: outputForSubscribers,\n      createSubscriberStream,\n      startBroadcast,\n    } = this.#withBroadcastStream(output, pubsub, key, streamId);\n    const record: AgentThreadRunRecord<OUTPUT> = {\n      agent,\n      output: outputForSubscribers,\n      runId: output.runId,\n      streamId,\n      streamSeq,\n      lifecycle: 'running',\n      threadId,\n      resourceId,\n      streamOptions: streamOptions as AgentThreadRunRecord<OUTPUT>['streamOptions'],\n      createSubscriberStream,\n    };\n\n    this.#clearSuspendedRun(state, output.runId);\n    state.threadRunsById.set(output.runId, record);\n    state.threadRunsByStreamId.set(streamId, record);\n    state.threadKeysByRunId.set(output.runId, key);\n    state.activeThreadRunIds.set(key, output.runId);\n    state.activeThreadStreamIds.set(key, streamId);\n    const resolvedPubSub = this.#getPubSub(pubsub);\n    const registered = (async () => {\n      // Every thread-bound run must hold the cross-process lease while it is\n      // live: the liveness checks (markActiveIfLive / #waitForRemoteRunToFinish)\n      // treat a lease-less run as a ghost, so a plain `agent.stream()` run that\n      // never acquired would let contending instances start competing runs\n      // instead of serializing behind it. Acquire BEFORE publishing\n      // `run-registered` so an observer that checks liveness on receipt finds\n      // the lease held. Same-owner acquire is an idempotent TTL refresh, so\n      // signal-woken runs that already hold the lease under this runId just\n      // renew. Fail-open on loss or error (simultaneous-start race): proceed\n      // and never roll back the local registration — matches pre-lease\n      // semantics and sendSignal's documented fail-open rationale. A thrown\n      // acquire (transient provider error) is treated as acquired so renewal\n      // starts: if the acquire landed server-side but the response failed,\n      // skipping renewal would let the lease expire mid-run; renewal\n      // self-stops when we don't own the key.\n      const lease = await this.#getLeaseProvider(resolvedPubSub)\n        .acquireLease(key, output.runId, AGENT_THREAD_LEASE_TTL_MS)\n        .catch(() => ({ acquired: true as boolean }));\n      if (lease.acquired) this.#startLeaseRenewal(resolvedPubSub, key, output.runId);\n      await this.#publishAndWait(pubsub, key, {\n        type: 'run-registered',\n        runId: output.runId,\n        streamId,\n        streamSeq,\n      });\n    })();\n    // Always drive the run's stream to completion, even when no caller consumes\n    // the returned output (e.g. a fire-and-forget schedule wake). The broadcast\n    // tee buffers every part, so a later/external subscriber still replays the\n    // full stream; without this pump the run never reaches a terminal state and\n    // its active-run record + thread lease would never release, permanently\n    // wedging the thread.\n    void registered.then(startBroadcast, startBroadcast);\n    this.#watchThreadRunCompletion(state, pubsub, key, record);\n    return registered;\n  }\n\n  #watchThreadRunCompletion(\n    state: AgentThreadRuntimeState,\n    pubsub: PubSub | undefined,\n    key: string,\n    record: AgentThreadRunRecord<any>,\n  ) {\n    if (state.watchedThreadStreamIds.has(record.streamId)) return;\n    state.watchedThreadStreamIds.add(record.streamId);\n\n    void record.output._waitUntilFinished().finally(() => {\n      state.watchedThreadStreamIds.delete(record.streamId);\n      this.#cleanupPreparedRun(state, record.runId);\n\n      if (record.output.status === 'suspended' && this.#isSuspendedRun(state, record.runId)) {\n        record.lifecycle = 'suspended';\n        // Leak fix: stamp when the run parked so the lazy TTL sweep\n        // (#sweepStaleSuspendedRecords) can evict it. The record stays fully intact\n        // for resume routing / thread-blocking / subscriber replay exactly as before\n        // — it is simply no longer retained for the life of the process. Mirrors the\n        // internal-workflow registry, which already bounds parked runs this way.\n        record.suspendedAt = Date.now();\n        this.#publish(pubsub, key, { type: 'run-suspended', runId: record.runId, streamId: record.streamId });\n        return;\n      }\n\n      record.lifecycle = 'completed';\n      this.#clearSuspendedRun(state, record.runId);\n      state.threadRunsByStreamId.delete(record.streamId);\n      if (state.threadRunsById.get(record.runId) === record) {\n        state.threadRunsById.delete(record.runId);\n        state.threadKeysByRunId.delete(record.runId);\n      }\n\n      if (\n        state.activeThreadRunIds.get(key) === record.runId &&\n        state.activeThreadStreamIds.get(key) === record.streamId\n      ) {\n        state.activeThreadRunIds.delete(key);\n        state.activeThreadStreamIds.delete(key);\n      }\n\n      // If queued follow-up work exists, keep the cross-process lease held by\n      // handing it to the next run instead of releasing it: releasing here\n      // would briefly empty the lease key, letting a racing process win it and\n      // start a competing run on this thread. The drain runs under the\n      // transferred lease and releases it only once every queue is empty. If\n      // there's no pending work, release as usual so other processes can wake\n      // the thread.\n      this.#publish(pubsub, key, { type: 'run-completed', runId: record.runId, streamId: record.streamId });\n      if (this.#hasPendingThreadWork(state, key)) {\n        void this.#drainPendingSignals(state, pubsub, key, record);\n      } else {\n        this.#releaseThreadLease(pubsub, key, record.runId);\n      }\n    });\n  }\n\n  async #drainPendingSignals(\n    state: AgentThreadRuntimeState,\n    pubsub: PubSub | undefined,\n    key: string,\n    previousRun: AgentThreadRunRecord<any>,\n  ) {\n    if (state.activeThreadRunIds.has(key)) {\n      return;\n    }\n\n    // A run can finish before its first model request drained its pre-run\n    // signals (e.g. it errored early). Don't strand them — fold them into the\n    // follow-up queue so the next run still picks them up.\n    const preRunLeftover = state.preRunSignalsByThread.get(key);\n    if (preRunLeftover?.length) {\n      state.preRunSignalsByThread.delete(key);\n      state.pendingSignalsByThread.set(key, [...preRunLeftover, ...(state.pendingSignalsByThread.get(key) ?? [])]);\n    }\n\n    const queue = state.pendingSignalsByThread.get(key);\n    const signal = queue?.shift();\n    if (signal && queue) {\n      if (queue.length === 0) {\n        state.pendingSignalsByThread.delete(key);\n      }\n\n      // Hand the lease from the finished run to this drained run before\n      // streaming, so the lease key never goes empty during the handoff. If the\n      // old owner already lost the lease (e.g. a pubsub blip let the TTL lapse\n      // and another process took over), forward the signal to the new winner\n      // instead of starting a competing run here.\n      const nextRunId = randomUUID();\n      state.activeThreadRunIds.set(key, nextRunId);\n      state.threadKeysByRunId.set(nextRunId, key);\n      const owns = await this.#acquireOrTransferThreadLease(pubsub, key, nextRunId, previousRun.runId);\n      if (!owns.acquired) {\n        if (state.activeThreadRunIds.get(key) === nextRunId) {\n          state.activeThreadRunIds.delete(key);\n        }\n        state.threadKeysByRunId.delete(nextRunId);\n        // Early follow-ups were already published as retained signal-enqueued\n        // events, so only discard this runtime's local pre-run copies.\n        state.preRunSignalsByThread.delete(key);\n        // Put the signal back at the head so a later drain (or the winner) runs\n        // it, and forward it to the current lease owner.\n        const restored = state.pendingSignalsByThread.get(key) ?? [];\n        state.pendingSignalsByThread.set(key, [signal, ...restored]);\n        if (owns.owner) {\n          await this.#publishAndWait(pubsub, key, {\n            type: 'signal-enqueued',\n            runId: owns.owner,\n            signal: this.#serializeSignal(signal),\n            sourceId: this.#getSourceId(),\n          }).catch(() => {});\n          state.pendingSignalsByThread.get(key)?.shift();\n          if ((state.pendingSignalsByThread.get(key)?.length ?? 0) === 0) {\n            state.pendingSignalsByThread.delete(key);\n          }\n        }\n        return;\n      }\n\n      const output = await previousRun.agent.stream(signal, {\n        ...(previousRun.streamOptions as any),\n        runId: nextRunId,\n        memory: withThreadMemory(\n          previousRun.streamOptions.memory,\n          previousRun.resourceId ?? '',\n          previousRun.threadId ?? '',\n        ),\n      });\n\n      if (queue.length > 0) {\n        const nextRecord = state.threadRunsById.get(output.runId);\n        if (nextRecord) {\n          this.#watchThreadRunCompletion(state, pubsub, key, nextRecord);\n        }\n      }\n      return;\n    }\n\n    if (await this.#drainPendingContinuations(state, pubsub, key, previousRun.runId)) {\n      return;\n    }\n\n    if (await this.#drainPendingIdleSignals(state, pubsub, key, previousRun.runId)) {\n      return;\n    }\n\n    // Nothing left to drain: release the lease we kept held for the drain.\n    this.#releaseThreadLease(pubsub, key, previousRun.runId);\n  }\n\n  async #drainPendingContinuations(\n    state: AgentThreadRuntimeState,\n    pubsub: PubSub | undefined,\n    key: string,\n    fromRunId?: string,\n  ) {\n    if (state.activeThreadRunIds.has(key)) {\n      return false;\n    }\n\n    const queue = state.pendingContinuationsByThread.get(key);\n    const pending = queue?.shift();\n    if (!pending || !queue) {\n      return false;\n    }\n    if (queue.length === 0) {\n      state.pendingContinuationsByThread.delete(key);\n    }\n\n    // A continuation only ever drains in the process that owned the finished\n    // run, so it always carries a `fromRunId` to hand the held lease to. If the\n    // old owner already lost the lease, re-queue the continuation and let the\n    // new lease owner take over rather than starting a competing run here.\n    if (fromRunId) {\n      state.activeThreadRunIds.set(key, pending.runId);\n      state.threadKeysByRunId.set(pending.runId, key);\n      const owns = await this.#acquireOrTransferThreadLease(pubsub, key, pending.runId, fromRunId);\n      if (!owns.acquired) {\n        if (state.activeThreadRunIds.get(key) === pending.runId) {\n          state.activeThreadRunIds.delete(key);\n        }\n        state.threadKeysByRunId.delete(pending.runId);\n        state.preRunSignalsByThread.delete(key);\n        const restored = state.pendingContinuationsByThread.get(key) ?? [];\n        state.pendingContinuationsByThread.set(key, [pending, ...restored]);\n        return false;\n      }\n    }\n\n    this.#startContinuation(state, pubsub, key, pending);\n    return true;\n  }\n\n  #startContinuation(\n    state: AgentThreadRuntimeState,\n    pubsub: PubSub | undefined,\n    key: string,\n    pending: PendingContinuation<any>,\n  ) {\n    state.activeThreadRunIds.set(key, pending.runId);\n    state.threadKeysByRunId.set(pending.runId, key);\n    void pending.agent\n      .stream(pending.messages, {\n        ...(pending.streamOptions as any),\n        runId: pending.runId,\n        memory: withThreadMemory(pending.streamOptions?.memory, pending.resourceId, pending.threadId),\n      })\n      .then(output => {\n        if ((state.pendingContinuationsByThread.get(key)?.length ?? 0) > 0) {\n          const nextRecord = state.threadRunsById.get(output.runId);\n          if (nextRecord) {\n            this.#watchThreadRunCompletion(state, pubsub, key, nextRecord);\n          }\n        }\n      })\n      .catch(err => {\n        state.threadKeysByRunId.delete(pending.runId);\n        this.#cleanupPreparedRun(state, pending.runId);\n        if (state.activeThreadRunIds.get(key) === pending.runId) {\n          state.activeThreadRunIds.delete(key);\n        }\n        this.#publish(pubsub, key, {\n          type: 'run-failed',\n          runId: pending.runId,\n          error: getErrorFromUnknown(err).message,\n        });\n        // Hand the lease to remaining queued work (transfer keeps the key from\n        // going empty); only release once nothing is left to drain.\n        void this.#drainPendingContinuations(state, pubsub, key, pending.runId).then(async started => {\n          if (started) return;\n          if (await this.#drainPendingIdleSignals(state, pubsub, key, pending.runId)) return;\n          this.#releaseThreadLease(pubsub, key, pending.runId);\n        });\n      });\n  }\n\n  continueWithMessages<OUTPUT = unknown>(\n    agent: Agent<any, any, any, any>,\n    messages: MessageListInput,\n    target: { resourceId: string; threadId: string; streamOptions?: AgentExecutionOptions<OUTPUT>; runId?: string },\n    pubsub?: PubSub,\n  ): { accepted: true; runId: string } {\n    const state = this.#getState(pubsub);\n    const key = this.#threadKey(target.resourceId, target.threadId);\n    const runId = target.runId ?? randomUUID();\n    const pending: PendingContinuation<OUTPUT> = {\n      agent,\n      messages,\n      runId,\n      resourceId: target.resourceId,\n      threadId: target.threadId,\n      streamOptions: target.streamOptions,\n    };\n\n    const activeRunId = state.activeThreadRunIds.get(key);\n    const activeRecord = activeRunId ? state.threadRunsById.get(activeRunId) : undefined;\n    if (state.activeThreadRunIds.has(key)) {\n      const queue = state.pendingContinuationsByThread.get(key) ?? [];\n      queue.push(pending);\n      state.pendingContinuationsByThread.set(key, queue);\n      if (activeRecord) {\n        this.#watchThreadRunCompletion(state, pubsub, key, activeRecord);\n      }\n      return { accepted: true, runId };\n    }\n\n    this.#startContinuation(state, pubsub, key, pending);\n    return { accepted: true, runId };\n  }\n\n  async #drainPendingIdleSignals(\n    state: AgentThreadRuntimeState,\n    pubsub: PubSub | undefined,\n    key: string,\n    fromRunId?: string,\n  ): Promise<boolean> {\n    if (state.activeThreadRunIds.has(key)) {\n      return false;\n    }\n\n    const idleQueue = state.pendingIdleSignalsByThread.get(key);\n    const pendingIdle = idleQueue?.shift();\n    if (!pendingIdle || !idleQueue) {\n      return false;\n    }\n    if (idleQueue.length === 0) {\n      state.pendingIdleSignalsByThread.delete(key);\n    }\n\n    state.activeThreadRunIds.set(key, pendingIdle.runId);\n    state.threadKeysByRunId.set(pendingIdle.runId, key);\n\n    // A queued idle signal may be draining either in the process that just\n    // finished a run (it still holds the lease — hand it over) or in a\n    // *different* process that observed the owner's run finish via pub/sub and\n    // now wants to wake the thread (it holds no lease — it must win one). Either\n    // way the run must only start if this process owns the cross-process lease,\n    // otherwise two processes could each start a competing idle run.\n    const owns = await this.#acquireOrTransferThreadLease(pubsub, key, pendingIdle.runId, fromRunId);\n    if (!owns.acquired) {\n      // Lost the wake race. Roll back the optimistic local reservation and\n      // forward the signal to the winner so it is not dropped, then try the\n      // next queued idle signal (which may belong to a different run we can win).\n      if (state.activeThreadRunIds.get(key) === pendingIdle.runId) {\n        state.activeThreadRunIds.delete(key);\n      }\n      state.threadKeysByRunId.delete(pendingIdle.runId);\n      state.preRunSignalsByThread.delete(key);\n      if (owns.owner) {\n        await this.#publishAndWait(pubsub, key, {\n          type: 'signal-enqueued',\n          runId: owns.owner,\n          signal: this.#serializeSignal(pendingIdle.signal),\n          sourceId: this.#getSourceId(),\n        }).catch(() => {});\n      }\n      await this.#drainPendingIdleSignals(state, pubsub, key, fromRunId);\n      return true;\n    }\n\n    try {\n      const output = await pendingIdle.agent.stream(pendingIdle.signal, {\n        ...(pendingIdle.streamOptions as any),\n        runId: pendingIdle.runId,\n        memory: withThreadMemory(pendingIdle.streamOptions?.memory, pendingIdle.resourceId, pendingIdle.threadId),\n      });\n\n      if ((idleQueue?.length ?? 0) > 0) {\n        const nextRecord = state.threadRunsById.get(output.runId);\n        if (nextRecord) {\n          this.#watchThreadRunCompletion(state, pubsub, key, nextRecord);\n        }\n      }\n    } catch (err) {\n      state.threadKeysByRunId.delete(pendingIdle.runId);\n      this.#cleanupPreparedRun(state, pendingIdle.runId);\n      if (state.activeThreadRunIds.get(key) === pendingIdle.runId) {\n        state.activeThreadRunIds.delete(key);\n      }\n      this.#publish(pubsub, key, {\n        type: 'run-failed',\n        runId: pendingIdle.runId,\n        error: getErrorFromUnknown(err).message,\n      });\n      // Hand the lease to remaining idle work; release only when none remains.\n      if (!(await this.#drainPendingIdleSignals(state, pubsub, key, pendingIdle.runId))) {\n        this.#releaseThreadLease(pubsub, key, pendingIdle.runId);\n      }\n    }\n    return true;\n  }\n\n  /**\n   * Drains queued signals for a run.\n   *\n   * - `scope: 'pending'` (default) returns active-run follow-up signals — each\n   *   becomes its own model turn via `signalDrainStep`.\n   * - `scope: 'pre-run'` returns signals queued before the run's first model\n   *   request — the first LLM step folds these into that request.\n   */\n  drainPendingSignals(runId: string, pubsub?: PubSub, scope: 'pending' | 'pre-run' = 'pending'): CreatedAgentSignal[] {\n    const state = this.#getState(pubsub);\n    const record = state.threadRunsById.get(runId);\n    const key = record ? this.#threadKey(record.resourceId, record.threadId) : state.threadKeysByRunId.get(runId);\n    if (!key) return [];\n\n    const signalsByThread = scope === 'pre-run' ? state.preRunSignalsByThread : state.pendingSignalsByThread;\n    const queue = signalsByThread.get(key);\n    if (!queue || queue.length === 0) {\n      return [];\n    }\n\n    signalsByThread.delete(key);\n    return queue;\n  }\n\n  async waitForCrossAgentThreadRun(\n    agent: Agent<any, any, any, any>,\n    options: { memory?: AgentExecutionOptions<any>['memory']; requestContext?: RequestContext },\n    pubsub?: PubSub,\n  ) {\n    const { threadId, resourceId } = this.#getThreadTarget(options);\n    if (!threadId) return;\n\n    const state = this.#getState(pubsub);\n    const key = this.#threadKey(resourceId, threadId);\n    while (true) {\n      const activeRunId = state.activeThreadRunIds.get(key);\n      if (!activeRunId) return;\n\n      const activeRecord = state.threadRunsById.get(activeRunId);\n      if (activeRecord) {\n        if (activeRecord.agent.id === agent.id || !this.#isThreadBlockingRun(state, activeRecord)) {\n          return;\n        }\n        await activeRecord.output._waitUntilFinished().catch(() => {});\n        continue;\n      }\n\n      if (state.threadKeysByRunId.get(activeRunId) === key) return;\n\n      await this.#waitForRemoteRunToFinish(pubsub, key, activeRunId);\n    }\n  }\n\n  async #waitForRemoteRunToFinish(pubsub: PubSub | undefined, key: string, runId: string) {\n    const resolvedPubSub = this.#getPubSub(pubsub);\n    const { provider, isFallback } = this.#resolveLeaseProvider(resolvedPubSub);\n    const topic = this.#threadTopic(key);\n    let timer: ReturnType<typeof setTimeout> | undefined;\n    let subscribed = false;\n    let settled = false;\n    let resolveWait!: () => void;\n    const wait = new Promise<void>(resolve => {\n      resolveWait = resolve;\n    });\n    const clearRemoteActive = (streamId?: string) => {\n      const state = this.#getState(resolvedPubSub);\n      if (\n        state.activeThreadRunIds.get(key) !== runId ||\n        (streamId && state.activeThreadStreamIds.get(key) !== streamId)\n      ) {\n        return;\n      }\n      state.activeThreadRunIds.delete(key);\n      state.activeThreadStreamIds.delete(key);\n      if (state.remoteThreadKeysByRunId.get(runId) === key) state.remoteThreadKeysByRunId.delete(runId);\n    };\n    const finish = () => {\n      if (settled) return;\n      settled = true;\n      if (timer) clearTimeout(timer);\n      resolveWait();\n    };\n    const checkLease = async () => {\n      if (settled) return;\n      if (isFallback) return;\n      const owner = await provider.getLeaseOwner(key).catch(() => undefined);\n      if (settled) return;\n      if (owner !== runId) {\n        clearRemoteActive();\n        finish();\n        return;\n      }\n      timer = setTimeout(() => void checkLease(), AGENT_THREAD_LEASE_TTL_MS);\n    };\n    const onEvent: EventCallback = event => {\n      const data = event.data as AgentThreadStreamRuntimeEvent | undefined;\n      if (\n        (data?.type === 'run-completed' || data?.type === 'run-aborted' || data?.type === 'run-failed') &&\n        data.runId === runId\n      ) {\n        clearRemoteActive(data.streamId);\n        finish();\n      }\n    };\n\n    try {\n      await resolvedPubSub.subscribe(topic, onEvent);\n      subscribed = true;\n      if (!isFallback) timer = setTimeout(() => void checkLease(), AGENT_THREAD_LEASE_TTL_MS);\n      await wait;\n    } catch {\n      finish();\n      await wait;\n    } finally {\n      if (timer) clearTimeout(timer);\n      if (subscribed) await resolvedPubSub.unsubscribe(topic, onEvent).catch(() => {});\n    }\n  }\n\n  async subscribeToThread<OUTPUT = unknown>(\n    agent: Agent<any, any, any, any>,\n    options: AgentSubscribeToThreadOptions,\n    pubsub?: PubSub,\n  ): Promise<AgentThreadSubscription<OUTPUT>> {\n    void agent;\n    const resolvedPubSub = this.#getPubSub(pubsub);\n    const state = this.#getState(resolvedPubSub);\n    const key = this.#threadKey(options.resourceId, options.threadId);\n    const topic = this.#threadTopic(key);\n    const seenStreamIds = new Set<string>();\n    const pendingRuns: AgentThreadRunRecord<any>[] = [];\n    const waiters: Array<() => void> = [];\n    const remoteRuns = new Map<\n      string,\n      {\n        parts: unknown[];\n        waiters: Array<() => void>;\n        finishWaiters: Array<() => void>;\n        done: boolean;\n        stream: ReadableStream<unknown>;\n      }\n    >();\n    let done = false;\n\n    const wake = () => {\n      while (waiters.length) waiters.shift()?.();\n    };\n\n    const activeRunId = () => {\n      const runId = state.activeThreadRunIds.get(key);\n      if (!runId) return null;\n      const record = state.threadRunsById.get(runId);\n      // No record yet means either a remote run (record never lives locally) or a local run\n      // that sendSignal has reserved but has not yet registered via registerRun. Both are\n      // in flight from the subscriber's perspective; treat them as active.\n      if (!record) return runId;\n      return this.#isThreadBlockingRun(state, record) ? runId : null;\n    };\n\n    const enqueueRun = (record: AgentThreadRunRecord<any>) => {\n      if (done || seenStreamIds.has(record.streamId)) return;\n      seenStreamIds.add(record.streamId);\n      pendingRuns.push(record);\n      wake();\n    };\n\n    const createRemoteRun = (runId: string, streamId: string, streamSeq: number): AgentThreadRunRecord<any> => {\n      const remoteRun = {\n        parts: [] as unknown[],\n        waiters: [] as Array<() => void>,\n        finishWaiters: [] as Array<() => void>,\n        done: false,\n        stream: undefined as unknown as ReadableStream<unknown>,\n        closed: false,\n      };\n      remoteRun.stream = new ReadableStream({\n        pull(controller) {\n          const drain = () => {\n            if (remoteRun.closed) return;\n            while (remoteRun.parts.length > 0) {\n              controller.enqueue(remoteRun.parts.shift());\n            }\n            if (remoteRun.done) {\n              remoteRun.closed = true;\n              controller.close();\n            }\n          };\n          drain();\n          if (!remoteRun.done && !remoteRun.closed) {\n            remoteRun.waiters.push(drain);\n          }\n        },\n        cancel() {\n          remoteRun.done = true;\n          remoteRun.closed = true;\n          remoteRun.waiters.length = 0;\n          while (remoteRun.finishWaiters.length) remoteRun.finishWaiters.shift()?.();\n        },\n      });\n      remoteRuns.set(streamId, remoteRun);\n      return {\n        agent,\n        output: {\n          runId,\n          status: 'running',\n          fullStream: remoteRun.stream,\n          _waitUntilFinished: async () => {\n            if (remoteRun.done) return;\n            await new Promise<void>(resolve => remoteRun.finishWaiters.push(resolve));\n          },\n        } as MastraModelOutput<any>,\n        runId,\n        streamId,\n        streamSeq,\n        lifecycle: 'running',\n        threadId: options.threadId,\n        resourceId: options.resourceId,\n        streamOptions: {},\n      };\n    };\n\n    const localStreamIds = new Set<string>();\n    const replayedStreamIds = new Set<string>();\n    let currentReader: ReadableStreamDefaultReader<any> | null = null;\n    let activeReaderRunId: string | null = null;\n    let cancelledByAbort = false;\n\n    const markActiveIfLive = async (runId: string, streamId: string, local: boolean) => {\n      if (!local && !(await this.#hasLiveThreadLease(resolvedPubSub, key, runId))) return;\n      state.activeThreadRunIds.set(key, runId);\n      state.activeThreadStreamIds.set(key, streamId);\n      if (!local) state.remoteThreadKeysByRunId.set(runId, key);\n    };\n\n    const clearActiveIfCurrent = (runId: string, streamId?: string) => {\n      if (\n        state.activeThreadRunIds.get(key) !== runId ||\n        (streamId && state.activeThreadStreamIds.get(key) !== streamId)\n      ) {\n        return;\n      }\n      state.activeThreadRunIds.delete(key);\n      state.activeThreadStreamIds.delete(key);\n      if (state.remoteThreadKeysByRunId.get(runId) === key) state.remoteThreadKeysByRunId.delete(runId);\n    };\n\n    const handleEvent = async (event: Parameters<EventCallback>[0]) => {\n      if (done) return;\n      const data = event.data as AgentThreadStreamRuntimeEvent | undefined;\n      if (!data) return;\n      if (data.type === 'run-registered') {\n        const localRecord = state.threadRunsByStreamId.get(data.streamId);\n        if (localRecord) {\n          localStreamIds.add(data.streamId);\n        } else {\n          replayedStreamIds.add(data.streamId);\n        }\n        await markActiveIfLive(data.runId, data.streamId, Boolean(localRecord));\n        const record = localRecord ?? createRemoteRun(data.runId, data.streamId, data.streamSeq);\n        enqueueRun(record);\n        wake();\n        return;\n      }\n      if (data.type === 'stream-part') {\n        if (\n          data.sourceId === this.#id &&\n          (localStreamIds.has(data.streamId) || !replayedStreamIds.has(data.streamId))\n        ) {\n          return;\n        }\n        if (\n          state.activeThreadRunIds.get(key) !== data.runId ||\n          state.activeThreadStreamIds.get(key) !== data.streamId\n        ) {\n          await markActiveIfLive(data.runId, data.streamId, false);\n        }\n        let remoteRun = remoteRuns.get(data.streamId);\n        if (!remoteRun) {\n          // A subscriber can attach after another runtime already broadcast run-registered.\n          // Treat the first stream-part on this thread topic as proof of the remote run and\n          // create the local proxy stream from that point forward.\n          enqueueRun(createRemoteRun(data.runId, data.streamId, state.streamSeqByRunId.get(data.runId) ?? 1));\n          remoteRun = remoteRuns.get(data.streamId);\n          if (!remoteRun) return;\n        }\n        remoteRun.parts.push(data.part);\n        while (remoteRun.waiters.length) remoteRun.waiters.shift()?.();\n        return;\n      }\n      if (data.type === 'signal-enqueued') {\n        if (data.sourceId === this.#id) return;\n        const signalsByThread = data.preRun ? state.preRunSignalsByThread : state.pendingSignalsByThread;\n        const queue = signalsByThread.get(key) ?? [];\n        queue.push(createSignal(data.signal));\n        signalsByThread.set(key, queue);\n        return;\n      }\n      if (data.type === 'run-abort-requested') {\n        if (\n          state.preparedRunsById.has(data.runId) &&\n          state.threadKeysByRunId.get(data.runId) === key &&\n          state.activeThreadRunIds.get(key) === data.runId &&\n          state.activeThreadStreamIds.get(key) === data.streamId &&\n          (await this.#hasLiveThreadLease(resolvedPubSub, key, data.runId))\n        ) {\n          this.abortRun(data.runId, resolvedPubSub);\n        }\n        return;\n      }\n      if (data.type === 'run-failed') {\n        const eventStreamId = data.streamId ?? data.runId;\n        clearActiveIfCurrent(data.runId, data.streamId);\n        let errorRun: AgentThreadRunRecord<any> | undefined;\n        let remoteRun = remoteRuns.get(eventStreamId);\n        if (!remoteRun) {\n          errorRun = createRemoteRun(data.runId, eventStreamId, state.streamSeqByRunId.get(data.runId) ?? 1);\n          remoteRun = remoteRuns.get(eventStreamId);\n        }\n        if (remoteRun) {\n          remoteRun.parts.push({ type: 'error', payload: { error: new Error(data.error) } });\n          remoteRun.done = true;\n          while (remoteRun.waiters.length) remoteRun.waiters.shift()?.();\n          while (remoteRun.finishWaiters.length) remoteRun.finishWaiters.shift()?.();\n          remoteRuns.delete(eventStreamId);\n          seenStreamIds.delete(eventStreamId);\n        }\n        if (errorRun) enqueueRun(errorRun);\n        await this.#drainPendingIdleSignals(state, resolvedPubSub, key, data.runId);\n        wake();\n        return;\n      }\n      if (data.type === 'run-completed' || data.type === 'run-aborted' || data.type === 'run-suspended') {\n        const eventStreamId = data.streamId ?? data.runId;\n        if (data.type === 'run-suspended') {\n          state.suspendedRunIds.add(data.runId);\n          const record = state.threadRunsByStreamId.get(eventStreamId) ?? state.threadRunsById.get(data.runId);\n          if (record) record.lifecycle = 'suspended';\n        } else {\n          clearActiveIfCurrent(data.runId, data.streamId);\n        }\n        if (data.type !== 'run-suspended') {\n          this.#clearSuspendedRun(state, data.runId);\n        }\n        const remoteRun = remoteRuns.get(eventStreamId);\n        if (remoteRun) {\n          remoteRun.done = true;\n          while (remoteRun.waiters.length) remoteRun.waiters.shift()?.();\n          while (remoteRun.finishWaiters.length) remoteRun.finishWaiters.shift()?.();\n          remoteRuns.delete(eventStreamId);\n          seenStreamIds.delete(eventStreamId);\n        }\n        // When a run is aborted, cancel the current subscriber stream reader so\n        // the generator's inner loop unblocks and can yield the synthetic abort.\n        if (data.type === 'run-aborted' && activeReaderRunId === data.runId && currentReader) {\n          cancelledByAbort = true;\n          try {\n            void currentReader.cancel();\n          } catch {}\n        }\n        if (data.type !== 'run-suspended') {\n          await this.#drainPendingIdleSignals(state, resolvedPubSub, key, data.runId);\n        }\n        wake();\n      }\n    };\n\n    let eventTail = Promise.resolve();\n    const onEvent: EventCallback = event => {\n      eventTail = eventTail.then(() => handleEvent(event)).catch(() => {});\n    };\n\n    await resolvedPubSub.subscribe(topic, onEvent);\n\n    const currentRunId = activeRunId();\n    const currentRecord = currentRunId ? state.threadRunsById.get(currentRunId) : undefined;\n    if (currentRecord) {\n      localStreamIds.add(currentRecord.streamId);\n      enqueueRun(currentRecord);\n    }\n\n    const unsubscribe = () => {\n      if (done) return;\n      done = true;\n      void resolvedPubSub.unsubscribe(topic, onEvent).catch(() => {});\n      // Cancel current reader so the generator's inner loop breaks.\n      if (currentReader) {\n        try {\n          void currentReader.cancel();\n        } catch {}\n      }\n      wake();\n    };\n\n    return {\n      activeRunId,\n      abort: () => this.abortThread(options, resolvedPubSub),\n      unsubscribe,\n      stream: (async function* () {\n        try {\n          while (!done || pendingRuns.length > 0) {\n            if (pendingRuns.length === 0) {\n              await new Promise<void>(resolve => waiters.push(resolve));\n              continue;\n            }\n            const run = pendingRuns.shift()!;\n            // Local registered runs expose createSubscriberStream, while remote runs are\n            // already per-subscription streams. Do not silently skip locked streams here:\n            // a locked fallback stream means a caller is sharing a non-multicast stream.\n            const subscriberStream = run.createSubscriberStream?.() ?? run.output.fullStream;\n            const reader = subscriberStream.getReader();\n            currentReader = reader as ReadableStreamDefaultReader<any>;\n            activeReaderRunId = run.runId;\n            let readerReleased = false;\n            try {\n              while (true) {\n                const { value: part, done: streamDone } = await reader.read();\n                if (streamDone) {\n                  break;\n                }\n                const typedPart = part as any;\n                const partWithRunId =\n                  typedPart && typeof typedPart === 'object' && !('runId' in typedPart)\n                    ? { ...typedPart, runId: run.runId }\n                    : typedPart;\n                yield partWithRunId;\n                if (done) break;\n                const finishReason = typedPart.finishReason ?? typedPart.payload?.finishReason;\n                const terminalBoundary =\n                  typedPart.type === 'error' ||\n                  typedPart.type === 'abort' ||\n                  (typedPart.type === 'finish' && finishReason !== 'tool-calls');\n                if (terminalBoundary) {\n                  // After a final terminal chunk, drain any non-visible trailing\n                  // data in the background to prevent upstream backpressure while\n                  // allowing the generator to immediately serve subsequent runs.\n                  readerReleased = true;\n                  void (async () => {\n                    try {\n                      while (true) {\n                        const { done: d } = await reader.read();\n                        if (d) break;\n                      }\n                    } catch {}\n                    reader.releaseLock();\n                  })();\n                  break;\n                }\n              }\n              // If the stream closed because we cancelled the reader after a\n              // run-aborted event, yield a synthetic abort so subscribers\n              // finalize the run.\n              if (!readerReleased && !done && cancelledByAbort) {\n                yield { type: 'abort', runId: run.runId } as any;\n                cancelledByAbort = false;\n              }\n            } finally {\n              currentReader = null;\n              activeReaderRunId = null;\n              if (!readerReleased) {\n                reader.releaseLock();\n              }\n            }\n          }\n        } finally {\n          unsubscribe();\n        }\n      })(),\n    };\n  }\n\n  sendMessage<OUTPUT = unknown>(\n    agent: Agent<any, any, any, any>,\n    message: AgentMessageInput,\n    target: SendAgentMessageOptions<OUTPUT>,\n    pubsub?: PubSub,\n  ): SendAgentMessageResult<OUTPUT> {\n    return this.sendSignal<OUTPUT>(agent, this.#createMessageSignalInput(message), target, pubsub);\n  }\n\n  queueMessage<OUTPUT = unknown>(\n    agent: Agent<any, any, any, any>,\n    message: AgentMessageInput,\n    target: QueueAgentMessageOptions<OUTPUT>,\n    pubsub?: PubSub,\n  ): QueueAgentMessageResult<OUTPUT> {\n    const state = this.#getState(pubsub);\n    const acceptedAt = new Date();\n    let key: string | undefined;\n    let runId = target.runId;\n    let activeRecord: AgentThreadRunRecord<any> | undefined;\n\n    if (target.resourceId && target.threadId) {\n      key = this.#threadKey(target.resourceId, target.threadId);\n      const activeRunId = state.activeThreadRunIds.get(key);\n      activeRecord = activeRunId ? state.threadRunsById.get(activeRunId) : undefined;\n      if (activeRecord && !this.#isThreadBlockingRun(state, activeRecord)) {\n        state.activeThreadRunIds.delete(key);\n        activeRecord = undefined;\n      }\n      runId ??= activeRunId;\n    }\n\n    if (runId) {\n      activeRecord ??= state.threadRunsById.get(runId);\n      if (activeRecord) {\n        key ??= this.#threadKey(activeRecord.resourceId, activeRecord.threadId);\n      }\n    }\n\n    const resourceId = target.resourceId ?? activeRecord?.resourceId;\n    const threadId = target.threadId ?? activeRecord?.threadId;\n    if (!resourceId || !threadId) {\n      throw new Error('resourceId and threadId are required to queue a message');\n    }\n\n    key ??= this.#threadKey(resourceId, threadId);\n    const signal = createMessageSignal(message, {\n      id: this.#generateSignalMessageId(agent, { resourceId, threadId }),\n      acceptedAt,\n    });\n    const queuedRunId = randomUUID();\n    const queuedStreamOptions = target.ifIdle?.streamOptions ?? activeRecord?.streamOptions;\n\n    if (activeRecord) {\n      const idleQueue = state.pendingIdleSignalsByThread.get(key) ?? [];\n      idleQueue.push({ agent, signal, runId: queuedRunId, resourceId, threadId, streamOptions: queuedStreamOptions });\n      state.pendingIdleSignalsByThread.set(key, idleQueue);\n      this.#watchThreadRunCompletion(state, pubsub, key, activeRecord);\n      return {\n        signal,\n        accepted: Promise.resolve({ action: 'deliver' as const, runId: queuedRunId }),\n      };\n    }\n\n    return this.sendSignal<OUTPUT>(\n      agent,\n      signal,\n      { ...target, runId, resourceId, threadId, ifIdle: { ...target.ifIdle, behavior: 'wake' } },\n      pubsub,\n    );\n  }\n\n  async sendStateSignal<OUTPUT = unknown>(\n    agent: Agent<any, any, any, any>,\n    stateInput: AgentStateSignalInput,\n    target: SendAgentStateSignalOptions<OUTPUT>,\n    pubsub?: PubSub,\n  ): Promise<SendAgentStateSignalResult<OUTPUT>> {\n    if (!target.resourceId || !target.threadId) {\n      throw new Error('resourceId and threadId are required to send a state signal');\n    }\n    const resourceId = target.resourceId;\n    const threadId = target.threadId;\n\n    const requestContext = target.ifIdle?.streamOptions?.requestContext;\n    const memoryContext = parseMemoryRequestContext(requestContext);\n    const memory = await agent.getMemory({ requestContext });\n    if (!memory) {\n      throw new Error('sendStateSignal requires Mastra memory');\n    }\n\n    const loadedThread = (await memory.getThreadById({ threadId })) ?? memoryContext?.thread;\n    if (!loadedThread) {\n      throw new Error(`sendStateSignal could not load thread ${threadId}`);\n    }\n\n    const thread = {\n      ...loadedThread,\n      id: threadId,\n      resourceId: loadedThread.resourceId ?? resourceId,\n      createdAt: loadedThread.createdAt ?? new Date(),\n      updatedAt: loadedThread.updatedAt ?? new Date(),\n      metadata: loadedThread.metadata,\n    };\n\n    const applied = await applyStateSignal({\n      input: stateInput,\n      memory,\n      thread,\n      resourceId,\n      threadId,\n      memoryConfig: memoryContext?.memoryConfig,\n      acceptedAt: new Date(),\n    });\n\n    if (applied.skipped) {\n      return { skipped: true, reason: 'unchanged' };\n    }\n\n    return this.sendSignal<OUTPUT>(agent, applied.signal, target, pubsub);\n  }\n\n  /**\n   * Routes a signal to an agent thread.\n   *\n   * Signals can land in three places:\n   * - an active same-agent run, where they are queued for the execution loop to drain;\n   * - a reserved thread run that has not registered its stream record yet;\n   * - a new idle-started run, when the caller opts into `ifIdle`.\n   *\n   * Cross-agent active runs are intentionally not interrupted here. They either finish first\n   * through `waitForCrossAgentThreadRun()` on the stream path, or this method falls through to\n   * the idle-start path when the caller provided a resource/thread target and `ifIdle` options.\n   */\n  sendSignal<OUTPUT = unknown>(\n    agent: Agent<any, any, any, any>,\n    signalInput: AgentSignal,\n    target: SendAgentSignalOptions<OUTPUT>,\n    pubsub?: PubSub,\n  ): SendAgentSignalResult<OUTPUT> {\n    const state = this.#getState(pubsub);\n    let key: string | undefined;\n    let runId = target.runId;\n    const activeBehavior = target.ifActive?.behavior ?? 'deliver';\n    const idleBehavior = target.ifIdle?.behavior ?? 'wake';\n\n    let activeRecord: AgentThreadRunRecord<any> | undefined;\n    if (target.resourceId && target.threadId) {\n      key = this.#threadKey(target.resourceId, target.threadId);\n      const activeRunId = state.activeThreadRunIds.get(key);\n      activeRecord = activeRunId ? state.threadRunsById.get(activeRunId) : undefined;\n      if (activeRecord && !this.#isThreadBlockingRun(state, activeRecord)) {\n        state.activeThreadRunIds.delete(key);\n        activeRecord = undefined;\n      }\n\n      // Prefer the active same-agent run for thread-targeted signals. This is the normal\n      // follow-up path used by clients that know the thread/resource but not the run id.\n      if (activeRecord && activeRecord.agent.id === agent.id) {\n        runId = activeRecord.runId;\n      } else if (activeRunId && !activeRecord) {\n        if (state.threadKeysByRunId.get(activeRunId) === key) {\n          // A run can be reserved before its stream record is registered. Keep the reserved\n          // id so early follow-ups still attach to the run that is starting.\n          runId = activeRunId;\n        } else {\n          // Stale cross-pod entry. Clean it up from the local map, then let the lease decide\n          state.activeThreadRunIds.delete(key);\n          state.activeThreadStreamIds.delete(key);\n        }\n      }\n    }\n\n    if (runId) {\n      activeRecord ??= state.threadRunsById.get(runId);\n      if (activeRecord) {\n        key ??= this.#threadKey(activeRecord.resourceId, activeRecord.threadId);\n      }\n    }\n\n    const resourceId = target.resourceId ?? activeRecord?.resourceId;\n    const threadId = target.threadId ?? activeRecord?.threadId;\n    if (!resourceId || !threadId) {\n      throw new Error('No active agent run found for signal target');\n    }\n\n    const isActiveTarget = Boolean(\n      runId && (activeRecord?.output.status === 'running' || (key && state.activeThreadRunIds.get(key) === runId)),\n    );\n    let signal = createSignal({\n      ...signalInput,\n      id: signalInput.id ?? this.#generateSignalMessageId(agent, { resourceId, threadId }),\n      acceptedAt: new Date(),\n    });\n\n    // Resolve conditional delivery attributes now that we know the delivery path.\n    signal = resolveDeliveryAttributes(\n      signal,\n      isActiveTarget ? target.ifActive?.attributes : target.ifIdle?.attributes,\n    );\n\n    if (isActiveTarget && activeBehavior !== 'deliver') {\n      if (activeBehavior === 'persist') {\n        if (!resourceId || !threadId) {\n          throw new Error('resourceId and threadId are required to persist an active signal');\n        }\n        // Transient signals are never written to storage, so a `persist` behavior has nothing\n        // to do with them — report the drop honestly as `discard` instead of `persist`.\n        if (signal.transient) {\n          return {\n            signal,\n            accepted: Promise.resolve({ action: 'discard' as const }),\n          };\n        }\n        const persisted = this.#persistSignal(\n          agent,\n          signal,\n          resourceId,\n          threadId,\n          target.ifIdle?.streamOptions?.requestContext,\n        );\n        void persisted.catch(() => {});\n        return {\n          signal,\n          persisted,\n          accepted: Promise.resolve({ action: 'persist' as const }),\n        };\n      }\n      return {\n        signal,\n        accepted: Promise.resolve({ action: 'discard' as const }),\n      };\n    }\n\n    if (runId) {\n      // A run is \"blocking\" while it is running or suspended awaiting tool approval. Both\n      // states mean the run has already made model requests, so a follow-up signal must be\n      // queued as a pending (next-turn) signal rather than folded into a not-yet-started\n      // first request via the pre-run path below.\n      if (activeRecord && this.#isThreadBlockingRun(state, activeRecord)) {\n        key ??= this.#threadKey(activeRecord.resourceId, activeRecord.threadId);\n        if (activeRecord.agent.id === agent.id) {\n          // Same-agent active run: queue the signal for in-loop draining so it becomes\n          // the next model input instead of waiting for the run to finish.\n          const queue = state.pendingSignalsByThread.get(key) ?? [];\n          queue.push(signal);\n          state.pendingSignalsByThread.set(key, queue);\n          this.#publish(pubsub, key, {\n            type: 'signal-enqueued',\n            runId,\n            signal: this.#serializeSignal(signal),\n            sourceId: this.#getSourceId(),\n          });\n          this.#watchThreadRunCompletion(state, pubsub, key, activeRecord);\n          return {\n            signal,\n            accepted: Promise.resolve({ action: 'deliver' as const, runId }),\n          };\n        }\n\n        return {\n          signal,\n          accepted: Promise.resolve({\n            action: 'blocked' as const,\n            reason: 'thread-blocked' as const,\n            runId: activeRecord.runId,\n          }),\n        };\n      }\n\n      if (key && state.activeThreadRunIds.get(key) === runId) {\n        // A local reserved run has not registered its stream record yet, so it\n        // has not made its first model request — queue the signal as a pre-run\n        // signal so the first LLM step folds it into that request. A run owned\n        // by another runtime instance is reached only via PubSub; treat it as a\n        // follow-up, since the sender cannot see the owner's request state.\n        const isLocalReservedRun = state.threadKeysByRunId.get(runId) === key;\n        if (isLocalReservedRun) {\n          const queue = state.preRunSignalsByThread.get(key) ?? [];\n          queue.push(signal);\n          state.preRunSignalsByThread.set(key, queue);\n        }\n        this.#publish(pubsub, key, {\n          type: 'signal-enqueued',\n          runId,\n          signal: this.#serializeSignal(signal),\n          sourceId: this.#getSourceId(),\n          preRun: isLocalReservedRun,\n        });\n        return {\n          signal,\n          accepted: Promise.resolve({ action: 'deliver' as const, runId }),\n        };\n      }\n    }\n\n    runId = randomUUID();\n    key ??= this.#threadKey(resourceId, threadId);\n    if (idleBehavior === 'persist') {\n      // Transient signals are never written to storage, so an idle `persist` behavior drops\n      // them entirely (no store, no broadcast) — report that as `discard`, not `persist`.\n      if (signal.transient) {\n        return {\n          signal,\n          accepted: Promise.resolve({ action: 'discard' as const }),\n        };\n      }\n      const persisted = this.#persistAndBroadcastIdleSignal(\n        state,\n        pubsub,\n        key,\n        runId,\n        agent,\n        signal,\n        resourceId,\n        threadId,\n        target.ifIdle?.streamOptions?.requestContext,\n      );\n      void persisted.catch(() => {});\n      return {\n        signal,\n        persisted,\n        accepted: Promise.resolve({ action: 'persist' as const }),\n      };\n    }\n    if (idleBehavior !== 'wake') {\n      return {\n        signal,\n        accepted: Promise.resolve({ action: 'discard' as const }),\n      };\n    }\n\n    if (state.activeThreadRunIds.has(key)) {\n      const blockingRunId = state.activeThreadRunIds.get(key)!;\n      const blockingRecord = activeRecord ?? state.threadRunsById.get(blockingRunId);\n      if (\n        this.#isSuspendedRun(state, blockingRunId) ||\n        blockingRecord?.output.status === 'suspended' ||\n        blockingRecord?.lifecycle === 'suspended'\n      ) {\n        return {\n          signal,\n          accepted: Promise.resolve({\n            action: 'blocked' as const,\n            reason: 'thread-blocked' as const,\n            runId: blockingRunId,\n          }),\n        };\n      }\n\n      // Another run owns the thread. Queue this idle-start request and let the watcher\n      // launch it only after the active run clears the thread reservation.\n      const idleQueue = state.pendingIdleSignalsByThread.get(key) ?? [];\n      idleQueue.push({ agent, signal, runId, resourceId, threadId, streamOptions: target.ifIdle?.streamOptions });\n      state.pendingIdleSignalsByThread.set(key, idleQueue);\n      if (activeRecord) {\n        this.#watchThreadRunCompletion(state, pubsub, key, activeRecord);\n      }\n      return {\n        signal,\n        accepted: Promise.resolve({ action: 'deliver' as const, runId }),\n      };\n    }\n\n    // No active same-agent run accepted the signal. Reserve the thread before starting\n    // the idle stream so concurrent callers do not launch duplicate runs.\n    state.activeThreadRunIds.set(key, runId);\n    state.threadKeysByRunId.set(runId, key);\n    const reservedKey = key;\n    const reservedRunId = runId;\n    const resolvedPubSub = this.#getPubSub(pubsub);\n    const leaseProvider = this.#getLeaseProvider(resolvedPubSub);\n    // First acquire the cross-process lease via pubsub; on win, kick off the stream and\n    // resolve a `wake` accepted result carrying the owned stream. On loss, hand the user\n    // signal off to the winning process via signal-enqueued and resolve a `deliver` result\n    // (the signal was queued onto the winning run, not run locally).\n    const accepted: Promise<SendAgentSignalAccepted<OUTPUT>> = (async () => {\n      // Fail-open on pubsub errors: if the lease backend is unreachable we treat the\n      // call as \"acquired\" so the caller still gets a response. The tradeoff is that\n      // if multiple processes hit the same pubsub failure simultaneously they can each\n      // start a stream for the same thread (the bug this lease is supposed to prevent),\n      // but failing closed would silently drop user messages on any Redis blip which\n      // is the worse failure mode. Lease TTL + renewal still bound the duplicate\n      // window to a single run, and the next clean acquireLease re-serializes callers.\n      const lease = await leaseProvider\n        .acquireLease(reservedKey, reservedRunId, AGENT_THREAD_LEASE_TTL_MS)\n        .catch(() => ({ acquired: true as boolean, owner: reservedRunId as string | undefined }));\n\n      if (!lease.acquired) {\n        // Lost the wake race to another process. Roll back our optimistic local reservation\n        // so we don't trip our own activeThreadRunIds check on a follow-up.\n        if (state.activeThreadRunIds.get(reservedKey) === reservedRunId) {\n          state.activeThreadRunIds.delete(reservedKey);\n        }\n        state.threadKeysByRunId.delete(reservedRunId);\n        state.preRunSignalsByThread.delete(reservedKey);\n\n        // Forward the user signal to the winning runId so the message is not dropped.\n        // Await the publish so that callers using `accepted` resolution as their\n        // \"safe to exit\" boundary (e.g. a serverless Lambda holding the request open\n        // via waitUntil) don't tear down before the enqueue lands on the broker.\n        const winnerRunId = lease.owner;\n        if (winnerRunId) {\n          await this.#publishAndWait(pubsub, reservedKey, {\n            type: 'signal-enqueued',\n            runId: winnerRunId,\n            signal: this.#serializeSignal(signal),\n            sourceId: this.#getSourceId(),\n          }).catch(() => {});\n        }\n        return { action: 'deliver' as const, runId: winnerRunId ?? reservedRunId };\n      }\n\n      // We own the lease. Start the renewal timer so it survives runs\n      // that outlive the TTL, then kick off the stream.\n      this.#startLeaseRenewal(resolvedPubSub, reservedKey, reservedRunId);\n      try {\n        const output = await agent.stream(signal, {\n          ...(target.ifIdle?.streamOptions as any),\n          untilIdle: true,\n          runId: reservedRunId,\n          memory: withThreadMemory(target.ifIdle?.streamOptions?.memory, resourceId, threadId),\n        });\n        return { action: 'wake' as const, runId: reservedRunId, output };\n      } catch (error) {\n        state.threadKeysByRunId.delete(reservedRunId);\n        this.#cleanupPreparedRun(state, reservedRunId);\n        if (state.activeThreadRunIds.get(reservedKey) === reservedRunId) {\n          state.activeThreadRunIds.delete(reservedKey);\n        }\n        this.#releaseThreadLease(pubsub, reservedKey, reservedRunId);\n        this.#publish(pubsub, reservedKey, {\n          type: 'run-failed',\n          runId: reservedRunId,\n          error: getErrorFromUnknown(error).message,\n        });\n        void this.#drainPendingIdleSignals(state, pubsub, reservedKey);\n        throw error;\n      }\n    })();\n    // Attach a detached no-op catch so that if stream setup throws (a misconfigured\n    // agent: no/unsupported model, FGA denial) and the caller never awaits\n    // `result.accepted`, the rejection does not surface as an unhandled rejection.\n    // Callers that opt in to `accepted` still see the rejection via their own\n    // await/catch — `accepted` itself remains rejectable; only this detached branch is\n    // swallowed.\n    void accepted.catch(() => {});\n\n    return {\n      signal,\n      accepted,\n    };\n  }\n}\n\nexport const agentThreadStreamRuntime = new AgentThreadStreamRuntime();\n","import type {\n  NotificationDeliveryAction,\n  NotificationDeliveryDecision,\n  NotificationDeliveryThreadState,\n  NotificationPriority,\n  NotificationRecord,\n  NotificationStatus,\n} from './types';\n\n/**\n * Delivery attempts allowed before a notification is marked `failed`. A failed\n * record is no longer due, so a deterministic delivery error (a missing model,\n * a rejected request context) stops being retried on every dispatch tick.\n */\nexport const MAX_NOTIFICATION_DELIVERY_ATTEMPTS = 5;\n\n/**\n * Attempt bookkeeping applied when a delivery attempt throws, spread into the\n * `updateNotification` call at each failure site.\n */\nexport function resolveDeliveryFailureUpdate(record: NotificationRecord): {\n  deliveryAttempts: number;\n  status?: NotificationStatus;\n} {\n  const attempts = record.deliveryAttempts ?? 0;\n  // Records written before the cap existed can already be past it. Terminalize\n  // them at their recorded count instead of inflating it further.\n  if (attempts >= MAX_NOTIFICATION_DELIVERY_ATTEMPTS) return { deliveryAttempts: attempts, status: 'failed' };\n  const deliveryAttempts = attempts + 1;\n  if (deliveryAttempts < MAX_NOTIFICATION_DELIVERY_ATTEMPTS) return { deliveryAttempts };\n  return { deliveryAttempts, status: 'failed' };\n}\n\nexport type NotificationDeliveryPolicyDecision = NotificationDeliveryAction | NotificationDeliveryDecision;\n\nexport type NotificationDeliveryPolicyInput = {\n  record: NotificationRecord;\n  threadState: NotificationDeliveryThreadState;\n  now: Date;\n};\n\nexport type NotificationDeliveryPolicyDecider = (\n  input: NotificationDeliveryPolicyInput,\n) => NotificationDeliveryPolicyDecision | undefined | Promise<NotificationDeliveryPolicyDecision | undefined>;\n\nexport type NotificationDeliveryPolicyConfig = {\n  default?: NotificationDeliveryPolicyDecision;\n  priorities?: Partial<Record<NotificationPriority, NotificationDeliveryPolicyDecision>>;\n  sources?: Record<string, NotificationDeliveryPolicyDecision>;\n  decide?: NotificationDeliveryPolicyDecider;\n};\n\nconst normalizeDecision = (decision: NotificationDeliveryPolicyDecision): NotificationDeliveryDecision => {\n  if (typeof decision === 'string') return { action: decision };\n  return decision;\n};\n\nexport function defaultNotificationDeliveryDecision(\n  input: NotificationDeliveryPolicyInput,\n): NotificationDeliveryDecision {\n  if (input.record.priority === 'urgent') {\n    return { action: 'deliver', reason: 'urgent' };\n  }\n\n  if (input.record.priority === 'high') {\n    return input.threadState === 'active'\n      ? { action: 'summarize', summaryAt: input.now, deliverAt: input.now, reason: 'active-high-summary-then-full' }\n      : { action: 'deliver', reason: 'idle-high' };\n  }\n\n  if (input.record.priority === 'medium') {\n    return input.threadState === 'active'\n      ? { action: 'summarize', summaryAt: input.now, reason: 'active-batch-summary' }\n      : { action: 'deliver', reason: 'idle-medium' };\n  }\n\n  return {\n    action: 'summarize',\n    summaryAt: input.now,\n    reason: input.threadState === 'active' ? 'active-batch-summary' : 'idle-low-summary',\n  };\n}\n\nexport async function resolveNotificationDeliveryDecision({\n  config,\n  ...input\n}: NotificationDeliveryPolicyInput & {\n  config?: NotificationDeliveryPolicyConfig;\n}): Promise<NotificationDeliveryDecision> {\n  const custom = await config?.decide?.(input);\n  if (custom) return normalizeDecision(custom);\n\n  const sourceDecision = config?.sources?.[input.record.source];\n  if (sourceDecision) return normalizeDecision(sourceDecision);\n\n  const priorityDecision = config?.priorities?.[input.record.priority];\n  if (priorityDecision) return normalizeDecision(priorityDecision);\n\n  if (config?.default) return normalizeDecision(config.default);\n\n  return defaultNotificationDeliveryDecision(input);\n}\n","import { createSignal } from '../agent/signals';\nimport type { AgentSignalAttributes, CreatedAgentSignal } from '../agent/signals';\nimport type {\n  NotificationPriority,\n  NotificationRecord,\n  NotificationSignalMetadata,\n  NotificationSummary,\n  NotificationSummarySignalMetadata,\n} from './types';\n\nexport function notificationSignalAttributes(notification: NotificationRecord): AgentSignalAttributes {\n  return {\n    ...notification.attributes,\n    id: notification.id,\n    source: notification.source,\n    type: notification.kind,\n    kind: notification.kind,\n    priority: notification.priority,\n    status: notification.status,\n    ...(notification.coalescedCount && notification.coalescedCount > 1\n      ? { coalescedCount: notification.coalescedCount }\n      : {}),\n  };\n}\n\nexport function notificationSummaryContents(summary: NotificationSummary): string {\n  const sources = Object.entries(summary.bySource)\n    .sort(([a], [b]) => a.localeCompare(b))\n    .map(([source, count]) => `${source}: ${count}`)\n    .join(', ');\n  return sources || 'No pending notifications';\n}\n\nconst priorityOrder: NotificationPriority[] = ['low', 'medium', 'high', 'urgent'];\n\nfunction highestPriority(byPriority: Partial<Record<NotificationPriority, number>>): NotificationPriority | undefined {\n  for (let i = priorityOrder.length - 1; i >= 0; i -= 1) {\n    const priority = priorityOrder[i];\n    if (priority && (byPriority[priority] ?? 0) > 0) return priority;\n  }\n  return undefined;\n}\n\nexport function notificationSignalMetadata(notification: NotificationRecord): NotificationSignalMetadata {\n  return {\n    signal: 'notification',\n    recordId: notification.id,\n    source: notification.source,\n    kind: notification.kind,\n    priority: notification.priority,\n    status: notification.status,\n    ...(notification.coalescedCount && notification.coalescedCount > 1\n      ? { coalescedCount: notification.coalescedCount }\n      : {}),\n    ...(notification.deliveredAt ? { deliveredAt: notification.deliveredAt.toISOString() } : {}),\n    ...(notification.seenAt ? { seenAt: notification.seenAt.toISOString() } : {}),\n  };\n}\n\nexport function notificationSummarySignalMetadata(summary: NotificationSummary): NotificationSummarySignalMetadata {\n  const priority = highestPriority(summary.byPriority);\n  return {\n    signal: 'summary',\n    pending: summary.pending,\n    groups: Object.entries(summary.bySource)\n      .sort(([a], [b]) => a.localeCompare(b))\n      .map(([source, count]) => ({ source, count })),\n    byPriority: summary.byPriority,\n    notificationIds: summary.notificationIds,\n    ...(priority ? { priority } : {}),\n  };\n}\n\nexport function createNotificationSignal(notification: NotificationRecord): CreatedAgentSignal {\n  return createSignal({\n    type: 'notification',\n    tagName: 'notification',\n    contents: notification.summary,\n    attributes: notificationSignalAttributes(notification),\n    metadata: { ...notification.metadata, notification: notificationSignalMetadata(notification) },\n  });\n}\n\nexport function createNotificationSummarySignal(summary: NotificationSummary): CreatedAgentSignal {\n  const notification = notificationSummarySignalMetadata(summary);\n  return createSignal({\n    type: 'notification',\n    tagName: 'notification-summary',\n    contents: notificationSummaryContents(summary),\n    attributes: {\n      pending: summary.pending,\n      ...(notification.priority ? { priority: notification.priority } : {}),\n    },\n    metadata: { notification, notificationSummary: summary, notificationIds: summary.notificationIds },\n  });\n}\n\nexport function summarizeNotifications(notifications: NotificationRecord[]): NotificationSummary {\n  const pendingNotifications = notifications.filter(notification => notification.status === 'pending');\n  const first = pendingNotifications[0] ?? notifications[0];\n  return pendingNotifications.reduce<NotificationSummary>(\n    (summary, notification) => {\n      summary.pending += 1;\n      summary.bySource[notification.source] = (summary.bySource[notification.source] ?? 0) + 1;\n      summary.byPriority[notification.priority] = (summary.byPriority[notification.priority] ?? 0) + 1;\n      summary.notificationIds.push(notification.id);\n      return summary;\n    },\n    {\n      threadId: first?.threadId ?? '',\n      resourceId: first?.resourceId,\n      agentId: first?.agentId,\n      pending: 0,\n      bySource: {},\n      byPriority: {},\n      notificationIds: [],\n    },\n  );\n}\n","import type { CreatedAgentSignal } from '../agent/signals';\nimport { agentThreadStreamRuntime } from '../agent/thread-stream-runtime';\nimport type { SendAgentSignalOptions, SendAgentSignalResult } from '../agent/types';\nimport type { PubSub } from '../events';\nimport type { Mastra } from '../mastra';\nimport { resolveDeliveryFailureUpdate } from './delivery-policy';\nimport { createNotificationSignal, createNotificationSummarySignal, summarizeNotifications } from './signals';\nimport type { NotificationsStorage } from './storage';\nimport type { NotificationDeliveryThreadState, NotificationRecord } from './types';\n\ntype NotificationDispatchAgent = {\n  id?: string;\n  getPubSub?: () => PubSub | undefined;\n  sendSignal: (signal: CreatedAgentSignal, target: SendAgentSignalOptions) => SendAgentSignalResult;\n};\n\nexport type DispatchDueNotificationsInput = {\n  mastra: Mastra;\n  storage: NotificationsStorage;\n  now?: Date;\n  limit?: number;\n};\n\nexport type DispatchDueNotificationsResult = {\n  delivered: NotificationRecord[];\n  failed: Array<{ record: NotificationRecord; error: string }>;\n  signals: CreatedAgentSignal[];\n};\n\nconst errorMessage = (error: unknown): string => (error instanceof Error ? error.message : String(error));\n\nconst isSummaryDue = (record: NotificationRecord, now: Date): boolean =>\n  Boolean(record.summaryAt && record.summaryAt.getTime() <= now.getTime());\n\nconst deliveryPriority: Record<NotificationRecord['priority'], number> = {\n  urgent: 0,\n  high: 1,\n  medium: 2,\n  low: 3,\n};\n\ntype DueDispatchGroup = {\n  key: string;\n  agentId: string;\n  resourceId: string;\n  threadId: string;\n  summaryRecords: NotificationRecord[];\n  individualRecords: NotificationRecord[];\n};\n\ntype DueDispatchItem =\n  | { type: 'summary'; records: NotificationRecord[]; priority: NotificationRecord['priority']; createdAt: Date }\n  | { type: 'individual'; record: NotificationRecord; priority: NotificationRecord['priority']; createdAt: Date };\n\nconst compareDueDispatchItems = (a: DueDispatchItem, b: DueDispatchItem): number =>\n  deliveryPriority[a.priority] - deliveryPriority[b.priority] || a.createdAt.getTime() - b.createdAt.getTime();\n\nconst getHighestPriority = (records: NotificationRecord[]): NotificationRecord['priority'] =>\n  records.reduce<NotificationRecord['priority']>(\n    (highest, record) => (deliveryPriority[record.priority] < deliveryPriority[highest] ? record.priority : highest),\n    'low',\n  );\n\nconst getEarliestCreatedAt = (records: NotificationRecord[]): Date =>\n  records.reduce(\n    (earliest, record) => (record.createdAt.getTime() < earliest.getTime() ? record.createdAt : earliest),\n    records[0]!.createdAt,\n  );\n\nconst groupKey = (record: NotificationRecord): string | undefined => {\n  if (!record.agentId || !record.resourceId || !record.threadId) return undefined;\n  return [record.agentId, record.resourceId, record.threadId].join('\\0');\n};\n\nasync function recordDeliveryFailure({\n  storage,\n  record,\n  now,\n  error,\n}: {\n  storage: NotificationsStorage;\n  record: NotificationRecord;\n  now: Date;\n  error: unknown;\n}) {\n  await storage.updateNotification({\n    id: record.id,\n    threadId: record.threadId,\n    ...resolveDeliveryFailureUpdate(record),\n    lastDeliveryAttemptAt: now,\n    lastDeliveryError: errorMessage(error),\n  });\n}\n\nasync function sendNotificationRecord({\n  mastra,\n  storage,\n  record,\n  now,\n  batchThreadState,\n}: {\n  mastra: Mastra;\n  storage: NotificationsStorage;\n  record: NotificationRecord;\n  now: Date;\n  batchThreadState?: NotificationDeliveryThreadState;\n}): Promise<{ record: NotificationRecord; signal: CreatedAgentSignal } | null> {\n  const current = await storage.getNotification({ threadId: record.threadId, id: record.id });\n  if (!current || current.status !== 'pending' || current.deliveredSignalId) return null;\n  if (!current.agentId) throw new Error(`Notification ${current.id} is missing agentId`);\n  if (!current.resourceId) throw new Error(`Notification ${current.id} is missing resourceId`);\n\n  const agent = (await mastra.getAgentById(current.agentId as never)) as NotificationDispatchAgent;\n  if (current.priority === 'high' && current.summarySignalId) {\n    const threadState =\n      batchThreadState ??\n      agentThreadStreamRuntime.getThreadState(\n        { resourceId: current.resourceId, threadId: current.threadId },\n        agent.getPubSub?.(),\n      );\n    if (threadState === 'active') return null;\n  }\n\n  const signal = createNotificationSignal({\n    ...current,\n    status: 'delivered',\n    deliveredAt: now,\n    lastDeliveryAttemptAt: now,\n  });\n  const target: SendAgentSignalOptions = { resourceId: current.resourceId, threadId: current.threadId };\n  const result = agent.sendSignal(signal, target);\n  // `accepted` rejects when the signal could not be routed/started (e.g. a\n  // misconfigured agent). Let that propagate so the caller records the\n  // notification as a failed delivery.\n  await result.accepted;\n  await result.persisted;\n  const updated = await storage.updateNotification({\n    id: current.id,\n    threadId: current.threadId,\n    status: 'delivered',\n    deliveredSignalId: result.signal.id,\n    lastDeliveryAttemptAt: now,\n  });\n  return { record: updated, signal: result.signal };\n}\n\nasync function sendNotificationSummary({\n  mastra,\n  storage,\n  records,\n  now,\n}: {\n  mastra: Mastra;\n  storage: NotificationsStorage;\n  records: NotificationRecord[];\n  now: Date;\n}): Promise<{ records: NotificationRecord[]; signal: CreatedAgentSignal }> {\n  const first = records[0];\n  if (!first?.agentId) throw new Error('Notification summary is missing agentId');\n  if (!first.resourceId) throw new Error('Notification summary is missing resourceId');\n\n  const agent = await mastra.getAgentById(first.agentId as never);\n  const summary = summarizeNotifications(records);\n  const signal = createNotificationSummarySignal(summary);\n  const target: SendAgentSignalOptions = records.every(record => record.priority === 'low')\n    ? { resourceId: first.resourceId, threadId: first.threadId, ifIdle: { behavior: 'persist' } }\n    : { resourceId: first.resourceId, threadId: first.threadId };\n  const result = (agent as NotificationDispatchAgent).sendSignal(signal, target);\n  // `accepted` rejects when the signal could not be routed/started; let it\n  // propagate so the caller records the notifications as failed deliveries.\n  await result.accepted;\n  await result.persisted;\n\n  const updatedRecords: NotificationRecord[] = [];\n  for (const record of records) {\n    updatedRecords.push(\n      await storage.updateNotification({\n        id: record.id,\n        threadId: record.threadId,\n        summaryAt: null,\n        summarySignalId: result.signal.id,\n        lastDeliveryAttemptAt: now,\n      }),\n    );\n  }\n  return { records: updatedRecords, signal: result.signal };\n}\n\nasync function getBatchThreadState({\n  mastra,\n  group,\n}: {\n  mastra: Mastra;\n  group: DueDispatchGroup;\n}): Promise<NotificationDeliveryThreadState> {\n  const agent = (await mastra.getAgentById(group.agentId as never)) as NotificationDispatchAgent;\n  return agentThreadStreamRuntime.getThreadState(\n    { resourceId: group.resourceId, threadId: group.threadId },\n    agent.getPubSub?.(),\n  );\n}\n\nexport async function dispatchDueNotifications({\n  mastra,\n  storage,\n  now = new Date(),\n  limit = 100,\n}: DispatchDueNotificationsInput): Promise<DispatchDueNotificationsResult> {\n  const due = await storage.listDueNotifications({ now, limit });\n  const delivered: NotificationRecord[] = [];\n  const failed: Array<{ record: NotificationRecord; error: string }> = [];\n  const signals: CreatedAgentSignal[] = [];\n  const groups = new Map<string, DueDispatchGroup>();\n  const ungroupedIndividual: NotificationRecord[] = [];\n\n  for (const record of due) {\n    const key = groupKey(record);\n    if (!key) {\n      if (isSummaryDue(record, now)) {\n        const error = new Error(\n          `Notification ${record.id} cannot be summarized without agentId, resourceId, and threadId`,\n        );\n        await recordDeliveryFailure({ storage, record, now, error });\n        failed.push({ record, error: error.message });\n      } else {\n        ungroupedIndividual.push(record);\n      }\n      continue;\n    }\n\n    const group = groups.get(key) ?? {\n      key,\n      agentId: record.agentId!,\n      resourceId: record.resourceId!,\n      threadId: record.threadId,\n      summaryRecords: [],\n      individualRecords: [],\n    };\n    if (isSummaryDue(record, now)) {\n      group.summaryRecords.push(record);\n    } else {\n      group.individualRecords.push(record);\n    }\n    groups.set(key, group);\n  }\n\n  for (const group of groups.values()) {\n    const records = [...group.summaryRecords, ...group.individualRecords];\n    let batchThreadState: NotificationDeliveryThreadState;\n    try {\n      batchThreadState = await getBatchThreadState({ mastra, group });\n    } catch (error) {\n      for (const record of records) {\n        await recordDeliveryFailure({ storage, record, now, error });\n        failed.push({ record, error: errorMessage(error) });\n      }\n      continue;\n    }\n\n    const items: DueDispatchItem[] = group.individualRecords.map(record => ({\n      type: 'individual',\n      record,\n      priority: record.priority,\n      createdAt: record.createdAt,\n    }));\n    if (group.summaryRecords.length > 0) {\n      items.push({\n        type: 'summary',\n        records: group.summaryRecords,\n        priority: getHighestPriority(group.summaryRecords),\n        createdAt: getEarliestCreatedAt(group.summaryRecords),\n      });\n    }\n    items.sort(compareDueDispatchItems);\n\n    for (const item of items) {\n      if (item.type === 'summary') {\n        try {\n          const result = await sendNotificationSummary({ mastra, storage, records: item.records, now });\n          delivered.push(...result.records);\n          signals.push(result.signal);\n        } catch (error) {\n          for (const record of item.records) {\n            await recordDeliveryFailure({ storage, record, now, error });\n            failed.push({ record, error: errorMessage(error) });\n          }\n        }\n        continue;\n      }\n\n      try {\n        const result = await sendNotificationRecord({ mastra, storage, record: item.record, now, batchThreadState });\n        if (!result) continue;\n        delivered.push(result.record);\n        signals.push(result.signal);\n      } catch (error) {\n        await recordDeliveryFailure({ storage, record: item.record, now, error });\n        failed.push({ record: item.record, error: errorMessage(error) });\n      }\n    }\n  }\n\n  ungroupedIndividual.sort((a, b) =>\n    compareDueDispatchItems(\n      { type: 'individual', record: a, priority: a.priority, createdAt: a.createdAt },\n      { type: 'individual', record: b, priority: b.priority, createdAt: b.createdAt },\n    ),\n  );\n  for (const record of ungroupedIndividual) {\n    try {\n      const result = await sendNotificationRecord({ mastra, storage, record, now });\n      if (!result) continue;\n      delivered.push(result.record);\n      signals.push(result.signal);\n    } catch (error) {\n      await recordDeliveryFailure({ storage, record, now, error });\n      failed.push({ record, error: errorMessage(error) });\n    }\n  }\n\n  return { delivered, failed, signals };\n}\n","import { randomUUID } from 'node:crypto';\nimport { StorageDomain } from '../storage/domains/base';\nimport type {\n  CreateNotificationInput,\n  ListDueNotificationsInput,\n  ListNotificationsInput,\n  NotificationRecord,\n  NotificationStatus,\n  UpdateNotificationInput,\n} from './types';\n\nexport abstract class NotificationsStorage extends StorageDomain {\n  constructor() {\n    super({ component: 'STORAGE', name: 'NOTIFICATIONS' });\n  }\n\n  abstract createNotification(input: CreateNotificationInput): Promise<NotificationRecord>;\n  abstract listNotifications(input: ListNotificationsInput): Promise<NotificationRecord[]>;\n  abstract listDueNotifications(input: ListDueNotificationsInput): Promise<NotificationRecord[]>;\n  abstract getNotification(input: { threadId: string; id: string }): Promise<NotificationRecord | null>;\n  abstract updateNotification(input: UpdateNotificationInput): Promise<NotificationRecord>;\n}\n\nconst cloneDate = (value?: Date) => (value ? new Date(value) : undefined);\nconst notificationKey = (threadId: string, id: string) => `${threadId}\\0${id}`;\nconst cloneValue = <T>(value: T | undefined): T | undefined =>\n  value === undefined ? undefined : structuredClone(value);\n\nconst cloneRecord = (record: NotificationRecord): NotificationRecord => ({\n  ...record,\n  createdAt: new Date(record.createdAt),\n  updatedAt: new Date(record.updatedAt),\n  deliveredAt: cloneDate(record.deliveredAt),\n  seenAt: cloneDate(record.seenAt),\n  dismissedAt: cloneDate(record.dismissedAt),\n  archivedAt: cloneDate(record.archivedAt),\n  discardedAt: cloneDate(record.discardedAt),\n  deliverAt: cloneDate(record.deliverAt),\n  summaryAt: cloneDate(record.summaryAt),\n  lastDeliveryAttemptAt: cloneDate(record.lastDeliveryAttemptAt),\n  payload: cloneValue(record.payload),\n  attributes: cloneValue(record.attributes),\n  metadata: cloneValue(record.metadata),\n});\n\nconst statusTimestamp = (status: NotificationStatus, now: Date) => {\n  if (status === 'delivered') return { deliveredAt: now };\n  if (status === 'seen') return { seenAt: now };\n  if (status === 'dismissed') return { dismissedAt: now };\n  if (status === 'archived') return { archivedAt: now };\n  if (status === 'discarded') return { discardedAt: now };\n  return {};\n};\n\nconst valueMatches = <T extends string>(value: T, filter?: T | T[]) => {\n  if (!filter) return true;\n  return Array.isArray(filter) ? filter.includes(value) : value === filter;\n};\n\nconst dueTime = (record: NotificationRecord): number => {\n  const deliverAt = record.deliverAt?.getTime();\n  const summaryAt = record.summaryAt?.getTime();\n  if (deliverAt !== undefined && summaryAt !== undefined) return Math.min(deliverAt, summaryAt);\n  return deliverAt ?? summaryAt ?? Number.POSITIVE_INFINITY;\n};\n\nexport class InMemoryNotificationsStorage extends NotificationsStorage {\n  #notifications = new Map<string, NotificationRecord>();\n\n  async createNotification(input: CreateNotificationInput): Promise<NotificationRecord> {\n    const existing = this.findCoalescable(input);\n    if (existing) {\n      const now = new Date();\n      const next: NotificationRecord = {\n        ...existing,\n        summary: input.summary,\n        payload: cloneValue(input.payload ?? existing.payload),\n        priority: input.priority ?? existing.priority,\n        attributes: input.attributes\n          ? { ...cloneValue(existing.attributes), ...cloneValue(input.attributes) }\n          : cloneValue(existing.attributes),\n        updatedAt: now,\n        deliverAt: input.deliverAt ?? existing.deliverAt,\n        summaryAt: input.summaryAt ?? existing.summaryAt,\n        deliveryReason: input.deliveryReason ?? existing.deliveryReason,\n        coalescedCount: (existing.coalescedCount ?? 1) + 1,\n        metadata: input.metadata\n          ? { ...cloneValue(existing.metadata), ...cloneValue(input.metadata) }\n          : cloneValue(existing.metadata),\n      };\n      this.#notifications.set(notificationKey(next.threadId, next.id), next);\n      return cloneRecord(next);\n    }\n\n    const now = input.createdAt ?? new Date();\n    const record: NotificationRecord = {\n      id: input.id ?? randomUUID(),\n      threadId: input.threadId,\n      source: input.source,\n      kind: input.kind,\n      priority: input.priority ?? 'medium',\n      status: 'pending',\n      summary: input.summary,\n      payload: cloneValue(input.payload),\n      resourceId: input.resourceId,\n      agentId: input.agentId,\n      sourceId: input.sourceId,\n      dedupeKey: input.dedupeKey,\n      coalesceKey: input.coalesceKey,\n      coalescedCount: 1,\n      attributes: cloneValue(input.attributes),\n      createdAt: now,\n      updatedAt: now,\n      deliverAt: input.deliverAt,\n      summaryAt: input.summaryAt,\n      deliveryReason: input.deliveryReason,\n      deliveryAttempts: 0,\n      metadata: cloneValue(input.metadata),\n    };\n    this.#notifications.set(notificationKey(record.threadId, record.id), record);\n    return cloneRecord(record);\n  }\n\n  async listNotifications(input: ListNotificationsInput): Promise<NotificationRecord[]> {\n    const search = input.search?.toLowerCase();\n    const results = [...this.#notifications.values()]\n      .filter(record => record.threadId === input.threadId)\n      .filter(record => valueMatches(record.status, input.status))\n      .filter(record => valueMatches(record.priority, input.priority))\n      .filter(record => !input.source || record.source === input.source)\n      .filter(record => !input.resourceId || record.resourceId === input.resourceId)\n      .filter(record => !input.agentId || record.agentId === input.agentId)\n      .filter(\n        record =>\n          !search ||\n          record.summary.toLowerCase().includes(search) ||\n          record.kind.toLowerCase().includes(search) ||\n          record.source.toLowerCase().includes(search),\n      )\n      .sort((a, b) => b.updatedAt.getTime() - a.updatedAt.getTime());\n    return results.slice(0, input.limit ?? results.length).map(cloneRecord);\n  }\n\n  async listDueNotifications(input: ListDueNotificationsInput): Promise<NotificationRecord[]> {\n    const now = input.now.getTime();\n    const results = [...this.#notifications.values()]\n      .filter(record => record.status === 'pending')\n      .filter(record => !input.agentId || record.agentId === input.agentId)\n      .filter(record => !input.resourceId || record.resourceId === input.resourceId)\n      .filter(record => dueTime(record) <= now)\n      .sort((a, b) => dueTime(a) - dueTime(b) || a.updatedAt.getTime() - b.updatedAt.getTime());\n    return results.slice(0, input.limit ?? results.length).map(cloneRecord);\n  }\n\n  async getNotification(input: { threadId: string; id: string }): Promise<NotificationRecord | null> {\n    const record = this.#notifications.get(notificationKey(input.threadId, input.id));\n    if (!record) return null;\n    return cloneRecord(record);\n  }\n\n  async updateNotification(input: UpdateNotificationInput): Promise<NotificationRecord> {\n    const existing = this.#notifications.get(notificationKey(input.threadId, input.id));\n    if (!existing) {\n      throw new Error(`Notification ${input.id} was not found for thread ${input.threadId}`);\n    }\n    const now = new Date();\n    const next: NotificationRecord = {\n      ...existing,\n      ...(input.status ? { status: input.status, ...statusTimestamp(input.status, now) } : {}),\n      ...(input.summary !== undefined ? { summary: input.summary } : {}),\n      ...(input.payload !== undefined ? { payload: cloneValue(input.payload) } : {}),\n      ...(input.attributes !== undefined ? { attributes: cloneValue(input.attributes) } : {}),\n      ...(input.metadata !== undefined ? { metadata: cloneValue(input.metadata) } : {}),\n      ...(input.deliverAt !== undefined ? { deliverAt: input.deliverAt ?? undefined } : {}),\n      ...(input.summaryAt !== undefined ? { summaryAt: input.summaryAt ?? undefined } : {}),\n      ...(input.deliveryReason !== undefined ? { deliveryReason: input.deliveryReason } : {}),\n      ...(input.deliveryAttempts !== undefined ? { deliveryAttempts: input.deliveryAttempts } : {}),\n      ...(input.lastDeliveryAttemptAt !== undefined ? { lastDeliveryAttemptAt: input.lastDeliveryAttemptAt } : {}),\n      ...(input.lastDeliveryError !== undefined ? { lastDeliveryError: input.lastDeliveryError } : {}),\n      ...(input.deliveredSignalId !== undefined ? { deliveredSignalId: input.deliveredSignalId } : {}),\n      ...(input.summarySignalId !== undefined ? { summarySignalId: input.summarySignalId } : {}),\n      updatedAt: now,\n    };\n    this.#notifications.set(notificationKey(next.threadId, next.id), next);\n    return cloneRecord(next);\n  }\n\n  async dangerouslyClearAll(): Promise<void> {\n    this.#notifications.clear();\n  }\n\n  private findCoalescable(input: CreateNotificationInput): NotificationRecord | undefined {\n    if (!input.dedupeKey && !input.coalesceKey) return undefined;\n    return [...this.#notifications.values()].find(record => {\n      if (\n        record.threadId !== input.threadId ||\n        record.source !== input.source ||\n        record.kind !== input.kind ||\n        record.status !== 'pending'\n      )\n        return false;\n      if (record.agentId !== input.agentId || record.resourceId !== input.resourceId) return false;\n      return Boolean(\n        (input.dedupeKey && record.dedupeKey === input.dedupeKey) ||\n        (input.coalesceKey && record.coalesceKey === input.coalesceKey),\n      );\n    });\n  }\n}\n"],"mappings":";;;;;;;;;AAiCA,MAAM,6BAA6B;AACnC,MAAM,mCAAmC;;;;;;;AAOzC,MAAM,4BAA4BA,cAAAA,mBAAmB,oCAAoC,IAAM;;;;;;AAM/F,MAAM,uCAAuCA,cAAAA,mBAC3C,+CACA,KAAK,MAAM,4BAA4B,CAAC,CAC1C;;;;;;;;;;AAUA,MAAM,6BAA6BA,cAAAA,mBAAmB,+BAA+B,OAAU,GAAI;AAEnG,IAAW,2BAAmC,IAAIC,sBAAAA,mBAAmB;AAErE,SAAS,iBAAiB,QAAiB,YAAoB,UAAkB;CAC/E,OAAO;EACL,GAAK,UAAU,OAAO,WAAW,WAAW,SAAS,CAAC;EACtD,UAAW,QAA8C,YAAY;EACrE,QAAS,QAA4C,UAAU;CACjE;AACF;AA6FA,SAAS,qBAA8C;CACrD,OAAO;EACL,gCAAgB,IAAI,IAAI;EACxB,sCAAsB,IAAI,IAAI;EAC9B,mCAAmB,IAAI,IAAI;EAC3B,yCAAyB,IAAI,IAAI;EACjC,oCAAoB,IAAI,IAAI;EAC5B,uCAAuB,IAAI,IAAI;EAC/B,kCAAkB,IAAI,IAAI;EAC1B,yCAAyB,IAAI,IAAI;EACjC,iCAAiB,IAAI,IAAI;EACzB,2CAA2B,IAAI,IAAI;EACnC,wCAAwB,IAAI,IAAI;EAChC,uCAAuB,IAAI,IAAI;EAC/B,4CAA4B,IAAI,IAAI;EACpC,8CAA8B,IAAI,IAAI;EACtC,wCAAwB,IAAI,IAAI;EAChC,kCAAkB,IAAI,IAAI;EAC1B,+BAAe,IAAI,IAAI;EACvB,oCAAoB,IAAI,IAAI;CAC9B;AACF;AAEA,IAAa,2BAAb,MAAsC;CACpC;CACA,kCAAkB,IAAI,QAAyC;CAE/D,WAAW,QAAyB;EAClC,OAAO,UAAU;CACnB;;;;;;;;;;;;;CAcA,kBAAkB,QAAgC;EAChD,MAAM,WAAW,KAAKC,WAAW,MAAM;EACvC,MAAM,SAAU,SAAoE;EACpF,IAAI,OAAO,WAAW,YAEpB,OADc,OAAO,KAAK,QACf,KAAKC,sBAAAA;EAElB,OAAOC,sBAAAA,gBAAgB,QAAQ,IAAI,WAAWD,sBAAAA;CAChD;CAEA,sBAAsB,QAAmE;EACvF,MAAM,WAAW,KAAKE,kBAAkB,MAAM;EAC9C,OAAO;GAAE;GAAU,YAAY,aAAaF,sBAAAA;EAAkB;CAChE;CAEA,MAAMG,oBAAoB,QAAgB,KAAa,OAAiC;EACtF,MAAM,EAAE,UAAU,eAAe,KAAKC,sBAAsB,MAAM;EAClE,IAAI,YAAY,OAAO;EACvB,OAAO,SACJ,cAAc,GAAG,CAAC,CAClB,MAAK,UAAS,UAAU,KAAK,CAAC,CAC9B,YAAY,KAAK;CACtB;CAEA,eAAuB;EACrB,KAAKC,SAAAA,GAAAA,OAAAA,WAAAA,CAAmB;EACxB,OAAO,KAAKA;CACd;;;;;;;;;CAUA,oBAAoB,QAA4B,KAAa,OAAqB;EAChF,MAAM,WAAW,KAAKN,WAAW,MAAM;EACvC,KAAKO,kBAAkB,UAAU,KAAK;EACtC,KAAUJ,kBAAkB,QAAQ,CAAC,CAClC,aAAa,KAAK,KAAK,CAAC,CACxB,YAAY,CAAC,CAAC;CACnB;;;;;;;;;CAUA,mBAAmB,QAAgB,KAAa,OAAqB;EACnE,MAAM,QAAQ,KAAKK,UAAU,MAAM;EACnC,IAAI,MAAM,mBAAmB,IAAI,KAAK,GAAG;EACzC,MAAM,gBAAgB,KAAKL,kBAAkB,MAAM;EACnD,MAAM,QAAQ,kBAAkB;GAC9B,cACG,WAAW,KAAK,OAAO,yBAAyB,CAAC,CACjD,MAAK,YAAW;IACf,IAAI,CAAC,SAGH,KAAKI,kBAAkB,QAAQ,KAAK;GAExC,CAAC,CAAC,CACD,YAAY,CAAC,CAAC;EACnB,GAAG,oCAAoC;EAEvC,IAAI,OAAO,UAAU,YAAY,SAAS,OAAQ,MAAc,UAAU,YACxE,MAAe,MAAM;EAEvB,MAAM,mBAAmB,IAAI,OAAO,KAAK;CAC3C;CAEA,kBAAkB,QAAgB,OAAqB;EACrD,MAAM,QAAQ,KAAKC,UAAU,MAAM;EACnC,MAAM,QAAQ,MAAM,mBAAmB,IAAI,KAAK;EAChD,IAAI,CAAC,OAAO;EACZ,cAAc,KAAK;EACnB,MAAM,mBAAmB,OAAO,KAAK;CACvC;;;;;;;;;;;;;CAcA,MAAMC,qBACJ,QACA,KACA,WACA,SACkB;EAClB,MAAM,WAAW,KAAKT,WAAW,MAAM;EAKvC,MAAM,OAAO,MAJS,KAAKG,kBAAkB,QAId,CAAC,CAC7B,cAAc,KAAK,WAAW,SAAS,yBAAyB,CAAC,CACjE,YAAY,KAAK;EAIpB,KAAKI,kBAAkB,UAAU,SAAS;EAC1C,IAAI,MACF,KAAKG,mBAAmB,UAAU,KAAK,OAAO;EAEhD,OAAO;CACT;;;;;;;;;;;;;;;CAgBA,MAAMC,8BACJ,QACA,KACA,SACA,WACgD;EAChD,MAAM,WAAW,KAAKX,WAAW,MAAM;EACvC,IAAI,WAEE;OAAA,MADsB,KAAKS,qBAAqB,QAAQ,KAAK,WAAW,OAAO,GAClE,OAAO;IAAE,UAAU;IAAM,OAAO;GAAQ;EAAA;EAI3D,MAAM,SAAS,MADO,KAAKN,kBAAkB,QACZ,CAAC,CAC/B,aAAa,KAAK,SAAS,yBAAyB,CAAC,CACrD,aAAa;GAAE,UAAU;GAAkB,OAAO,KAAA;EAAgC,EAAE;EACvF,IAAI,OAAO,UAAU;GACnB,KAAKO,mBAAmB,UAAU,KAAK,OAAO;GAC9C,OAAO;IAAE,UAAU;IAAM,OAAO;GAAQ;EAC1C;EACA,OAAO;GAAE,UAAU;GAAO,OAAO,OAAO;EAAM;CAChD;;;;;;;CAQA,sBAAsB,OAAgC,KAAsB;EAC1E,QACG,MAAM,uBAAuB,IAAI,GAAG,CAAC,EAAE,UAAU,KAAK,MACtD,MAAM,sBAAsB,IAAI,GAAG,CAAC,EAAE,UAAU,KAAK,MACrD,MAAM,6BAA6B,IAAI,GAAG,CAAC,EAAE,UAAU,KAAK,MAC5D,MAAM,2BAA2B,IAAI,GAAG,CAAC,EAAE,UAAU,KAAK;CAE/D;CAEA,UAAU,QAA0C;EAClD,MAAM,iBAAiB,KAAKV,WAAW,MAAM;EAC7C,IAAI,QAAQ,KAAKY,gBAAgB,IAAI,cAAc;EACnD,IAAI,CAAC,OAAO;GACV,QAAQ,mBAAmB;GAC3B,KAAKA,gBAAgB,IAAI,gBAAgB,KAAK;EAChD;EACA,OAAO;CACT;CAEA,WAAW,YAAgC,UAA0B;EACnE,OAAO,CAAC,cAAc,IAAI,QAAQ,CAAC,CAAC,KAAK,0BAA0B;CACrE;CAEA,aAAa,KAAqB;EAChC,OAAO,GAAG,iCAAiC,GAAG,mBAAmB,GAAG;CACtE;CAEA,wBAAwB,OAAgC,OAAe;EACrE,OAAO,MAAM,wBAAwB,IAAI,KAAK;CAChD;CAEA,gBAAgB,OAAgC,OAAe;EAC7D,OAAO,MAAM,gBAAgB,IAAI,KAAK,KAAK,KAAKC,wBAAwB,OAAO,KAAK;CACtF;CAEA,qBAAqB,OAAgC,QAAmC;EACtF,OACE,OAAO,OAAO,WAAW,aACzB,OAAO,OAAO,WAAW,eACzB,OAAO,cAAc,gBACrB,OAAO,cAAc,eACrB,CAAC,CAAC,OAAO,cACT,KAAKC,gBAAgB,OAAO,OAAO,KAAK;CAE5C;CAEA,iBAAiB,QAAqD;EACpE,OAAO;CACT;CAEA,oBAAoB,OAAgC,OAAe;EACjE,MAAM,aAAa,MAAM,iBAAiB,IAAI,KAAK,KAAK,KAAK;EAC7D,MAAM,iBAAiB,IAAI,OAAO,SAAS;EAC3C,OAAO;GAAE,WAAA,GAAA,OAAA,WAAA,CAAqB;GAAG;EAAU;CAC7C;CAEA,mBACE,OACA,OACA,UACA,YACA;EACA,MAAM,gBAAgB,IAAI,KAAK;EAC/B,MAAM,0BAA0B,IAAI,OAAO,UAAU;EACrD,MAAM,SAAS,MAAM,qBAAqB,IAAI,QAAQ,KAAK,MAAM,eAAe,IAAI,KAAK;EACzF,IAAI,QAAQ;GACV,OAAO,YAAY;GACnB,OAAO,aAAa;EACtB;EACA,IAAI,WAAW,SAAS,YACtB,MAAM,wBAAwB,IAAI,KAAK;CAE3C;CAEA,mBAAmB,OAAgC,OAAe;EAChE,MAAM,gBAAgB,OAAO,KAAK;EAClC,MAAM,0BAA0B,OAAO,KAAK;EAC5C,MAAM,wBAAwB,OAAO,KAAK;CAC5C;CAEA,yBACE,OACA,QACQ;EACR,OACE,MAAM,oBAAoB,CAAC,EAAE,WAAW;GACtC,QAAQ;GACR,QAAQ;GACR,UAAU,MAAM;GAChB,UAAU,OAAO;GACjB,YAAY,OAAO;EACrB,CAAC,MAAA,GAAA,OAAA,WAAA,CAAgB;CAErB;CAEA,0BAA0B,SAAyC;EAEjE,OAAO;GACL,GAFwB,OAAO,YAAY,YAAY,MAAM,QAAQ,OAAO,IAAI,EAAE,UAAU,QAAQ,IAAI;GAGxG,MAAM;GACN,SAAS;EACX;CACF;CAEA,eAAe,SAAoD,QAAmC;EACpG,MAAM,QAAQ,KAAKN,UAAU,MAAM;EACnC,MAAM,MAAM,KAAKO,WAAW,QAAQ,YAAY,QAAQ,QAAQ;EAChE,MAAM,cAAc,MAAM,mBAAmB,IAAI,GAAG;EACpD,IAAI,CAAC,aAAa,OAAO;EAEzB,MAAM,eAAe,MAAM,eAAe,IAAI,WAAW;EACzD,IAAI,gBAAgB,CAAC,KAAKC,qBAAqB,OAAO,YAAY,GAAG;GACnE,MAAM,mBAAmB,OAAO,GAAG;GACnC,OAAO;EACT;EAEA,OAAO;CACT;CAEA,SAAS,QAA4B,KAAa,OAAsC;EACtF,KAAUC,gBAAgB,QAAQ,KAAK,KAAK,CAAC,CAAC,YAAY,CAAC,CAAC;CAC9D;CAEA,MAAMA,gBAAgB,QAA4B,KAAa,OAAsC;EACnG,MAAM,KAAKjB,WAAW,MAAM,CAAC,CAAC,QAAQ,KAAKkB,aAAa,GAAG,GAAG;GAC5D,MAAM,MAAM;GACZ,OAAO,MAAM;GACb,MAAM;EACR,CAAC;CACH;CAEA,qBACE,QACA,QACA,KACA,UACA;EACA,MAAM,UAAU;EAEhB,MAAM,QAAmB,CAAC;EAC1B,MAAM,0BAAU,IAAI,IAAgB;EACpC,IAAI,UAAU;EACd,IAAI,OAAO;EACX,IAAI;EAEJ,MAAM,aAAa;GACjB,MAAM,UAAU,CAAC,GAAG,OAAO;GAC3B,QAAQ,MAAM;GACd,KAAK,MAAM,UAAU,SAAS,OAAO;EACvC;EAEA,MAAM,WAAW,OAAO,SAAkB;GACxC,IAAI,QAAQ,OAAO,SAAS,YAAY,UAAU,MAAM;IACtD,MAAM,YAAY;IAClB,IAAI,UAAU,SAAS,wBAAwB,UAAU,SAAS,uBAChE,QAAQC,mBAAmB,QAAQX,UAAU,MAAM,GAAG,OAAO,OAAO,UAAU;KAC5E,YAAY,UAAU,SAAS;KAC/B,UAAU,UAAU,SAAS;KAC7B,MAAM,UAAU,SAAS,uBAAuB,aAAa;IAC/D,CAAC;GAEL;GACA,MAAM,KAAK,IAAI;GACf,MAAM,QAAQS,gBAAgB,QAAQ,KAAK;IACzC,MAAM;IACN,OAAO,OAAO;IACd;IACA;IACA,UAAU,QAAQG,aAAa;GACjC,CAAC;GACD,KAAK;EACP;EAEA,MAAM,cAAc;GAClB,IAAI,SAAS;GACb,UAAU;GACV,CAAM,YAAY;IAChB,IAAI;KACF,MAAM,SAAS,OAAO;KACtB,IAAI,CAAC,QAAQ;KAEb,IAAI,OAAO,OAAO,cAAc,YAAY;MAC1C,MAAM,SAAS,OAAO,UAAU;MAChC,IAAI;OACF,OAAO,MAAM;QACX,MAAM,EAAE,OAAO,MAAM,MAAM,eAAe,MAAM,OAAO,KAAK;QAC5D,IAAI,YAAY;QAChB,MAAM,SAAS,IAAI;OACrB;MACF,UAAU;OACR,OAAO,YAAY;MACrB;KACF,OACE,WAAW,MAAM,QAAQ,QACvB,MAAM,SAAS,IAAI;IAGzB,SAAS,QAAQ;KACf,QAAQ;IACV,UAAU;KACR,OAAO;KACP,KAAK;IACP;GACF,EAAA,CAAG;EACL;EAEA,MAAM,qBAAqB;GACzB,IAAI,QAAQ;GACZ,IAAI,SAAS;GACb,IAAI;GACJ,OAAO,IAAI,eAAe;IACxB,MAAM,KAAK,YAAY;KACrB,MAAM;KACN,OAAO,CAAC,QAAQ;MACd,IAAI,QAAQ,MAAM,QAAQ;OACxB,WAAW,QAAQ,MAAM,QAAQ;OACjC;MACF;MACA,IAAI,OAAO;OACT,WAAW,MAAM,KAAK;OACtB;MACF;MACA,IAAI,MAAM;OACR,WAAW,MAAM;OACjB;MACF;MACA,MAAM,IAAI,SAAc,YAAW;OACjC,SAAS;OACT,QAAQ,IAAI,OAAO;MACrB,CAAC;MACD,IAAI,QAAQ;OACV,QAAQ,OAAO,MAAM;OACrB,SAAS,KAAA;MACX;KACF;IACF;IACA,SAAS;KACP,SAAS;KACT,IAAI,QAAQ;MACV,QAAQ,OAAO,MAAM;MACrB,OAAO;MACP,SAAS,KAAA;KACX;IACF;GACF,CAAC;EACH;EAEA,OAAO;GAAE;GAAQ,wBAAwB;GAAc,gBAAgB;EAAM;CAC/E;CAEA,iBAAiB,SAA8F;EAC7G,MAAM,SAAS,SAAS,QAAQ;EAOhC,OAAO;GAAE,UALN,SAAS,gBAAgB,IAAA,kBAAwB,MACjD,OAAO,WAAW,WAAW,SAAS,QAAQ;GAI9B,YAFhB,SAAS,gBAAgB,IAAA,oBAA0B,KAA4B,SAAS,QAAQ;EAErE;CAChC;CAEA,kBAA0B,SAAwC,QAAgD;EAChH,MAAM,EAAE,aAAa,KAAKC,iBAAiB,OAAO;EAClD,IAAI,CAAC,YAAY,CAAC,QAAQ,OAAO,OAAO;EAExC,MAAM,QAAQ,KAAKb,UAAU,MAAM;EACnC,MAAM,kBAAkB,IAAI,gBAAgB;EAC5C,MAAM,sBAAsB,QAAQ;EACpC,MAAM,cAAc,gBAAgB,MAAM;EAC1C,IAAI,qBAAqB,SACvB,MAAM;OAEN,qBAAqB,iBAAiB,SAAS,OAAO,EAAE,MAAM,KAAK,CAAC;EAGtE,MAAM,iBAAiB,IAAI,QAAQ,OAAO;GACxC;GACA,eAAe,qBAAqB,oBAAoB,SAAS,KAAK;EACxE,CAAC;EAED,IAAI,MAAM,cAAc,IAAI,QAAQ,KAAK,GACvC,MAAM;EAGR,OAAO;GACL,GAAG;GACH,aAAa,gBAAgB;EAC/B;CACF;CAEA,SAAS,OAAe,QAA0B;EAChD,MAAM,QAAQ,KAAKA,UAAU,MAAM;EACnC,MAAM,cAAc,MAAM,iBAAiB,IAAI,KAAK;EACpD,IAAI,CAAC,aAAa;GAChB,MAAM,cAAc,IAAI,KAAK;GAC7B,OAAO;EACT;EAEA,YAAY,gBAAgB,MAAM;EAClC,MAAM,cAAc,IAAI,KAAK;EAE7B,MAAM,MAAM,MAAM,kBAAkB,IAAI,KAAK;EAC7C,IAAI,KAAK;GACP,MAAM,WAAW,MAAM,mBAAmB,IAAI,GAAG,MAAM,QAAQ,MAAM,sBAAsB,IAAI,GAAG,IAAI,KAAA;GACtG,KAAKc,oBAAoB,QAAQ,KAAK,KAAK;GAC3C,KAAKC,SAAS,QAAQ,KAAK;IAAE,MAAM;IAAe;IAAO;GAAS,CAAC;EACrE;EAEA,OAAO;CACT;CAEA,qBAAqB,SAAwC,QAAqC;EAChG,MAAM,QAAQ,KAAKf,UAAU,MAAM;EACnC,MAAM,MAAM,KAAKO,WAAW,QAAQ,YAAY,QAAQ,QAAQ;EAChE,MAAM,cAAc,MAAM,mBAAmB,IAAI,GAAG;EACpD,IAAI,CAAC,aAAa,OAAO,KAAA;EAEzB,MAAM,SAAS,MAAM,eAAe,IAAI,WAAW;EACnD,IAAI,UAAU,CAAC,KAAKC,qBAAqB,OAAO,MAAM,GAAG,OAAO,KAAA;EAEhE,OAAO;CACT;CAEA,sBACE,SACA,QACoD;EACpD,MAAM,QAAQ,KAAKR,UAAU,MAAM;EACnC,MAAM,MAAM,KAAKO,WAAW,QAAQ,YAAY,QAAQ,QAAQ;EAChE,MAAM,SAAS,MAAM,eAAe,IAAI,QAAQ,KAAK;EACrD,MAAM,cAAc,KAAKD,gBAAgB,OAAO,QAAQ,KAAK;EAC7D,IAAI,CAAC,UAAU,MAAM,kBAAkB,IAAI,QAAQ,KAAK,MAAM,OAAO,CAAC,aACpE;EAGF,MAAM,aAAa,OAAO,cAAc,MAAM,0BAA0B,IAAI,QAAQ,KAAK;EACzF,IAAI,QAAQ,cAAc,YAAY,cAAc,WAAW,eAAe,QAAQ,YACpF;EAGF,OAAO;GAAE,OAAO,QAAQ;GAAO,YAAY,QAAQ,cAAc,YAAY;EAAW;CAC1F;CAEA,YAAY,SAAwC,QAA0B;EAC5E,MAAM,iBAAiB,KAAKd,WAAW,MAAM;EAC7C,MAAM,QAAQ,KAAKQ,UAAU,cAAc;EAC3C,MAAM,MAAM,KAAKO,WAAW,QAAQ,YAAY,QAAQ,QAAQ;EAChE,MAAM,QAAQ,KAAK,qBAAqB,SAAS,cAAc;EAC/D,IAAI,CAAC,OAAO,OAAO;EACnB,IAAI,MAAM,iBAAiB,IAAI,KAAK,GAAG,OAAO,KAAK,SAAS,OAAO,cAAc;EACjF,IAAI,MAAM,kBAAkB,IAAI,KAAK,MAAM,KAAK;GAI9C,KAAK,SAAS,OAAO,cAAc;GACnC,OAAO;EACT;EACA,IAAI,MAAM,wBAAwB,IAAI,KAAK,MAAM,KAAK,OAAO;EAC7D,MAAM,WAAW,MAAM,sBAAsB,IAAI,GAAG;EACpD,IAAI,CAAC,UAAU,OAAO;EACtB,KAAKQ,SAAS,gBAAgB,KAAK;GAAE,MAAM;GAAuB;GAAO;EAAS,CAAC;EACnF,OAAO;CACT;;CAGA,gBAAgB;EACd,KAAK,MAAM,UAAU,CAAC,wBAAwB,GAAG;GAC/C,KAAKC,YAAY,MAAM;GACvB,OAAiD,QAAQ;EAC3D;EACA,2BAA2B,IAAIzB,sBAAAA,mBAAmB;CACpD;CAEA,YAAY,QAAgB;EAC1B,MAAM,QAAQ,KAAKa,gBAAgB,IAAI,MAAM;EAC7C,IAAI,CAAC,OAAO;EAEZ,MAAM,iBAAiB,SAAQ,gBAAe;GAC5C,YAAY,gBAAgB,MAAM;GAClC,YAAY,QAAQ;EACtB,CAAC;EACD,MAAM,mBAAmB,SAAQ,UAAS,cAAc,KAAK,CAAC;EAC9D,MAAM,mBAAmB,MAAM;EAC/B,MAAM,eAAe,MAAM;EAC3B,MAAM,qBAAqB,MAAM;EACjC,MAAM,kBAAkB,MAAM;EAC9B,MAAM,wBAAwB,MAAM;EACpC,MAAM,mBAAmB,MAAM;EAC/B,MAAM,wBAAwB,MAAM;EACpC,MAAM,gBAAgB,MAAM;EAC5B,MAAM,0BAA0B,MAAM;EACtC,MAAM,uBAAuB,MAAM;EACnC,MAAM,sBAAsB,MAAM;EAClC,MAAM,2BAA2B,MAAM;EACvC,MAAM,6BAA6B,MAAM;EACzC,MAAM,sBAAsB,MAAM;EAClC,MAAM,iBAAiB,MAAM;EAC7B,MAAM,uBAAuB,MAAM;EACnC,MAAM,iBAAiB,MAAM;EAC7B,MAAM,cAAc,MAAM;CAC5B;CAEA,oBAAoB,OAAgC,OAAe;EACjE,MAAM,iBAAiB,IAAI,KAAK,CAAC,EAAE,QAAQ;EAC3C,MAAM,iBAAiB,OAAO,KAAK;EACnC,MAAM,cAAc,OAAO,KAAK;CAClC;CAEA,MAAMa,eACJ,OACA,QACA,YACA,UACA,gBACA;EAIA,IAAI,OAAO,WAAW;EACtB,MAAM,SAAS,MAAM,MAAM,UAAU,EAAE,eAAe,CAAC;EACvD,IAAI,CAAC,QAAQ;EACb,MAAM,OAAO,aAAa,EACxB,UAAU,CAAC,OAAO,YAAY;GAAE;GAAY;EAAS,CAAC,CAAC,EACzD,CAAC;CACH;CAEA,0BACE,OACA,QACA,KACA,OACA,QACA,YACA,UACA;EACA,IAAI;EACJ,MAAM,WAAW,IAAI,SAAc,YAAW;GAC5C,SAAS;EACX,CAAC;EACD,MAAM,QAAe;GACnB;IAAE,MAAM;IAAS;GAAM;GACvB;IAAE,GAAG,OAAO,WAAW;IAAG;GAAM;GAChC;IACE,MAAM;IACN;IACA,SAAS;KACP,YAAY,EAAE,QAAQ,OAAO;KAC7B,QAAQ,EACN,OAAO;MAAE,aAAa;MAAG,cAAc;MAAG,aAAa;KAAE,EAC3D;IACF;GACF;EACF;EACA,MAAM,SAAS;GACb;GACA,QAAQ;GACR,YAAY,IAAI,eAAe,EAC7B,MAAM,YAAY;IAChB,KAAK,MAAM,QAAQ,OAAO,WAAW,QAAQ,IAAI;IACjD,WAAW,MAAM;IACjB,OAAO;GACT,EACF,CAAC;GACD,0BAA0B;EAC5B;EACA,MAAM,EAAE,UAAU,cAAc,KAAKC,oBAAoB,OAAO,KAAK;EACrE,MAAM,EACJ,QAAQ,sBACR,wBACA,mBACE,KAAKC,qBAAqB,QAAQ,QAAQ,KAAK,QAAQ;EAC3D,MAAM,SAAoC;GACxC,OAAO,EAAE,IAAI,oBAAoB,OAAO,KAAK;GAC7C,QAAQ;GACR;GACA;GACA;GACA,WAAW;GACX;GACA;GACA,eAAe,CAAC;GAChB;EACF;EAEA,MAAM,eAAe,IAAI,OAAO,MAAM;EACtC,MAAM,qBAAqB,IAAI,UAAU,MAAM;EAC/C,MAAM,kBAAkB,IAAI,OAAO,GAAG;EACtC,MAAM,sBAAsB,IAAI,KAAK,QAAQ;EAE7C,KADwBV,gBAAgB,QAAQ,KAAK;GAAE,MAAM;GAAkB;GAAO;GAAU;EAAU,CAC5F,CAAC,CAAC,KAAK,gBAAgB,cAAc;EACnD,qBAA0B,mBAAmB,CAAC,CAAC,cAAc;GAC3D,iBAAiB;IACf,MAAM,qBAAqB,OAAO,QAAQ;IAC1C,IAAI,MAAM,eAAe,IAAI,KAAK,MAAM,QAAQ;KAC9C,MAAM,eAAe,OAAO,KAAK;KACjC,MAAM,kBAAkB,OAAO,KAAK;IACtC;IACA,IAAI,MAAM,mBAAmB,IAAI,GAAG,MAAM,SAAS,MAAM,sBAAsB,IAAI,GAAG,MAAM,UAAU;KACpG,MAAM,mBAAmB,OAAO,GAAG;KACnC,MAAM,sBAAsB,OAAO,GAAG;IACxC;IACA,KAAKK,oBAAoB,QAAQ,KAAK,KAAK;IAC3C,KAAKC,SAAS,QAAQ,KAAK;KAAE,MAAM;KAAiB;KAAO;IAAS,CAAC;GACvE,GAAG,CAAC;EACN,CAAC;CACH;CAEA,MAAMK,+BACJ,OACA,QACA,KACA,OACA,OACA,QACA,YACA,UACA,gBACA;EACA,IAAI,OAAO,WAAW;EAEtB,MAAM,KAAKH,eAAe,OAAO,QAAQ,YAAY,UAAU,cAAc;EAC7E,KAAKI,0BAA0B,OAAO,QAAQ,KAAK,OAAO,QAAQ,YAAY,QAAQ;CACxF;;;;;;;;;;;;;;;;CAiBA,4BAA4B,OAAgC,QAA4B;EACtF,MAAM,MAAM,KAAK,IAAI;EACrB,KAAK,MAAM,CAAC,UAAU,WAAW,MAAM,sBAAsB;GAC3D,IAAI,OAAO,cAAc,eAAe,OAAO,gBAAgB,KAAA,GAAW;GAC1E,IAAI,MAAM,OAAO,eAAe,4BAA4B;GAC5D,MAAM,qBAAqB,OAAO,QAAQ;GAC1C,MAAM,uBAAuB,OAAO,QAAQ;GAK5C,IAAI,MAAM,eAAe,IAAI,OAAO,KAAK,MAAM,QAAQ;GACvD,MAAM,WAAW,KAAKd,WAAW,OAAO,YAAY,OAAO,QAAQ;GACnE,MAAM,eAAe,OAAO,OAAO,KAAK;GACxC,MAAM,kBAAkB,OAAO,OAAO,KAAK;GAC3C,KAAKe,mBAAmB,OAAO,OAAO,KAAK;GAG3C,KAAKR,oBAAoB,QAAQ,UAAU,OAAO,KAAK;GACvD,IACE,MAAM,mBAAmB,IAAI,QAAQ,MAAM,OAAO,SAClD,MAAM,sBAAsB,IAAI,QAAQ,MAAM,UAC9C;IACA,MAAM,mBAAmB,OAAO,QAAQ;IACxC,MAAM,sBAAsB,OAAO,QAAQ;GAC7C;GACA,KAAKC,SAAS,QAAQ,UAAU;IAAE,MAAM;IAAiB,OAAO,OAAO;IAAO;GAAS,CAAC;EAC1F;CACF;CAEA,YACE,OACA,QACA,eACA,QAC2B;EAC3B,MAAM,EAAE,UAAU,eAAe,KAAKF,iBAAiB,aAAa;EACpE,IAAI,CAAC,UAAU;EAEf,MAAM,QAAQ,KAAKb,UAAU,MAAM;EACnC,KAAKuB,4BAA4B,OAAO,MAAM;EAC9C,MAAM,MAAM,KAAKhB,WAAW,YAAY,QAAQ;EAChD,MAAM,EAAE,UAAU,cAAc,KAAKW,oBAAoB,OAAO,OAAO,KAAK;EAC5E,MAAM,EACJ,QAAQ,sBACR,wBACA,mBACE,KAAKC,qBAAqB,QAAQ,QAAQ,KAAK,QAAQ;EAC3D,MAAM,SAAuC;GAC3C;GACA,QAAQ;GACR,OAAO,OAAO;GACd;GACA;GACA,WAAW;GACX;GACA;GACe;GACf;EACF;EAEA,KAAKG,mBAAmB,OAAO,OAAO,KAAK;EAC3C,MAAM,eAAe,IAAI,OAAO,OAAO,MAAM;EAC7C,MAAM,qBAAqB,IAAI,UAAU,MAAM;EAC/C,MAAM,kBAAkB,IAAI,OAAO,OAAO,GAAG;EAC7C,MAAM,mBAAmB,IAAI,KAAK,OAAO,KAAK;EAC9C,MAAM,sBAAsB,IAAI,KAAK,QAAQ;EAC7C,MAAM,iBAAiB,KAAK9B,WAAW,MAAM;EAC7C,MAAM,cAAc,YAAY;GAmB9B,KAAI,MAHgB,KAAKG,kBAAkB,cAAc,CAAC,CACvD,aAAa,KAAK,OAAO,OAAO,yBAAyB,CAAC,CAC1D,aAAa,EAAE,UAAU,KAAgB,EAAE,EAAA,CACpC,UAAU,KAAKO,mBAAmB,gBAAgB,KAAK,OAAO,KAAK;GAC7E,MAAM,KAAKO,gBAAgB,QAAQ,KAAK;IACtC,MAAM;IACN,OAAO,OAAO;IACd;IACA;GACF,CAAC;EACH,EAAA,CAAG;EAOH,WAAgB,KAAK,gBAAgB,cAAc;EACnD,KAAKe,0BAA0B,OAAO,QAAQ,KAAK,MAAM;EACzD,OAAO;CACT;CAEA,0BACE,OACA,QACA,KACA,QACA;EACA,IAAI,MAAM,uBAAuB,IAAI,OAAO,QAAQ,GAAG;EACvD,MAAM,uBAAuB,IAAI,OAAO,QAAQ;EAEhD,OAAY,OAAO,mBAAmB,CAAC,CAAC,cAAc;GACpD,MAAM,uBAAuB,OAAO,OAAO,QAAQ;GACnD,KAAKC,oBAAoB,OAAO,OAAO,KAAK;GAE5C,IAAI,OAAO,OAAO,WAAW,eAAe,KAAKnB,gBAAgB,OAAO,OAAO,KAAK,GAAG;IACrF,OAAO,YAAY;IAMnB,OAAO,cAAc,KAAK,IAAI;IAC9B,KAAKS,SAAS,QAAQ,KAAK;KAAE,MAAM;KAAiB,OAAO,OAAO;KAAO,UAAU,OAAO;IAAS,CAAC;IACpG;GACF;GAEA,OAAO,YAAY;GACnB,KAAKO,mBAAmB,OAAO,OAAO,KAAK;GAC3C,MAAM,qBAAqB,OAAO,OAAO,QAAQ;GACjD,IAAI,MAAM,eAAe,IAAI,OAAO,KAAK,MAAM,QAAQ;IACrD,MAAM,eAAe,OAAO,OAAO,KAAK;IACxC,MAAM,kBAAkB,OAAO,OAAO,KAAK;GAC7C;GAEA,IACE,MAAM,mBAAmB,IAAI,GAAG,MAAM,OAAO,SAC7C,MAAM,sBAAsB,IAAI,GAAG,MAAM,OAAO,UAChD;IACA,MAAM,mBAAmB,OAAO,GAAG;IACnC,MAAM,sBAAsB,OAAO,GAAG;GACxC;GASA,KAAKP,SAAS,QAAQ,KAAK;IAAE,MAAM;IAAiB,OAAO,OAAO;IAAO,UAAU,OAAO;GAAS,CAAC;GACpG,IAAI,KAAKW,sBAAsB,OAAO,GAAG,GACvC,KAAUC,qBAAqB,OAAO,QAAQ,KAAK,MAAM;QAEzD,KAAKb,oBAAoB,QAAQ,KAAK,OAAO,KAAK;EAEtD,CAAC;CACH;CAEA,MAAMa,qBACJ,OACA,QACA,KACA,aACA;EACA,IAAI,MAAM,mBAAmB,IAAI,GAAG,GAClC;EAMF,MAAM,iBAAiB,MAAM,sBAAsB,IAAI,GAAG;EAC1D,IAAI,gBAAgB,QAAQ;GAC1B,MAAM,sBAAsB,OAAO,GAAG;GACtC,MAAM,uBAAuB,IAAI,KAAK,CAAC,GAAG,gBAAgB,GAAI,MAAM,uBAAuB,IAAI,GAAG,KAAK,CAAC,CAAE,CAAC;EAC7G;EAEA,MAAM,QAAQ,MAAM,uBAAuB,IAAI,GAAG;EAClD,MAAM,SAAS,OAAO,MAAM;EAC5B,IAAI,UAAU,OAAO;GACnB,IAAI,MAAM,WAAW,GACnB,MAAM,uBAAuB,OAAO,GAAG;GAQzC,MAAM,aAAA,GAAA,OAAA,WAAA,CAAuB;GAC7B,MAAM,mBAAmB,IAAI,KAAK,SAAS;GAC3C,MAAM,kBAAkB,IAAI,WAAW,GAAG;GAC1C,MAAM,OAAO,MAAM,KAAKxB,8BAA8B,QAAQ,KAAK,WAAW,YAAY,KAAK;GAC/F,IAAI,CAAC,KAAK,UAAU;IAClB,IAAI,MAAM,mBAAmB,IAAI,GAAG,MAAM,WACxC,MAAM,mBAAmB,OAAO,GAAG;IAErC,MAAM,kBAAkB,OAAO,SAAS;IAGxC,MAAM,sBAAsB,OAAO,GAAG;IAGtC,MAAM,WAAW,MAAM,uBAAuB,IAAI,GAAG,KAAK,CAAC;IAC3D,MAAM,uBAAuB,IAAI,KAAK,CAAC,QAAQ,GAAG,QAAQ,CAAC;IAC3D,IAAI,KAAK,OAAO;KACd,MAAM,KAAKM,gBAAgB,QAAQ,KAAK;MACtC,MAAM;MACN,OAAO,KAAK;MACZ,QAAQ,KAAKmB,iBAAiB,MAAM;MACpC,UAAU,KAAKhB,aAAa;KAC9B,CAAC,CAAC,CAAC,YAAY,CAAC,CAAC;KACjB,MAAM,uBAAuB,IAAI,GAAG,CAAC,EAAE,MAAM;KAC7C,KAAK,MAAM,uBAAuB,IAAI,GAAG,CAAC,EAAE,UAAU,OAAO,GAC3D,MAAM,uBAAuB,OAAO,GAAG;IAE3C;IACA;GACF;GAEA,MAAM,SAAS,MAAM,YAAY,MAAM,OAAO,QAAQ;IACpD,GAAI,YAAY;IAChB,OAAO;IACP,QAAQ,iBACN,YAAY,cAAc,QAC1B,YAAY,cAAc,IAC1B,YAAY,YAAY,EAC1B;GACF,CAAC;GAED,IAAI,MAAM,SAAS,GAAG;IACpB,MAAM,aAAa,MAAM,eAAe,IAAI,OAAO,KAAK;IACxD,IAAI,YACF,KAAKY,0BAA0B,OAAO,QAAQ,KAAK,UAAU;GAEjE;GACA;EACF;EAEA,IAAI,MAAM,KAAKK,2BAA2B,OAAO,QAAQ,KAAK,YAAY,KAAK,GAC7E;EAGF,IAAI,MAAM,KAAKC,yBAAyB,OAAO,QAAQ,KAAK,YAAY,KAAK,GAC3E;EAIF,KAAKhB,oBAAoB,QAAQ,KAAK,YAAY,KAAK;CACzD;CAEA,MAAMe,2BACJ,OACA,QACA,KACA,WACA;EACA,IAAI,MAAM,mBAAmB,IAAI,GAAG,GAClC,OAAO;EAGT,MAAM,QAAQ,MAAM,6BAA6B,IAAI,GAAG;EACxD,MAAM,UAAU,OAAO,MAAM;EAC7B,IAAI,CAAC,WAAW,CAAC,OACf,OAAO;EAET,IAAI,MAAM,WAAW,GACnB,MAAM,6BAA6B,OAAO,GAAG;EAO/C,IAAI,WAAW;GACb,MAAM,mBAAmB,IAAI,KAAK,QAAQ,KAAK;GAC/C,MAAM,kBAAkB,IAAI,QAAQ,OAAO,GAAG;GAE9C,IAAI,EAAC,MADc,KAAK1B,8BAA8B,QAAQ,KAAK,QAAQ,OAAO,SAAS,EAAA,CACjF,UAAU;IAClB,IAAI,MAAM,mBAAmB,IAAI,GAAG,MAAM,QAAQ,OAChD,MAAM,mBAAmB,OAAO,GAAG;IAErC,MAAM,kBAAkB,OAAO,QAAQ,KAAK;IAC5C,MAAM,sBAAsB,OAAO,GAAG;IACtC,MAAM,WAAW,MAAM,6BAA6B,IAAI,GAAG,KAAK,CAAC;IACjE,MAAM,6BAA6B,IAAI,KAAK,CAAC,SAAS,GAAG,QAAQ,CAAC;IAClE,OAAO;GACT;EACF;EAEA,KAAK4B,mBAAmB,OAAO,QAAQ,KAAK,OAAO;EACnD,OAAO;CACT;CAEA,mBACE,OACA,QACA,KACA,SACA;EACA,MAAM,mBAAmB,IAAI,KAAK,QAAQ,KAAK;EAC/C,MAAM,kBAAkB,IAAI,QAAQ,OAAO,GAAG;EAC9C,QAAa,MACV,OAAO,QAAQ,UAAU;GACxB,GAAI,QAAQ;GACZ,OAAO,QAAQ;GACf,QAAQ,iBAAiB,QAAQ,eAAe,QAAQ,QAAQ,YAAY,QAAQ,QAAQ;EAC9F,CAAC,CAAC,CACD,MAAK,WAAU;GACd,KAAK,MAAM,6BAA6B,IAAI,GAAG,CAAC,EAAE,UAAU,KAAK,GAAG;IAClE,MAAM,aAAa,MAAM,eAAe,IAAI,OAAO,KAAK;IACxD,IAAI,YACF,KAAKP,0BAA0B,OAAO,QAAQ,KAAK,UAAU;GAEjE;EACF,CAAC,CAAC,CACD,OAAM,QAAO;GACZ,MAAM,kBAAkB,OAAO,QAAQ,KAAK;GAC5C,KAAKC,oBAAoB,OAAO,QAAQ,KAAK;GAC7C,IAAI,MAAM,mBAAmB,IAAI,GAAG,MAAM,QAAQ,OAChD,MAAM,mBAAmB,OAAO,GAAG;GAErC,KAAKV,SAAS,QAAQ,KAAK;IACzB,MAAM;IACN,OAAO,QAAQ;IACf,OAAOiB,cAAAA,oBAAoB,GAAG,CAAC,CAAC;GAClC,CAAC;GAGD,KAAUH,2BAA2B,OAAO,QAAQ,KAAK,QAAQ,KAAK,CAAC,CAAC,KAAK,OAAM,YAAW;IAC5F,IAAI,SAAS;IACb,IAAI,MAAM,KAAKC,yBAAyB,OAAO,QAAQ,KAAK,QAAQ,KAAK,GAAG;IAC5E,KAAKhB,oBAAoB,QAAQ,KAAK,QAAQ,KAAK;GACrD,CAAC;EACH,CAAC;CACL;CAEA,qBACE,OACA,UACA,QACA,QACmC;EACnC,MAAM,QAAQ,KAAKd,UAAU,MAAM;EACnC,MAAM,MAAM,KAAKO,WAAW,OAAO,YAAY,OAAO,QAAQ;EAC9D,MAAM,QAAQ,OAAO,UAAA,GAAA,OAAA,WAAA,CAAoB;EACzC,MAAM,UAAuC;GAC3C;GACA;GACA;GACA,YAAY,OAAO;GACnB,UAAU,OAAO;GACjB,eAAe,OAAO;EACxB;EAEA,MAAM,cAAc,MAAM,mBAAmB,IAAI,GAAG;EACpD,MAAM,eAAe,cAAc,MAAM,eAAe,IAAI,WAAW,IAAI,KAAA;EAC3E,IAAI,MAAM,mBAAmB,IAAI,GAAG,GAAG;GACrC,MAAM,QAAQ,MAAM,6BAA6B,IAAI,GAAG,KAAK,CAAC;GAC9D,MAAM,KAAK,OAAO;GAClB,MAAM,6BAA6B,IAAI,KAAK,KAAK;GACjD,IAAI,cACF,KAAKiB,0BAA0B,OAAO,QAAQ,KAAK,YAAY;GAEjE,OAAO;IAAE,UAAU;IAAM;GAAM;EACjC;EAEA,KAAKO,mBAAmB,OAAO,QAAQ,KAAK,OAAO;EACnD,OAAO;GAAE,UAAU;GAAM;EAAM;CACjC;CAEA,MAAMD,yBACJ,OACA,QACA,KACA,WACkB;EAClB,IAAI,MAAM,mBAAmB,IAAI,GAAG,GAClC,OAAO;EAGT,MAAM,YAAY,MAAM,2BAA2B,IAAI,GAAG;EAC1D,MAAM,cAAc,WAAW,MAAM;EACrC,IAAI,CAAC,eAAe,CAAC,WACnB,OAAO;EAET,IAAI,UAAU,WAAW,GACvB,MAAM,2BAA2B,OAAO,GAAG;EAG7C,MAAM,mBAAmB,IAAI,KAAK,YAAY,KAAK;EACnD,MAAM,kBAAkB,IAAI,YAAY,OAAO,GAAG;EAQlD,MAAM,OAAO,MAAM,KAAK3B,8BAA8B,QAAQ,KAAK,YAAY,OAAO,SAAS;EAC/F,IAAI,CAAC,KAAK,UAAU;GAIlB,IAAI,MAAM,mBAAmB,IAAI,GAAG,MAAM,YAAY,OACpD,MAAM,mBAAmB,OAAO,GAAG;GAErC,MAAM,kBAAkB,OAAO,YAAY,KAAK;GAChD,MAAM,sBAAsB,OAAO,GAAG;GACtC,IAAI,KAAK,OACP,MAAM,KAAKM,gBAAgB,QAAQ,KAAK;IACtC,MAAM;IACN,OAAO,KAAK;IACZ,QAAQ,KAAKmB,iBAAiB,YAAY,MAAM;IAChD,UAAU,KAAKhB,aAAa;GAC9B,CAAC,CAAC,CAAC,YAAY,CAAC,CAAC;GAEnB,MAAM,KAAKkB,yBAAyB,OAAO,QAAQ,KAAK,SAAS;GACjE,OAAO;EACT;EAEA,IAAI;GACF,MAAM,SAAS,MAAM,YAAY,MAAM,OAAO,YAAY,QAAQ;IAChE,GAAI,YAAY;IAChB,OAAO,YAAY;IACnB,QAAQ,iBAAiB,YAAY,eAAe,QAAQ,YAAY,YAAY,YAAY,QAAQ;GAC1G,CAAC;GAED,KAAK,WAAW,UAAU,KAAK,GAAG;IAChC,MAAM,aAAa,MAAM,eAAe,IAAI,OAAO,KAAK;IACxD,IAAI,YACF,KAAKN,0BAA0B,OAAO,QAAQ,KAAK,UAAU;GAEjE;EACF,SAAS,KAAK;GACZ,MAAM,kBAAkB,OAAO,YAAY,KAAK;GAChD,KAAKC,oBAAoB,OAAO,YAAY,KAAK;GACjD,IAAI,MAAM,mBAAmB,IAAI,GAAG,MAAM,YAAY,OACpD,MAAM,mBAAmB,OAAO,GAAG;GAErC,KAAKV,SAAS,QAAQ,KAAK;IACzB,MAAM;IACN,OAAO,YAAY;IACnB,OAAOiB,cAAAA,oBAAoB,GAAG,CAAC,CAAC;GAClC,CAAC;GAED,IAAI,CAAE,MAAM,KAAKF,yBAAyB,OAAO,QAAQ,KAAK,YAAY,KAAK,GAC7E,KAAKhB,oBAAoB,QAAQ,KAAK,YAAY,KAAK;EAE3D;EACA,OAAO;CACT;;;;;;;;;CAUA,oBAAoB,OAAe,QAAiB,QAA+B,WAAiC;EAClH,MAAM,QAAQ,KAAKd,UAAU,MAAM;EACnC,MAAM,SAAS,MAAM,eAAe,IAAI,KAAK;EAC7C,MAAM,MAAM,SAAS,KAAKO,WAAW,OAAO,YAAY,OAAO,QAAQ,IAAI,MAAM,kBAAkB,IAAI,KAAK;EAC5G,IAAI,CAAC,KAAK,OAAO,CAAC;EAElB,MAAM,kBAAkB,UAAU,YAAY,MAAM,wBAAwB,MAAM;EAClF,MAAM,QAAQ,gBAAgB,IAAI,GAAG;EACrC,IAAI,CAAC,SAAS,MAAM,WAAW,GAC7B,OAAO,CAAC;EAGV,gBAAgB,OAAO,GAAG;EAC1B,OAAO;CACT;CAEA,MAAM,2BACJ,OACA,SACA,QACA;EACA,MAAM,EAAE,UAAU,eAAe,KAAKM,iBAAiB,OAAO;EAC9D,IAAI,CAAC,UAAU;EAEf,MAAM,QAAQ,KAAKb,UAAU,MAAM;EACnC,MAAM,MAAM,KAAKO,WAAW,YAAY,QAAQ;EAChD,OAAO,MAAM;GACX,MAAM,cAAc,MAAM,mBAAmB,IAAI,GAAG;GACpD,IAAI,CAAC,aAAa;GAElB,MAAM,eAAe,MAAM,eAAe,IAAI,WAAW;GACzD,IAAI,cAAc;IAChB,IAAI,aAAa,MAAM,OAAO,MAAM,MAAM,CAAC,KAAKC,qBAAqB,OAAO,YAAY,GACtF;IAEF,MAAM,aAAa,OAAO,mBAAmB,CAAC,CAAC,YAAY,CAAC,CAAC;IAC7D;GACF;GAEA,IAAI,MAAM,kBAAkB,IAAI,WAAW,MAAM,KAAK;GAEtD,MAAM,KAAKyB,0BAA0B,QAAQ,KAAK,WAAW;EAC/D;CACF;CAEA,MAAMA,0BAA0B,QAA4B,KAAa,OAAe;EACtF,MAAM,iBAAiB,KAAKzC,WAAW,MAAM;EAC7C,MAAM,EAAE,UAAU,eAAe,KAAKK,sBAAsB,cAAc;EAC1E,MAAM,QAAQ,KAAKa,aAAa,GAAG;EACnC,IAAI;EACJ,IAAI,aAAa;EACjB,IAAI,UAAU;EACd,IAAI;EACJ,MAAM,OAAO,IAAI,SAAc,YAAW;GACxC,cAAc;EAChB,CAAC;EACD,MAAM,qBAAqB,aAAsB;GAC/C,MAAM,QAAQ,KAAKV,UAAU,cAAc;GAC3C,IACE,MAAM,mBAAmB,IAAI,GAAG,MAAM,SACrC,YAAY,MAAM,sBAAsB,IAAI,GAAG,MAAM,UAEtD;GAEF,MAAM,mBAAmB,OAAO,GAAG;GACnC,MAAM,sBAAsB,OAAO,GAAG;GACtC,IAAI,MAAM,wBAAwB,IAAI,KAAK,MAAM,KAAK,MAAM,wBAAwB,OAAO,KAAK;EAClG;EACA,MAAM,eAAe;GACnB,IAAI,SAAS;GACb,UAAU;GACV,IAAI,OAAO,aAAa,KAAK;GAC7B,YAAY;EACd;EACA,MAAM,aAAa,YAAY;GAC7B,IAAI,SAAS;GACb,IAAI,YAAY;GAChB,MAAM,QAAQ,MAAM,SAAS,cAAc,GAAG,CAAC,CAAC,YAAY,KAAA,CAAS;GACrE,IAAI,SAAS;GACb,IAAI,UAAU,OAAO;IACnB,kBAAkB;IAClB,OAAO;IACP;GACF;GACA,QAAQ,iBAAiB,KAAK,WAAW,GAAG,yBAAyB;EACvE;EACA,MAAM,WAAyB,UAAS;GACtC,MAAM,OAAO,MAAM;GACnB,KACG,MAAM,SAAS,mBAAmB,MAAM,SAAS,iBAAiB,MAAM,SAAS,iBAClF,KAAK,UAAU,OACf;IACA,kBAAkB,KAAK,QAAQ;IAC/B,OAAO;GACT;EACF;EAEA,IAAI;GACF,MAAM,eAAe,UAAU,OAAO,OAAO;GAC7C,aAAa;GACb,IAAI,CAAC,YAAY,QAAQ,iBAAiB,KAAK,WAAW,GAAG,yBAAyB;GACtF,MAAM;EACR,QAAQ;GACN,OAAO;GACP,MAAM;EACR,UAAU;GACR,IAAI,OAAO,aAAa,KAAK;GAC7B,IAAI,YAAY,MAAM,eAAe,YAAY,OAAO,OAAO,CAAC,CAAC,YAAY,CAAC,CAAC;EACjF;CACF;CAEA,MAAM,kBACJ,OACA,SACA,QAC0C;EAE1C,MAAM,iBAAiB,KAAKR,WAAW,MAAM;EAC7C,MAAM,QAAQ,KAAKQ,UAAU,cAAc;EAC3C,MAAM,MAAM,KAAKO,WAAW,QAAQ,YAAY,QAAQ,QAAQ;EAChE,MAAM,QAAQ,KAAKG,aAAa,GAAG;EACnC,MAAM,gCAAgB,IAAI,IAAY;EACtC,MAAM,cAA2C,CAAC;EAClD,MAAM,UAA6B,CAAC;EACpC,MAAM,6BAAa,IAAI,IASrB;EACF,IAAI,OAAO;EAEX,MAAM,aAAa;GACjB,OAAO,QAAQ,QAAQ,QAAQ,MAAM,CAAC,GAAG;EAC3C;EAEA,MAAM,oBAAoB;GACxB,MAAM,QAAQ,MAAM,mBAAmB,IAAI,GAAG;GAC9C,IAAI,CAAC,OAAO,OAAO;GACnB,MAAM,SAAS,MAAM,eAAe,IAAI,KAAK;GAI7C,IAAI,CAAC,QAAQ,OAAO;GACpB,OAAO,KAAKF,qBAAqB,OAAO,MAAM,IAAI,QAAQ;EAC5D;EAEA,MAAM,cAAc,WAAsC;GACxD,IAAI,QAAQ,cAAc,IAAI,OAAO,QAAQ,GAAG;GAChD,cAAc,IAAI,OAAO,QAAQ;GACjC,YAAY,KAAK,MAAM;GACvB,KAAK;EACP;EAEA,MAAM,mBAAmB,OAAe,UAAkB,cAAiD;GACzG,MAAM,YAAY;IAChB,OAAO,CAAC;IACR,SAAS,CAAC;IACV,eAAe,CAAC;IAChB,MAAM;IACN,QAAQ,KAAA;IACR,QAAQ;GACV;GACA,UAAU,SAAS,IAAI,eAAe;IACpC,KAAK,YAAY;KACf,MAAM,cAAc;MAClB,IAAI,UAAU,QAAQ;MACtB,OAAO,UAAU,MAAM,SAAS,GAC9B,WAAW,QAAQ,UAAU,MAAM,MAAM,CAAC;MAE5C,IAAI,UAAU,MAAM;OAClB,UAAU,SAAS;OACnB,WAAW,MAAM;MACnB;KACF;KACA,MAAM;KACN,IAAI,CAAC,UAAU,QAAQ,CAAC,UAAU,QAChC,UAAU,QAAQ,KAAK,KAAK;IAEhC;IACA,SAAS;KACP,UAAU,OAAO;KACjB,UAAU,SAAS;KACnB,UAAU,QAAQ,SAAS;KAC3B,OAAO,UAAU,cAAc,QAAQ,UAAU,cAAc,MAAM,CAAC,GAAG;IAC3E;GACF,CAAC;GACD,WAAW,IAAI,UAAU,SAAS;GAClC,OAAO;IACL;IACA,QAAQ;KACN;KACA,QAAQ;KACR,YAAY,UAAU;KACtB,oBAAoB,YAAY;MAC9B,IAAI,UAAU,MAAM;MACpB,MAAM,IAAI,SAAc,YAAW,UAAU,cAAc,KAAK,OAAO,CAAC;KAC1E;IACF;IACA;IACA;IACA;IACA,WAAW;IACX,UAAU,QAAQ;IAClB,YAAY,QAAQ;IACpB,eAAe,CAAC;GAClB;EACF;EAEA,MAAM,iCAAiB,IAAI,IAAY;EACvC,MAAM,oCAAoB,IAAI,IAAY;EAC1C,IAAI,gBAAyD;EAC7D,IAAI,oBAAmC;EACvC,IAAI,mBAAmB;EAEvB,MAAM,mBAAmB,OAAO,OAAe,UAAkB,UAAmB;GAClF,IAAI,CAAC,SAAS,CAAE,MAAM,KAAKZ,oBAAoB,gBAAgB,KAAK,KAAK,GAAI;GAC7E,MAAM,mBAAmB,IAAI,KAAK,KAAK;GACvC,MAAM,sBAAsB,IAAI,KAAK,QAAQ;GAC7C,IAAI,CAAC,OAAO,MAAM,wBAAwB,IAAI,OAAO,GAAG;EAC1D;EAEA,MAAM,wBAAwB,OAAe,aAAsB;GACjE,IACE,MAAM,mBAAmB,IAAI,GAAG,MAAM,SACrC,YAAY,MAAM,sBAAsB,IAAI,GAAG,MAAM,UAEtD;GAEF,MAAM,mBAAmB,OAAO,GAAG;GACnC,MAAM,sBAAsB,OAAO,GAAG;GACtC,IAAI,MAAM,wBAAwB,IAAI,KAAK,MAAM,KAAK,MAAM,wBAAwB,OAAO,KAAK;EAClG;EAEA,MAAM,cAAc,OAAO,UAAwC;GACjE,IAAI,MAAM;GACV,MAAM,OAAO,MAAM;GACnB,IAAI,CAAC,MAAM;GACX,IAAI,KAAK,SAAS,kBAAkB;IAClC,MAAM,cAAc,MAAM,qBAAqB,IAAI,KAAK,QAAQ;IAChE,IAAI,aACF,eAAe,IAAI,KAAK,QAAQ;SAEhC,kBAAkB,IAAI,KAAK,QAAQ;IAErC,MAAM,iBAAiB,KAAK,OAAO,KAAK,UAAU,QAAQ,WAAW,CAAC;IACtE,MAAM,SAAS,eAAe,gBAAgB,KAAK,OAAO,KAAK,UAAU,KAAK,SAAS;IACvF,WAAW,MAAM;IACjB,KAAK;IACL;GACF;GACA,IAAI,KAAK,SAAS,eAAe;IAC/B,IACE,KAAK,aAAa,KAAKE,QACtB,eAAe,IAAI,KAAK,QAAQ,KAAK,CAAC,kBAAkB,IAAI,KAAK,QAAQ,IAE1E;IAEF,IACE,MAAM,mBAAmB,IAAI,GAAG,MAAM,KAAK,SAC3C,MAAM,sBAAsB,IAAI,GAAG,MAAM,KAAK,UAE9C,MAAM,iBAAiB,KAAK,OAAO,KAAK,UAAU,KAAK;IAEzD,IAAI,YAAY,WAAW,IAAI,KAAK,QAAQ;IAC5C,IAAI,CAAC,WAAW;KAId,WAAW,gBAAgB,KAAK,OAAO,KAAK,UAAU,MAAM,iBAAiB,IAAI,KAAK,KAAK,KAAK,CAAC,CAAC;KAClG,YAAY,WAAW,IAAI,KAAK,QAAQ;KACxC,IAAI,CAAC,WAAW;IAClB;IACA,UAAU,MAAM,KAAK,KAAK,IAAI;IAC9B,OAAO,UAAU,QAAQ,QAAQ,UAAU,QAAQ,MAAM,CAAC,GAAG;IAC7D;GACF;GACA,IAAI,KAAK,SAAS,mBAAmB;IACnC,IAAI,KAAK,aAAa,KAAKA,KAAK;IAChC,MAAM,kBAAkB,KAAK,SAAS,MAAM,wBAAwB,MAAM;IAC1E,MAAM,QAAQ,gBAAgB,IAAI,GAAG,KAAK,CAAC;IAC3C,MAAM,KAAKoC,gBAAAA,aAAa,KAAK,MAAM,CAAC;IACpC,gBAAgB,IAAI,KAAK,KAAK;IAC9B;GACF;GACA,IAAI,KAAK,SAAS,uBAAuB;IACvC,IACE,MAAM,iBAAiB,IAAI,KAAK,KAAK,KACrC,MAAM,kBAAkB,IAAI,KAAK,KAAK,MAAM,OAC5C,MAAM,mBAAmB,IAAI,GAAG,MAAM,KAAK,SAC3C,MAAM,sBAAsB,IAAI,GAAG,MAAM,KAAK,YAC7C,MAAM,KAAKtC,oBAAoB,gBAAgB,KAAK,KAAK,KAAK,GAE/D,KAAK,SAAS,KAAK,OAAO,cAAc;IAE1C;GACF;GACA,IAAI,KAAK,SAAS,cAAc;IAC9B,MAAM,gBAAgB,KAAK,YAAY,KAAK;IAC5C,qBAAqB,KAAK,OAAO,KAAK,QAAQ;IAC9C,IAAI;IACJ,IAAI,YAAY,WAAW,IAAI,aAAa;IAC5C,IAAI,CAAC,WAAW;KACd,WAAW,gBAAgB,KAAK,OAAO,eAAe,MAAM,iBAAiB,IAAI,KAAK,KAAK,KAAK,CAAC;KACjG,YAAY,WAAW,IAAI,aAAa;IAC1C;IACA,IAAI,WAAW;KACb,UAAU,MAAM,KAAK;MAAE,MAAM;MAAS,SAAS,EAAE,OAAO,IAAI,MAAM,KAAK,KAAK,EAAE;KAAE,CAAC;KACjF,UAAU,OAAO;KACjB,OAAO,UAAU,QAAQ,QAAQ,UAAU,QAAQ,MAAM,CAAC,GAAG;KAC7D,OAAO,UAAU,cAAc,QAAQ,UAAU,cAAc,MAAM,CAAC,GAAG;KACzE,WAAW,OAAO,aAAa;KAC/B,cAAc,OAAO,aAAa;IACpC;IACA,IAAI,UAAU,WAAW,QAAQ;IACjC,MAAM,KAAKkC,yBAAyB,OAAO,gBAAgB,KAAK,KAAK,KAAK;IAC1E,KAAK;IACL;GACF;GACA,IAAI,KAAK,SAAS,mBAAmB,KAAK,SAAS,iBAAiB,KAAK,SAAS,iBAAiB;IACjG,MAAM,gBAAgB,KAAK,YAAY,KAAK;IAC5C,IAAI,KAAK,SAAS,iBAAiB;KACjC,MAAM,gBAAgB,IAAI,KAAK,KAAK;KACpC,MAAM,SAAS,MAAM,qBAAqB,IAAI,aAAa,KAAK,MAAM,eAAe,IAAI,KAAK,KAAK;KACnG,IAAI,QAAQ,OAAO,YAAY;IACjC,OACE,qBAAqB,KAAK,OAAO,KAAK,QAAQ;IAEhD,IAAI,KAAK,SAAS,iBAChB,KAAKR,mBAAmB,OAAO,KAAK,KAAK;IAE3C,MAAM,YAAY,WAAW,IAAI,aAAa;IAC9C,IAAI,WAAW;KACb,UAAU,OAAO;KACjB,OAAO,UAAU,QAAQ,QAAQ,UAAU,QAAQ,MAAM,CAAC,GAAG;KAC7D,OAAO,UAAU,cAAc,QAAQ,UAAU,cAAc,MAAM,CAAC,GAAG;KACzE,WAAW,OAAO,aAAa;KAC/B,cAAc,OAAO,aAAa;IACpC;IAGA,IAAI,KAAK,SAAS,iBAAiB,sBAAsB,KAAK,SAAS,eAAe;KACpF,mBAAmB;KACnB,IAAI;MACF,cAAmB,OAAO;KAC5B,QAAQ,CAAC;IACX;IACA,IAAI,KAAK,SAAS,iBAChB,MAAM,KAAKQ,yBAAyB,OAAO,gBAAgB,KAAK,KAAK,KAAK;IAE5E,KAAK;GACP;EACF;EAEA,IAAI,YAAY,QAAQ,QAAQ;EAChC,MAAM,WAAyB,UAAS;GACtC,YAAY,UAAU,WAAW,YAAY,KAAK,CAAC,CAAC,CAAC,YAAY,CAAC,CAAC;EACrE;EAEA,MAAM,eAAe,UAAU,OAAO,OAAO;EAE7C,MAAM,eAAe,YAAY;EACjC,MAAM,gBAAgB,eAAe,MAAM,eAAe,IAAI,YAAY,IAAI,KAAA;EAC9E,IAAI,eAAe;GACjB,eAAe,IAAI,cAAc,QAAQ;GACzC,WAAW,aAAa;EAC1B;EAEA,MAAM,oBAAoB;GACxB,IAAI,MAAM;GACV,OAAO;GACP,eAAoB,YAAY,OAAO,OAAO,CAAC,CAAC,YAAY,CAAC,CAAC;GAE9D,IAAI,eACF,IAAI;IACF,cAAmB,OAAO;GAC5B,QAAQ,CAAC;GAEX,KAAK;EACP;EAEA,OAAO;GACL;GACA,aAAa,KAAK,YAAY,SAAS,cAAc;GACrD;GACA,SAAS,mBAAmB;IAC1B,IAAI;KACF,OAAO,CAAC,QAAQ,YAAY,SAAS,GAAG;MACtC,IAAI,YAAY,WAAW,GAAG;OAC5B,MAAM,IAAI,SAAc,YAAW,QAAQ,KAAK,OAAO,CAAC;OACxD;MACF;MACA,MAAM,MAAM,YAAY,MAAM;MAK9B,MAAM,UADmB,IAAI,yBAAyB,KAAK,IAAI,OAAO,WAAA,CACtC,UAAU;MAC1C,gBAAgB;MAChB,oBAAoB,IAAI;MACxB,IAAI,iBAAiB;MACrB,IAAI;OACF,OAAO,MAAM;QACX,MAAM,EAAE,OAAO,MAAM,MAAM,eAAe,MAAM,OAAO,KAAK;QAC5D,IAAI,YACF;QAEF,MAAM,YAAY;QAKlB,MAHE,aAAa,OAAO,cAAc,YAAY,EAAE,WAAW,aACvD;SAAE,GAAG;SAAW,OAAO,IAAI;QAAM,IACjC;QAEN,IAAI,MAAM;QACV,MAAM,eAAe,UAAU,gBAAgB,UAAU,SAAS;QAKlE,IAHE,UAAU,SAAS,WACnB,UAAU,SAAS,WAClB,UAAU,SAAS,YAAY,iBAAiB,cAC7B;SAIpB,iBAAiB;SACjB,CAAM,YAAY;UAChB,IAAI;WACF,OAAO,MAAM;YACX,MAAM,EAAE,MAAM,MAAM,MAAM,OAAO,KAAK;YACtC,IAAI,GAAG;WACT;UACF,QAAQ,CAAC;UACT,OAAO,YAAY;SACrB,EAAA,CAAG;SACH;QACF;OACF;OAIA,IAAI,CAAC,kBAAkB,CAAC,QAAQ,kBAAkB;QAChD,MAAM;SAAE,MAAM;SAAS,OAAO,IAAI;QAAM;QACxC,mBAAmB;OACrB;MACF,UAAU;OACR,gBAAgB;OAChB,oBAAoB;OACpB,IAAI,CAAC,gBACH,OAAO,YAAY;MAEvB;KACF;IACF,UAAU;KACR,YAAY;IACd;GACF,EAAA,CAAG;EACL;CACF;CAEA,YACE,OACA,SACA,QACA,QACgC;EAChC,OAAO,KAAK,WAAmB,OAAO,KAAKK,0BAA0B,OAAO,GAAG,QAAQ,MAAM;CAC/F;CAEA,aACE,OACA,SACA,QACA,QACiC;EACjC,MAAM,QAAQ,KAAKnC,UAAU,MAAM;EACnC,MAAM,6BAAa,IAAI,KAAK;EAC5B,IAAI;EACJ,IAAI,QAAQ,OAAO;EACnB,IAAI;EAEJ,IAAI,OAAO,cAAc,OAAO,UAAU;GACxC,MAAM,KAAKO,WAAW,OAAO,YAAY,OAAO,QAAQ;GACxD,MAAM,cAAc,MAAM,mBAAmB,IAAI,GAAG;GACpD,eAAe,cAAc,MAAM,eAAe,IAAI,WAAW,IAAI,KAAA;GACrE,IAAI,gBAAgB,CAAC,KAAKC,qBAAqB,OAAO,YAAY,GAAG;IACnE,MAAM,mBAAmB,OAAO,GAAG;IACnC,eAAe,KAAA;GACjB;GACA,UAAU;EACZ;EAEA,IAAI,OAAO;GACT,iBAAiB,MAAM,eAAe,IAAI,KAAK;GAC/C,IAAI,cACF,QAAQ,KAAKD,WAAW,aAAa,YAAY,aAAa,QAAQ;EAE1E;EAEA,MAAM,aAAa,OAAO,cAAc,cAAc;EACtD,MAAM,WAAW,OAAO,YAAY,cAAc;EAClD,IAAI,CAAC,cAAc,CAAC,UAClB,MAAM,IAAI,MAAM,yDAAyD;EAG3E,QAAQ,KAAKA,WAAW,YAAY,QAAQ;EAC5C,MAAM,SAAS6B,gBAAAA,oBAAoB,SAAS;GAC1C,IAAI,KAAKC,yBAAyB,OAAO;IAAE;IAAY;GAAS,CAAC;GACjE;EACF,CAAC;EACD,MAAM,eAAA,GAAA,OAAA,WAAA,CAAyB;EAC/B,MAAM,sBAAsB,OAAO,QAAQ,iBAAiB,cAAc;EAE1E,IAAI,cAAc;GAChB,MAAM,YAAY,MAAM,2BAA2B,IAAI,GAAG,KAAK,CAAC;GAChE,UAAU,KAAK;IAAE;IAAO;IAAQ,OAAO;IAAa;IAAY;IAAU,eAAe;GAAoB,CAAC;GAC9G,MAAM,2BAA2B,IAAI,KAAK,SAAS;GACnD,KAAKb,0BAA0B,OAAO,QAAQ,KAAK,YAAY;GAC/D,OAAO;IACL;IACA,UAAU,QAAQ,QAAQ;KAAE,QAAQ;KAAoB,OAAO;IAAY,CAAC;GAC9E;EACF;EAEA,OAAO,KAAK,WACV,OACA,QACA;GAAE,GAAG;GAAQ;GAAO;GAAY;GAAU,QAAQ;IAAE,GAAG,OAAO;IAAQ,UAAU;GAAO;EAAE,GACzF,MACF;CACF;CAEA,MAAM,gBACJ,OACA,YACA,QACA,QAC6C;EAC7C,IAAI,CAAC,OAAO,cAAc,CAAC,OAAO,UAChC,MAAM,IAAI,MAAM,6DAA6D;EAE/E,MAAM,aAAa,OAAO;EAC1B,MAAM,WAAW,OAAO;EAExB,MAAM,iBAAiB,OAAO,QAAQ,eAAe;EACrD,MAAM,gBAAgBc,cAAAA,0BAA0B,cAAc;EAC9D,MAAM,SAAS,MAAM,MAAM,UAAU,EAAE,eAAe,CAAC;EACvD,IAAI,CAAC,QACH,MAAM,IAAI,MAAM,wCAAwC;EAG1D,MAAM,eAAgB,MAAM,OAAO,cAAc,EAAE,SAAS,CAAC,KAAM,eAAe;EAClF,IAAI,CAAC,cACH,MAAM,IAAI,MAAM,yCAAyC,UAAU;EAYrE,MAAM,UAAU,MAAMC,cAAAA,iBAAiB;GACrC,OAAO;GACP;GACA,QAAA;IAXA,GAAG;IACH,IAAI;IACJ,YAAY,aAAa,cAAc;IACvC,WAAW,aAAa,6BAAa,IAAI,KAAK;IAC9C,WAAW,aAAa,6BAAa,IAAI,KAAK;IAC9C,UAAU,aAAa;GAMlB;GACL;GACA;GACA,cAAc,eAAe;GAC7B,4BAAY,IAAI,KAAK;EACvB,CAAC;EAED,IAAI,QAAQ,SACV,OAAO;GAAE,SAAS;GAAM,QAAQ;EAAY;EAG9C,OAAO,KAAK,WAAmB,OAAO,QAAQ,QAAQ,QAAQ,MAAM;CACtE;;;;;;;;;;;;;CAcA,WACE,OACA,aACA,QACA,QAC+B;EAC/B,MAAM,QAAQ,KAAKvC,UAAU,MAAM;EACnC,IAAI;EACJ,IAAI,QAAQ,OAAO;EACnB,MAAM,iBAAiB,OAAO,UAAU,YAAY;EACpD,MAAM,eAAe,OAAO,QAAQ,YAAY;EAEhD,IAAI;EACJ,IAAI,OAAO,cAAc,OAAO,UAAU;GACxC,MAAM,KAAKO,WAAW,OAAO,YAAY,OAAO,QAAQ;GACxD,MAAM,cAAc,MAAM,mBAAmB,IAAI,GAAG;GACpD,eAAe,cAAc,MAAM,eAAe,IAAI,WAAW,IAAI,KAAA;GACrE,IAAI,gBAAgB,CAAC,KAAKC,qBAAqB,OAAO,YAAY,GAAG;IACnE,MAAM,mBAAmB,OAAO,GAAG;IACnC,eAAe,KAAA;GACjB;GAIA,IAAI,gBAAgB,aAAa,MAAM,OAAO,MAAM,IAClD,QAAQ,aAAa;QAChB,IAAI,eAAe,CAAC,cACzB,IAAI,MAAM,kBAAkB,IAAI,WAAW,MAAM,KAG/C,QAAQ;QACH;IAEL,MAAM,mBAAmB,OAAO,GAAG;IACnC,MAAM,sBAAsB,OAAO,GAAG;GACxC;EAEJ;EAEA,IAAI,OAAO;GACT,iBAAiB,MAAM,eAAe,IAAI,KAAK;GAC/C,IAAI,cACF,QAAQ,KAAKD,WAAW,aAAa,YAAY,aAAa,QAAQ;EAE1E;EAEA,MAAM,aAAa,OAAO,cAAc,cAAc;EACtD,MAAM,WAAW,OAAO,YAAY,cAAc;EAClD,IAAI,CAAC,cAAc,CAAC,UAClB,MAAM,IAAI,MAAM,6CAA6C;EAG/D,MAAM,iBAAiB,QACrB,UAAU,cAAc,OAAO,WAAW,aAAc,OAAO,MAAM,mBAAmB,IAAI,GAAG,MAAM,MACvG;EACA,IAAI,SAAS2B,gBAAAA,aAAa;GACxB,GAAG;GACH,IAAI,YAAY,MAAM,KAAKG,yBAAyB,OAAO;IAAE;IAAY;GAAS,CAAC;GACnF,4BAAY,IAAI,KAAK;EACvB,CAAC;EAGD,SAASG,gBAAAA,0BACP,QACA,iBAAiB,OAAO,UAAU,aAAa,OAAO,QAAQ,UAChE;EAEA,IAAI,kBAAkB,mBAAmB,WAAW;GAClD,IAAI,mBAAmB,WAAW;IAChC,IAAI,CAAC,cAAc,CAAC,UAClB,MAAM,IAAI,MAAM,kEAAkE;IAIpF,IAAI,OAAO,WACT,OAAO;KACL;KACA,UAAU,QAAQ,QAAQ,EAAE,QAAQ,UAAmB,CAAC;IAC1D;IAEF,MAAM,YAAY,KAAKvB,eACrB,OACA,QACA,YACA,UACA,OAAO,QAAQ,eAAe,cAChC;IACA,UAAe,YAAY,CAAC,CAAC;IAC7B,OAAO;KACL;KACA;KACA,UAAU,QAAQ,QAAQ,EAAE,QAAQ,UAAmB,CAAC;IAC1D;GACF;GACA,OAAO;IACL;IACA,UAAU,QAAQ,QAAQ,EAAE,QAAQ,UAAmB,CAAC;GAC1D;EACF;EAEA,IAAI,OAAO;GAKT,IAAI,gBAAgB,KAAKT,qBAAqB,OAAO,YAAY,GAAG;IAClE,QAAQ,KAAKD,WAAW,aAAa,YAAY,aAAa,QAAQ;IACtE,IAAI,aAAa,MAAM,OAAO,MAAM,IAAI;KAGtC,MAAM,QAAQ,MAAM,uBAAuB,IAAI,GAAG,KAAK,CAAC;KACxD,MAAM,KAAK,MAAM;KACjB,MAAM,uBAAuB,IAAI,KAAK,KAAK;KAC3C,KAAKQ,SAAS,QAAQ,KAAK;MACzB,MAAM;MACN;MACA,QAAQ,KAAKa,iBAAiB,MAAM;MACpC,UAAU,KAAKhB,aAAa;KAC9B,CAAC;KACD,KAAKY,0BAA0B,OAAO,QAAQ,KAAK,YAAY;KAC/D,OAAO;MACL;MACA,UAAU,QAAQ,QAAQ;OAAE,QAAQ;OAAoB;MAAM,CAAC;KACjE;IACF;IAEA,OAAO;KACL;KACA,UAAU,QAAQ,QAAQ;MACxB,QAAQ;MACR,QAAQ;MACR,OAAO,aAAa;KACtB,CAAC;IACH;GACF;GAEA,IAAI,OAAO,MAAM,mBAAmB,IAAI,GAAG,MAAM,OAAO;IAMtD,MAAM,qBAAqB,MAAM,kBAAkB,IAAI,KAAK,MAAM;IAClE,IAAI,oBAAoB;KACtB,MAAM,QAAQ,MAAM,sBAAsB,IAAI,GAAG,KAAK,CAAC;KACvD,MAAM,KAAK,MAAM;KACjB,MAAM,sBAAsB,IAAI,KAAK,KAAK;IAC5C;IACA,KAAKT,SAAS,QAAQ,KAAK;KACzB,MAAM;KACN;KACA,QAAQ,KAAKa,iBAAiB,MAAM;KACpC,UAAU,KAAKhB,aAAa;KAC5B,QAAQ;IACV,CAAC;IACD,OAAO;KACL;KACA,UAAU,QAAQ,QAAQ;MAAE,QAAQ;MAAoB;KAAM,CAAC;IACjE;GACF;EACF;EAEA,SAAA,GAAA,OAAA,WAAA,CAAmB;EACnB,QAAQ,KAAKL,WAAW,YAAY,QAAQ;EAC5C,IAAI,iBAAiB,WAAW;GAG9B,IAAI,OAAO,WACT,OAAO;IACL;IACA,UAAU,QAAQ,QAAQ,EAAE,QAAQ,UAAmB,CAAC;GAC1D;GAEF,MAAM,YAAY,KAAKa,+BACrB,OACA,QACA,KACA,OACA,OACA,QACA,YACA,UACA,OAAO,QAAQ,eAAe,cAChC;GACA,UAAe,YAAY,CAAC,CAAC;GAC7B,OAAO;IACL;IACA;IACA,UAAU,QAAQ,QAAQ,EAAE,QAAQ,UAAmB,CAAC;GAC1D;EACF;EACA,IAAI,iBAAiB,QACnB,OAAO;GACL;GACA,UAAU,QAAQ,QAAQ,EAAE,QAAQ,UAAmB,CAAC;EAC1D;EAGF,IAAI,MAAM,mBAAmB,IAAI,GAAG,GAAG;GACrC,MAAM,gBAAgB,MAAM,mBAAmB,IAAI,GAAG;GACtD,MAAM,iBAAiB,gBAAgB,MAAM,eAAe,IAAI,aAAa;GAC7E,IACE,KAAKd,gBAAgB,OAAO,aAAa,KACzC,gBAAgB,OAAO,WAAW,eAClC,gBAAgB,cAAc,aAE9B,OAAO;IACL;IACA,UAAU,QAAQ,QAAQ;KACxB,QAAQ;KACR,QAAQ;KACR,OAAO;IACT,CAAC;GACH;GAKF,MAAM,YAAY,MAAM,2BAA2B,IAAI,GAAG,KAAK,CAAC;GAChE,UAAU,KAAK;IAAE;IAAO;IAAQ;IAAO;IAAY;IAAU,eAAe,OAAO,QAAQ;GAAc,CAAC;GAC1G,MAAM,2BAA2B,IAAI,KAAK,SAAS;GACnD,IAAI,cACF,KAAKkB,0BAA0B,OAAO,QAAQ,KAAK,YAAY;GAEjE,OAAO;IACL;IACA,UAAU,QAAQ,QAAQ;KAAE,QAAQ;KAAoB;IAAM,CAAC;GACjE;EACF;EAIA,MAAM,mBAAmB,IAAI,KAAK,KAAK;EACvC,MAAM,kBAAkB,IAAI,OAAO,GAAG;EACtC,MAAM,cAAc;EACpB,MAAM,gBAAgB;EACtB,MAAM,iBAAiB,KAAKhC,WAAW,MAAM;EAC7C,MAAM,gBAAgB,KAAKG,kBAAkB,cAAc;EAK3D,MAAM,YAAsD,YAAY;GAQtE,MAAM,QAAQ,MAAM,cACjB,aAAa,aAAa,eAAe,yBAAyB,CAAC,CACnE,aAAa;IAAE,UAAU;IAAiB,OAAO;GAAoC,EAAE;GAE1F,IAAI,CAAC,MAAM,UAAU;IAGnB,IAAI,MAAM,mBAAmB,IAAI,WAAW,MAAM,eAChD,MAAM,mBAAmB,OAAO,WAAW;IAE7C,MAAM,kBAAkB,OAAO,aAAa;IAC5C,MAAM,sBAAsB,OAAO,WAAW;IAM9C,MAAM,cAAc,MAAM;IAC1B,IAAI,aACF,MAAM,KAAKc,gBAAgB,QAAQ,aAAa;KAC9C,MAAM;KACN,OAAO;KACP,QAAQ,KAAKmB,iBAAiB,MAAM;KACpC,UAAU,KAAKhB,aAAa;IAC9B,CAAC,CAAC,CAAC,YAAY,CAAC,CAAC;IAEnB,OAAO;KAAE,QAAQ;KAAoB,OAAO,eAAe;IAAc;GAC3E;GAIA,KAAKV,mBAAmB,gBAAgB,aAAa,aAAa;GAClE,IAAI;IACF,MAAM,SAAS,MAAM,MAAM,OAAO,QAAQ;KACxC,GAAI,OAAO,QAAQ;KACnB,WAAW;KACX,OAAO;KACP,QAAQ,iBAAiB,OAAO,QAAQ,eAAe,QAAQ,YAAY,QAAQ;IACrF,CAAC;IACD,OAAO;KAAE,QAAQ;KAAiB,OAAO;KAAe;IAAO;GACjE,SAAS,OAAO;IACd,MAAM,kBAAkB,OAAO,aAAa;IAC5C,KAAKuB,oBAAoB,OAAO,aAAa;IAC7C,IAAI,MAAM,mBAAmB,IAAI,WAAW,MAAM,eAChD,MAAM,mBAAmB,OAAO,WAAW;IAE7C,KAAKX,oBAAoB,QAAQ,aAAa,aAAa;IAC3D,KAAKC,SAAS,QAAQ,aAAa;KACjC,MAAM;KACN,OAAO;KACP,OAAOiB,cAAAA,oBAAoB,KAAK,CAAC,CAAC;IACpC,CAAC;IACD,KAAUF,yBAAyB,OAAO,QAAQ,WAAW;IAC7D,MAAM;GACR;EACF,EAAA,CAAG;EAOH,SAAc,YAAY,CAAC,CAAC;EAE5B,OAAO;GACL;GACA;EACF;CACF;AACF;AAEA,MAAa,2BAA2B,IAAI,yBAAyB;;;;;;;;ACruErE,MAAa,qCAAqC;;;;;AAMlD,SAAgB,6BAA6B,QAG3C;CACA,MAAM,WAAW,OAAO,oBAAoB;CAG5C,IAAI,YAAA,GAAgD,OAAO;EAAE,kBAAkB;EAAU,QAAQ;CAAS;CAC1G,MAAM,mBAAmB,WAAW;CACpC,IAAI,mBAAA,GAAuD,OAAO,EAAE,iBAAiB;CACrF,OAAO;EAAE;EAAkB,QAAQ;CAAS;AAC9C;AAqBA,MAAM,qBAAqB,aAA+E;CACxG,IAAI,OAAO,aAAa,UAAU,OAAO,EAAE,QAAQ,SAAS;CAC5D,OAAO;AACT;AAEA,SAAgB,oCACd,OAC8B;CAC9B,IAAI,MAAM,OAAO,aAAa,UAC5B,OAAO;EAAE,QAAQ;EAAW,QAAQ;CAAS;CAG/C,IAAI,MAAM,OAAO,aAAa,QAC5B,OAAO,MAAM,gBAAgB,WACzB;EAAE,QAAQ;EAAa,WAAW,MAAM;EAAK,WAAW,MAAM;EAAK,QAAQ;CAAgC,IAC3G;EAAE,QAAQ;EAAW,QAAQ;CAAY;CAG/C,IAAI,MAAM,OAAO,aAAa,UAC5B,OAAO,MAAM,gBAAgB,WACzB;EAAE,QAAQ;EAAa,WAAW,MAAM;EAAK,QAAQ;CAAuB,IAC5E;EAAE,QAAQ;EAAW,QAAQ;CAAc;CAGjD,OAAO;EACL,QAAQ;EACR,WAAW,MAAM;EACjB,QAAQ,MAAM,gBAAgB,WAAW,yBAAyB;CACpE;AACF;AAEA,eAAsB,oCAAoC,EACxD,QACA,GAAG,SAGqC;CACxC,MAAM,SAAS,MAAM,QAAQ,SAAS,KAAK;CAC3C,IAAI,QAAQ,OAAO,kBAAkB,MAAM;CAE3C,MAAM,iBAAiB,QAAQ,UAAU,MAAM,OAAO;CACtD,IAAI,gBAAgB,OAAO,kBAAkB,cAAc;CAE3D,MAAM,mBAAmB,QAAQ,aAAa,MAAM,OAAO;CAC3D,IAAI,kBAAkB,OAAO,kBAAkB,gBAAgB;CAE/D,IAAI,QAAQ,SAAS,OAAO,kBAAkB,OAAO,OAAO;CAE5D,OAAO,oCAAoC,KAAK;AAClD;;;AC3FA,SAAgB,6BAA6B,cAAyD;CACpG,OAAO;EACL,GAAG,aAAa;EAChB,IAAI,aAAa;EACjB,QAAQ,aAAa;EACrB,MAAM,aAAa;EACnB,MAAM,aAAa;EACnB,UAAU,aAAa;EACvB,QAAQ,aAAa;EACrB,GAAI,aAAa,kBAAkB,aAAa,iBAAiB,IAC7D,EAAE,gBAAgB,aAAa,eAAe,IAC9C,CAAC;CACP;AACF;AAEA,SAAgB,4BAA4B,SAAsC;CAKhF,OAJgB,OAAO,QAAQ,QAAQ,QAAQ,CAAC,CAC7C,MAAM,CAAC,IAAI,CAAC,OAAO,EAAE,cAAc,CAAC,CAAC,CAAC,CACtC,KAAK,CAAC,QAAQ,WAAW,GAAG,OAAO,IAAI,OAAO,CAAC,CAC/C,KAAK,IACK,KAAK;AACpB;AAEA,MAAM,gBAAwC;CAAC;CAAO;CAAU;CAAQ;AAAQ;AAEhF,SAAS,gBAAgB,YAA6F;CACpH,KAAK,IAAI,IAAI,cAAc,SAAS,GAAG,KAAK,GAAG,KAAK,GAAG;EACrD,MAAM,WAAW,cAAc;EAC/B,IAAI,aAAa,WAAW,aAAa,KAAK,GAAG,OAAO;CAC1D;AAEF;AAEA,SAAgB,2BAA2B,cAA8D;CACvG,OAAO;EACL,QAAQ;EACR,UAAU,aAAa;EACvB,QAAQ,aAAa;EACrB,MAAM,aAAa;EACnB,UAAU,aAAa;EACvB,QAAQ,aAAa;EACrB,GAAI,aAAa,kBAAkB,aAAa,iBAAiB,IAC7D,EAAE,gBAAgB,aAAa,eAAe,IAC9C,CAAC;EACL,GAAI,aAAa,cAAc,EAAE,aAAa,aAAa,YAAY,YAAY,EAAE,IAAI,CAAC;EAC1F,GAAI,aAAa,SAAS,EAAE,QAAQ,aAAa,OAAO,YAAY,EAAE,IAAI,CAAC;CAC7E;AACF;AAEA,SAAgB,kCAAkC,SAAiE;CACjH,MAAM,WAAW,gBAAgB,QAAQ,UAAU;CACnD,OAAO;EACL,QAAQ;EACR,SAAS,QAAQ;EACjB,QAAQ,OAAO,QAAQ,QAAQ,QAAQ,CAAC,CACrC,MAAM,CAAC,IAAI,CAAC,OAAO,EAAE,cAAc,CAAC,CAAC,CAAC,CACtC,KAAK,CAAC,QAAQ,YAAY;GAAE;GAAQ;EAAM,EAAE;EAC/C,YAAY,QAAQ;EACpB,iBAAiB,QAAQ;EACzB,GAAI,WAAW,EAAE,SAAS,IAAI,CAAC;CACjC;AACF;AAEA,SAAgB,yBAAyB,cAAsD;CAC7F,OAAOW,gBAAAA,aAAa;EAClB,MAAM;EACN,SAAS;EACT,UAAU,aAAa;EACvB,YAAY,6BAA6B,YAAY;EACrD,UAAU;GAAE,GAAG,aAAa;GAAU,cAAc,2BAA2B,YAAY;EAAE;CAC/F,CAAC;AACH;AAEA,SAAgB,gCAAgC,SAAkD;CAChG,MAAM,eAAe,kCAAkC,OAAO;CAC9D,OAAOA,gBAAAA,aAAa;EAClB,MAAM;EACN,SAAS;EACT,UAAU,4BAA4B,OAAO;EAC7C,YAAY;GACV,SAAS,QAAQ;GACjB,GAAI,aAAa,WAAW,EAAE,UAAU,aAAa,SAAS,IAAI,CAAC;EACrE;EACA,UAAU;GAAE;GAAc,qBAAqB;GAAS,iBAAiB,QAAQ;EAAgB;CACnG,CAAC;AACH;AAEA,SAAgB,uBAAuB,eAA0D;CAC/F,MAAM,uBAAuB,cAAc,QAAO,iBAAgB,aAAa,WAAW,SAAS;CACnG,MAAM,QAAQ,qBAAqB,MAAM,cAAc;CACvD,OAAO,qBAAqB,QACzB,SAAS,iBAAiB;EACzB,QAAQ,WAAW;EACnB,QAAQ,SAAS,aAAa,WAAW,QAAQ,SAAS,aAAa,WAAW,KAAK;EACvF,QAAQ,WAAW,aAAa,aAAa,QAAQ,WAAW,aAAa,aAAa,KAAK;EAC/F,QAAQ,gBAAgB,KAAK,aAAa,EAAE;EAC5C,OAAO;CACT,GACA;EACE,UAAU,OAAO,YAAY;EAC7B,YAAY,OAAO;EACnB,SAAS,OAAO;EAChB,SAAS;EACT,UAAU,CAAC;EACX,YAAY,CAAC;EACb,iBAAiB,CAAC;CACpB,CACF;AACF;;;ACzFA,MAAM,gBAAgB,UAA4B,iBAAiB,QAAQ,MAAM,UAAU,OAAO,KAAK;AAEvG,MAAM,gBAAgB,QAA4B,QAChD,QAAQ,OAAO,aAAa,OAAO,UAAU,QAAQ,KAAK,IAAI,QAAQ,CAAC;AAEzE,MAAM,mBAAmE;CACvE,QAAQ;CACR,MAAM;CACN,QAAQ;CACR,KAAK;AACP;AAeA,MAAM,2BAA2B,GAAoB,MACnD,iBAAiB,EAAE,YAAY,iBAAiB,EAAE,aAAa,EAAE,UAAU,QAAQ,IAAI,EAAE,UAAU,QAAQ;AAE7G,MAAM,sBAAsB,YAC1B,QAAQ,QACL,SAAS,WAAY,iBAAiB,OAAO,YAAY,iBAAiB,WAAW,OAAO,WAAW,SACxG,KACF;AAEF,MAAM,wBAAwB,YAC5B,QAAQ,QACL,UAAU,WAAY,OAAO,UAAU,QAAQ,IAAI,SAAS,QAAQ,IAAI,OAAO,YAAY,UAC5F,QAAQ,EAAE,CAAE,SACd;AAEF,MAAM,YAAY,WAAmD;CACnE,IAAI,CAAC,OAAO,WAAW,CAAC,OAAO,cAAc,CAAC,OAAO,UAAU,OAAO,KAAA;CACtE,OAAO;EAAC,OAAO;EAAS,OAAO;EAAY,OAAO;CAAQ,CAAC,CAAC,KAAK,IAAI;AACvE;AAEA,eAAe,sBAAsB,EACnC,SACA,QACA,KACA,SAMC;CACD,MAAM,QAAQ,mBAAmB;EAC/B,IAAI,OAAO;EACX,UAAU,OAAO;EACjB,GAAG,6BAA6B,MAAM;EACtC,uBAAuB;EACvB,mBAAmB,aAAa,KAAK;CACvC,CAAC;AACH;AAEA,eAAe,uBAAuB,EACpC,QACA,SACA,QACA,KACA,oBAO6E;CAC7E,MAAM,UAAU,MAAM,QAAQ,gBAAgB;EAAE,UAAU,OAAO;EAAU,IAAI,OAAO;CAAG,CAAC;CAC1F,IAAI,CAAC,WAAW,QAAQ,WAAW,aAAa,QAAQ,mBAAmB,OAAO;CAClF,IAAI,CAAC,QAAQ,SAAS,MAAM,IAAI,MAAM,gBAAgB,QAAQ,GAAG,oBAAoB;CACrF,IAAI,CAAC,QAAQ,YAAY,MAAM,IAAI,MAAM,gBAAgB,QAAQ,GAAG,uBAAuB;CAE3F,MAAM,QAAS,MAAM,OAAO,aAAa,QAAQ,OAAgB;CACjE,IAAI,QAAQ,aAAa,UAAU,QAAQ,iBAEvC;OAAA,oBACA,yBAAyB,eACvB;GAAE,YAAY,QAAQ;GAAY,UAAU,QAAQ;EAAS,GAC7D,MAAM,YAAY,CACpB,OACkB,UAAU,OAAO;CAAA;CAGvC,MAAM,SAAS,yBAAyB;EACtC,GAAG;EACH,QAAQ;EACR,aAAa;EACb,uBAAuB;CACzB,CAAC;CACD,MAAM,SAAiC;EAAE,YAAY,QAAQ;EAAY,UAAU,QAAQ;CAAS;CACpG,MAAM,SAAS,MAAM,WAAW,QAAQ,MAAM;CAI9C,MAAM,OAAO;CACb,MAAM,OAAO;CAQb,OAAO;EAAE,QAAQ,MAPK,QAAQ,mBAAmB;GAC/C,IAAI,QAAQ;GACZ,UAAU,QAAQ;GAClB,QAAQ;GACR,mBAAmB,OAAO,OAAO;GACjC,uBAAuB;EACzB,CAAC;EACyB,QAAQ,OAAO;CAAO;AAClD;AAEA,eAAe,wBAAwB,EACrC,QACA,SACA,SACA,OAMyE;CACzE,MAAM,QAAQ,QAAQ;CACtB,IAAI,CAAC,OAAO,SAAS,MAAM,IAAI,MAAM,yCAAyC;CAC9E,IAAI,CAAC,MAAM,YAAY,MAAM,IAAI,MAAM,4CAA4C;CAEnF,MAAM,QAAQ,MAAM,OAAO,aAAa,MAAM,OAAgB;CAE9D,MAAM,SAAS,gCADC,uBAAuB,OACc,CAAC;CACtD,MAAM,SAAiC,QAAQ,OAAM,WAAU,OAAO,aAAa,KAAK,IACpF;EAAE,YAAY,MAAM;EAAY,UAAU,MAAM;EAAU,QAAQ,EAAE,UAAU,UAAU;CAAE,IAC1F;EAAE,YAAY,MAAM;EAAY,UAAU,MAAM;CAAS;CAC7D,MAAM,SAAU,MAAoC,WAAW,QAAQ,MAAM;CAG7E,MAAM,OAAO;CACb,MAAM,OAAO;CAEb,MAAM,iBAAuC,CAAC;CAC9C,KAAK,MAAM,UAAU,SACnB,eAAe,KACb,MAAM,QAAQ,mBAAmB;EAC/B,IAAI,OAAO;EACX,UAAU,OAAO;EACjB,WAAW;EACX,iBAAiB,OAAO,OAAO;EAC/B,uBAAuB;CACzB,CAAC,CACH;CAEF,OAAO;EAAE,SAAS;EAAgB,QAAQ,OAAO;CAAO;AAC1D;AAEA,eAAe,oBAAoB,EACjC,QACA,SAI2C;CAC3C,MAAM,QAAS,MAAM,OAAO,aAAa,MAAM,OAAgB;CAC/D,OAAO,yBAAyB,eAC9B;EAAE,YAAY,MAAM;EAAY,UAAU,MAAM;CAAS,GACzD,MAAM,YAAY,CACpB;AACF;AAEA,eAAsB,yBAAyB,EAC7C,QACA,SACA,sBAAM,IAAI,KAAK,GACf,QAAQ,OACiE;CACzE,MAAM,MAAM,MAAM,QAAQ,qBAAqB;EAAE;EAAK;CAAM,CAAC;CAC7D,MAAM,YAAkC,CAAC;CACzC,MAAM,SAA+D,CAAC;CACtE,MAAM,UAAgC,CAAC;CACvC,MAAM,yBAAS,IAAI,IAA8B;CACjD,MAAM,sBAA4C,CAAC;CAEnD,KAAK,MAAM,UAAU,KAAK;EACxB,MAAM,MAAM,SAAS,MAAM;EAC3B,IAAI,CAAC,KAAK;GACR,IAAI,aAAa,QAAQ,GAAG,GAAG;IAC7B,MAAM,wBAAQ,IAAI,MAChB,gBAAgB,OAAO,GAAG,gEAC5B;IACA,MAAM,sBAAsB;KAAE;KAAS;KAAQ;KAAK;IAAM,CAAC;IAC3D,OAAO,KAAK;KAAE;KAAQ,OAAO,MAAM;IAAQ,CAAC;GAC9C,OACE,oBAAoB,KAAK,MAAM;GAEjC;EACF;EAEA,MAAM,QAAQ,OAAO,IAAI,GAAG,KAAK;GAC/B;GACA,SAAS,OAAO;GAChB,YAAY,OAAO;GACnB,UAAU,OAAO;GACjB,gBAAgB,CAAC;GACjB,mBAAmB,CAAC;EACtB;EACA,IAAI,aAAa,QAAQ,GAAG,GAC1B,MAAM,eAAe,KAAK,MAAM;OAEhC,MAAM,kBAAkB,KAAK,MAAM;EAErC,OAAO,IAAI,KAAK,KAAK;CACvB;CAEA,KAAK,MAAM,SAAS,OAAO,OAAO,GAAG;EACnC,MAAM,UAAU,CAAC,GAAG,MAAM,gBAAgB,GAAG,MAAM,iBAAiB;EACpE,IAAI;EACJ,IAAI;GACF,mBAAmB,MAAM,oBAAoB;IAAE;IAAQ;GAAM,CAAC;EAChE,SAAS,OAAO;GACd,KAAK,MAAM,UAAU,SAAS;IAC5B,MAAM,sBAAsB;KAAE;KAAS;KAAQ;KAAK;IAAM,CAAC;IAC3D,OAAO,KAAK;KAAE;KAAQ,OAAO,aAAa,KAAK;IAAE,CAAC;GACpD;GACA;EACF;EAEA,MAAM,QAA2B,MAAM,kBAAkB,KAAI,YAAW;GACtE,MAAM;GACN;GACA,UAAU,OAAO;GACjB,WAAW,OAAO;EACpB,EAAE;EACF,IAAI,MAAM,eAAe,SAAS,GAChC,MAAM,KAAK;GACT,MAAM;GACN,SAAS,MAAM;GACf,UAAU,mBAAmB,MAAM,cAAc;GACjD,WAAW,qBAAqB,MAAM,cAAc;EACtD,CAAC;EAEH,MAAM,KAAK,uBAAuB;EAElC,KAAK,MAAM,QAAQ,OAAO;GACxB,IAAI,KAAK,SAAS,WAAW;IAC3B,IAAI;KACF,MAAM,SAAS,MAAM,wBAAwB;MAAE;MAAQ;MAAS,SAAS,KAAK;MAAS;KAAI,CAAC;KAC5F,UAAU,KAAK,GAAG,OAAO,OAAO;KAChC,QAAQ,KAAK,OAAO,MAAM;IAC5B,SAAS,OAAO;KACd,KAAK,MAAM,UAAU,KAAK,SAAS;MACjC,MAAM,sBAAsB;OAAE;OAAS;OAAQ;OAAK;MAAM,CAAC;MAC3D,OAAO,KAAK;OAAE;OAAQ,OAAO,aAAa,KAAK;MAAE,CAAC;KACpD;IACF;IACA;GACF;GAEA,IAAI;IACF,MAAM,SAAS,MAAM,uBAAuB;KAAE;KAAQ;KAAS,QAAQ,KAAK;KAAQ;KAAK;IAAiB,CAAC;IAC3G,IAAI,CAAC,QAAQ;IACb,UAAU,KAAK,OAAO,MAAM;IAC5B,QAAQ,KAAK,OAAO,MAAM;GAC5B,SAAS,OAAO;IACd,MAAM,sBAAsB;KAAE;KAAS,QAAQ,KAAK;KAAQ;KAAK;IAAM,CAAC;IACxE,OAAO,KAAK;KAAE,QAAQ,KAAK;KAAQ,OAAO,aAAa,KAAK;IAAE,CAAC;GACjE;EACF;CACF;CAEA,oBAAoB,MAAM,GAAG,MAC3B,wBACE;EAAE,MAAM;EAAc,QAAQ;EAAG,UAAU,EAAE;EAAU,WAAW,EAAE;CAAU,GAC9E;EAAE,MAAM;EAAc,QAAQ;EAAG,UAAU,EAAE;EAAU,WAAW,EAAE;CAAU,CAChF,CACF;CACA,KAAK,MAAM,UAAU,qBACnB,IAAI;EACF,MAAM,SAAS,MAAM,uBAAuB;GAAE;GAAQ;GAAS;GAAQ;EAAI,CAAC;EAC5E,IAAI,CAAC,QAAQ;EACb,UAAU,KAAK,OAAO,MAAM;EAC5B,QAAQ,KAAK,OAAO,MAAM;CAC5B,SAAS,OAAO;EACd,MAAM,sBAAsB;GAAE;GAAS;GAAQ;GAAK;EAAM,CAAC;EAC3D,OAAO,KAAK;GAAE;GAAQ,OAAO,aAAa,KAAK;EAAE,CAAC;CACpD;CAGF,OAAO;EAAE;EAAW;EAAQ;CAAQ;AACtC;;;ACtTA,IAAsB,uBAAtB,cAAmDC,aAAAA,cAAc;CAC/D,cAAc;EACZ,MAAM;GAAE,WAAW;GAAW,MAAM;EAAgB,CAAC;CACvD;AAOF;AAEA,MAAM,aAAa,UAAkB,QAAQ,IAAI,KAAK,KAAK,IAAI,KAAA;AAC/D,MAAM,mBAAmB,UAAkB,OAAe,GAAG,SAAS,IAAI;AAC1E,MAAM,cAAiB,UACrB,UAAU,KAAA,IAAY,KAAA,IAAY,gBAAgB,KAAK;AAEzD,MAAM,eAAe,YAAoD;CACvE,GAAG;CACH,WAAW,IAAI,KAAK,OAAO,SAAS;CACpC,WAAW,IAAI,KAAK,OAAO,SAAS;CACpC,aAAa,UAAU,OAAO,WAAW;CACzC,QAAQ,UAAU,OAAO,MAAM;CAC/B,aAAa,UAAU,OAAO,WAAW;CACzC,YAAY,UAAU,OAAO,UAAU;CACvC,aAAa,UAAU,OAAO,WAAW;CACzC,WAAW,UAAU,OAAO,SAAS;CACrC,WAAW,UAAU,OAAO,SAAS;CACrC,uBAAuB,UAAU,OAAO,qBAAqB;CAC7D,SAAS,WAAW,OAAO,OAAO;CAClC,YAAY,WAAW,OAAO,UAAU;CACxC,UAAU,WAAW,OAAO,QAAQ;AACtC;AAEA,MAAM,mBAAmB,QAA4B,QAAc;CACjE,IAAI,WAAW,aAAa,OAAO,EAAE,aAAa,IAAI;CACtD,IAAI,WAAW,QAAQ,OAAO,EAAE,QAAQ,IAAI;CAC5C,IAAI,WAAW,aAAa,OAAO,EAAE,aAAa,IAAI;CACtD,IAAI,WAAW,YAAY,OAAO,EAAE,YAAY,IAAI;CACpD,IAAI,WAAW,aAAa,OAAO,EAAE,aAAa,IAAI;CACtD,OAAO,CAAC;AACV;AAEA,MAAM,gBAAkC,OAAU,WAAqB;CACrE,IAAI,CAAC,QAAQ,OAAO;CACpB,OAAO,MAAM,QAAQ,MAAM,IAAI,OAAO,SAAS,KAAK,IAAI,UAAU;AACpE;AAEA,MAAM,WAAW,WAAuC;CACtD,MAAM,YAAY,OAAO,WAAW,QAAQ;CAC5C,MAAM,YAAY,OAAO,WAAW,QAAQ;CAC5C,IAAI,cAAc,KAAA,KAAa,cAAc,KAAA,GAAW,OAAO,KAAK,IAAI,WAAW,SAAS;CAC5F,OAAO,aAAa,aAAa,OAAO;AAC1C;AAEA,IAAa,+BAAb,cAAkD,qBAAqB;CACrE,iCAAiB,IAAI,IAAgC;CAErD,MAAM,mBAAmB,OAA6D;EACpF,MAAM,WAAW,KAAK,gBAAgB,KAAK;EAC3C,IAAI,UAAU;GACZ,MAAM,sBAAM,IAAI,KAAK;GACrB,MAAM,OAA2B;IAC/B,GAAG;IACH,SAAS,MAAM;IACf,SAAS,WAAW,MAAM,WAAW,SAAS,OAAO;IACrD,UAAU,MAAM,YAAY,SAAS;IACrC,YAAY,MAAM,aACd;KAAE,GAAG,WAAW,SAAS,UAAU;KAAG,GAAG,WAAW,MAAM,UAAU;IAAE,IACtE,WAAW,SAAS,UAAU;IAClC,WAAW;IACX,WAAW,MAAM,aAAa,SAAS;IACvC,WAAW,MAAM,aAAa,SAAS;IACvC,gBAAgB,MAAM,kBAAkB,SAAS;IACjD,iBAAiB,SAAS,kBAAkB,KAAK;IACjD,UAAU,MAAM,WACZ;KAAE,GAAG,WAAW,SAAS,QAAQ;KAAG,GAAG,WAAW,MAAM,QAAQ;IAAE,IAClE,WAAW,SAAS,QAAQ;GAClC;GACA,KAAKC,eAAe,IAAI,gBAAgB,KAAK,UAAU,KAAK,EAAE,GAAG,IAAI;GACrE,OAAO,YAAY,IAAI;EACzB;EAEA,MAAM,MAAM,MAAM,6BAAa,IAAI,KAAK;EACxC,MAAM,SAA6B;GACjC,IAAI,MAAM,OAAA,GAAA,OAAA,WAAA,CAAiB;GAC3B,UAAU,MAAM;GAChB,QAAQ,MAAM;GACd,MAAM,MAAM;GACZ,UAAU,MAAM,YAAY;GAC5B,QAAQ;GACR,SAAS,MAAM;GACf,SAAS,WAAW,MAAM,OAAO;GACjC,YAAY,MAAM;GAClB,SAAS,MAAM;GACf,UAAU,MAAM;GAChB,WAAW,MAAM;GACjB,aAAa,MAAM;GACnB,gBAAgB;GAChB,YAAY,WAAW,MAAM,UAAU;GACvC,WAAW;GACX,WAAW;GACX,WAAW,MAAM;GACjB,WAAW,MAAM;GACjB,gBAAgB,MAAM;GACtB,kBAAkB;GAClB,UAAU,WAAW,MAAM,QAAQ;EACrC;EACA,KAAKA,eAAe,IAAI,gBAAgB,OAAO,UAAU,OAAO,EAAE,GAAG,MAAM;EAC3E,OAAO,YAAY,MAAM;CAC3B;CAEA,MAAM,kBAAkB,OAA8D;EACpF,MAAM,SAAS,MAAM,QAAQ,YAAY;EACzC,MAAM,UAAU,CAAC,GAAG,KAAKA,eAAe,OAAO,CAAC,CAAC,CAC9C,QAAO,WAAU,OAAO,aAAa,MAAM,QAAQ,CAAC,CACpD,QAAO,WAAU,aAAa,OAAO,QAAQ,MAAM,MAAM,CAAC,CAAC,CAC3D,QAAO,WAAU,aAAa,OAAO,UAAU,MAAM,QAAQ,CAAC,CAAC,CAC/D,QAAO,WAAU,CAAC,MAAM,UAAU,OAAO,WAAW,MAAM,MAAM,CAAC,CACjE,QAAO,WAAU,CAAC,MAAM,cAAc,OAAO,eAAe,MAAM,UAAU,CAAC,CAC7E,QAAO,WAAU,CAAC,MAAM,WAAW,OAAO,YAAY,MAAM,OAAO,CAAC,CACpE,QACC,WACE,CAAC,UACD,OAAO,QAAQ,YAAY,CAAC,CAAC,SAAS,MAAM,KAC5C,OAAO,KAAK,YAAY,CAAC,CAAC,SAAS,MAAM,KACzC,OAAO,OAAO,YAAY,CAAC,CAAC,SAAS,MAAM,CAC/C,CAAC,CACA,MAAM,GAAG,MAAM,EAAE,UAAU,QAAQ,IAAI,EAAE,UAAU,QAAQ,CAAC;EAC/D,OAAO,QAAQ,MAAM,GAAG,MAAM,SAAS,QAAQ,MAAM,CAAC,CAAC,IAAI,WAAW;CACxE;CAEA,MAAM,qBAAqB,OAAiE;EAC1F,MAAM,MAAM,MAAM,IAAI,QAAQ;EAC9B,MAAM,UAAU,CAAC,GAAG,KAAKA,eAAe,OAAO,CAAC,CAAC,CAC9C,QAAO,WAAU,OAAO,WAAW,SAAS,CAAC,CAC7C,QAAO,WAAU,CAAC,MAAM,WAAW,OAAO,YAAY,MAAM,OAAO,CAAC,CACpE,QAAO,WAAU,CAAC,MAAM,cAAc,OAAO,eAAe,MAAM,UAAU,CAAC,CAC7E,QAAO,WAAU,QAAQ,MAAM,KAAK,GAAG,CAAC,CACxC,MAAM,GAAG,MAAM,QAAQ,CAAC,IAAI,QAAQ,CAAC,KAAK,EAAE,UAAU,QAAQ,IAAI,EAAE,UAAU,QAAQ,CAAC;EAC1F,OAAO,QAAQ,MAAM,GAAG,MAAM,SAAS,QAAQ,MAAM,CAAC,CAAC,IAAI,WAAW;CACxE;CAEA,MAAM,gBAAgB,OAA6E;EACjG,MAAM,SAAS,KAAKA,eAAe,IAAI,gBAAgB,MAAM,UAAU,MAAM,EAAE,CAAC;EAChF,IAAI,CAAC,QAAQ,OAAO;EACpB,OAAO,YAAY,MAAM;CAC3B;CAEA,MAAM,mBAAmB,OAA6D;EACpF,MAAM,WAAW,KAAKA,eAAe,IAAI,gBAAgB,MAAM,UAAU,MAAM,EAAE,CAAC;EAClF,IAAI,CAAC,UACH,MAAM,IAAI,MAAM,gBAAgB,MAAM,GAAG,4BAA4B,MAAM,UAAU;EAEvF,MAAM,sBAAM,IAAI,KAAK;EACrB,MAAM,OAA2B;GAC/B,GAAG;GACH,GAAI,MAAM,SAAS;IAAE,QAAQ,MAAM;IAAQ,GAAG,gBAAgB,MAAM,QAAQ,GAAG;GAAE,IAAI,CAAC;GACtF,GAAI,MAAM,YAAY,KAAA,IAAY,EAAE,SAAS,MAAM,QAAQ,IAAI,CAAC;GAChE,GAAI,MAAM,YAAY,KAAA,IAAY,EAAE,SAAS,WAAW,MAAM,OAAO,EAAE,IAAI,CAAC;GAC5E,GAAI,MAAM,eAAe,KAAA,IAAY,EAAE,YAAY,WAAW,MAAM,UAAU,EAAE,IAAI,CAAC;GACrF,GAAI,MAAM,aAAa,KAAA,IAAY,EAAE,UAAU,WAAW,MAAM,QAAQ,EAAE,IAAI,CAAC;GAC/E,GAAI,MAAM,cAAc,KAAA,IAAY,EAAE,WAAW,MAAM,aAAa,KAAA,EAAU,IAAI,CAAC;GACnF,GAAI,MAAM,cAAc,KAAA,IAAY,EAAE,WAAW,MAAM,aAAa,KAAA,EAAU,IAAI,CAAC;GACnF,GAAI,MAAM,mBAAmB,KAAA,IAAY,EAAE,gBAAgB,MAAM,eAAe,IAAI,CAAC;GACrF,GAAI,MAAM,qBAAqB,KAAA,IAAY,EAAE,kBAAkB,MAAM,iBAAiB,IAAI,CAAC;GAC3F,GAAI,MAAM,0BAA0B,KAAA,IAAY,EAAE,uBAAuB,MAAM,sBAAsB,IAAI,CAAC;GAC1G,GAAI,MAAM,sBAAsB,KAAA,IAAY,EAAE,mBAAmB,MAAM,kBAAkB,IAAI,CAAC;GAC9F,GAAI,MAAM,sBAAsB,KAAA,IAAY,EAAE,mBAAmB,MAAM,kBAAkB,IAAI,CAAC;GAC9F,GAAI,MAAM,oBAAoB,KAAA,IAAY,EAAE,iBAAiB,MAAM,gBAAgB,IAAI,CAAC;GACxF,WAAW;EACb;EACA,KAAKA,eAAe,IAAI,gBAAgB,KAAK,UAAU,KAAK,EAAE,GAAG,IAAI;EACrE,OAAO,YAAY,IAAI;CACzB;CAEA,MAAM,sBAAqC;EACzC,KAAKA,eAAe,MAAM;CAC5B;CAEA,gBAAwB,OAAgE;EACtF,IAAI,CAAC,MAAM,aAAa,CAAC,MAAM,aAAa,OAAO,KAAA;EACnD,OAAO,CAAC,GAAG,KAAKA,eAAe,OAAO,CAAC,CAAC,CAAC,MAAK,WAAU;GACtD,IACE,OAAO,aAAa,MAAM,YAC1B,OAAO,WAAW,MAAM,UACxB,OAAO,SAAS,MAAM,QACtB,OAAO,WAAW,WAElB,OAAO;GACT,IAAI,OAAO,YAAY,MAAM,WAAW,OAAO,eAAe,MAAM,YAAY,OAAO;GACvF,OAAO,QACJ,MAAM,aAAa,OAAO,cAAc,MAAM,aAC9C,MAAM,eAAe,OAAO,gBAAgB,MAAM,WACrD;EACF,CAAC;CACH;AACF"}