import { ThreadStream } from "../client/stream/index.cjs";
import { StreamStore } from "./store.cjs";
import { RootSnapshot, StreamControllerOptions, StreamRespondAllOptions, StreamRespondOptions, StreamStopOptions, StreamSubmitOptions } from "./types.cjs";
import { ChannelRegistry } from "./channel-registry.cjs";
import { SubagentMap } from "./discovery/subagents.cjs";
import { SubgraphByNodeMap, SubgraphMap } from "./discovery/subgraphs.cjs";
import { MessageMetadata as MessageMetadata$1, MessageMetadataMap } from "./message-metadata-tracker.cjs";
import { SubmissionQueueEntry, SubmissionQueueSnapshot } from "./submit-coordinator.cjs";
import { Channel } from "@langchain/protocol";

//#region src/stream/controller.d.ts
/**
 * Channel set covered by the always-on root subscription. Exported so
 * projections (and transports) can reason about what the root pump
 * already delivers before opening additional server subscriptions.
 */
declare const ROOT_PUMP_CHANNELS: readonly Channel[];
/**
 * Coordinates one thread's protocol-v2 stream and exposes stable
 * observable projections for framework bindings.
 *
 * The controller owns the root subscription, lazily binds scoped
 * projections through {@link ChannelRegistry}, and normalizes protocol
 * events into class-message, tool-call, discovery, interrupt, and queue
 * stores.
 *
 * @typeParam StateType - Shape of the graph state exposed on `values`.
 * @typeParam InterruptType - Shape of protocol interrupt payloads.
 * @typeParam ConfigurableType - Shape of `config.configurable` accepted by submit.
 */
declare class StreamController<StateType extends object = Record<string, unknown>, InterruptType = unknown, ConfigurableType extends object = Record<string, unknown>> {
  #private;
  readonly rootStore: StreamStore<RootSnapshot<StateType, InterruptType>>;
  readonly subagentStore: StreamStore<SubagentMap>;
  readonly subgraphStore: StreamStore<SubgraphMap>;
  readonly subgraphByNodeStore: StreamStore<SubgraphByNodeMap>;
  readonly messageMetadataStore: StreamStore<MessageMetadataMap>;
  readonly queueStore: StreamStore<SubmissionQueueSnapshot<StateType>>;
  readonly registry: ChannelRegistry;
  /**
   * Create a controller around a LangGraph client and optional initial thread.
   *
   * @param options - Runtime configuration, client, thread id, and initial state.
   */
  constructor(options: StreamControllerOptions<StateType>);
  /**
   * Promise that settles the first time {@link hydrate} finishes on
   * the current thread. Resolves on a clean hydration, rejects when
   * the thread-state fetch errors. A fresh promise is installed on
   * every thread swap so `<Suspense>` wrappers re-suspend on
   * `switchThread`.
   */
  get hydrationPromise(): Promise<void>;
  /**
   * Fetch the checkpointed thread state and seed the root snapshot.
   * Re-calling with a different `threadId` swaps the underlying
   * {@link ThreadStream}, rewires the registry to the new thread, and
   * resets assemblers.
   *
   * @param threadId - Optional replacement thread id; `null` clears the active thread.
   */
  hydrate(threadId?: string | null): Promise<void>;
  /**
   * Lazily resolve a single subagent's execution namespace from
   * checkpoint history. Intended call site: the first scoped
   * `useMessages` / `useToolCalls` mount for a subagent whose namespace
   * is still the default `tools:<toolCallId>`. A fallback for the
   * hydrate-time bulk seed ({@link #seedDiscoveryFromHistory}) — most
   * subagents are already promoted by the time a panel opens.
   *
   * Skips ids already promoted past default-only (SSE replay or a prior
   * resolve). Concurrent calls for the same id share one `getHistory`
   * walk via {@link #namespaceResolves}.
   *
   * @param toolCallId - Parent `task` tool-call id (the subagent's discovery key).
   */
  resolveSubagentNamespace(toolCallId: string): Promise<void>;
  /**
   * Submit input to the active thread.
   *
   * To resume a pending interrupt, use {@link respond} instead.
   *
   * @param input - Input payload for a new run.
   * @param options - Per-run config, metadata, multitask behavior, and callbacks.
   */
  submit(input: unknown, options?: StreamSubmitOptions<StateType, ConfigurableType>): Promise<void>;
  /**
   * Disconnect the client from the active run and mark the controller
   * idle. By default also cancels the run server-side; pass
   * `{ cancel: false }` or call {@link disconnect} to keep the agent
   * running (join/rejoin).
   */
  stop(options?: StreamStopOptions): Promise<void>;
  /**
   * Disconnect the client without cancelling the run server-side.
   * Alias for `stop({ cancel: false })`.
   */
  disconnect(): Promise<void>;
  /**
   * Cancel a queued submission by id. Returns `true` when the entry
   * was found and removed, `false` otherwise.
   *
   * Today this only removes the entry from the client-side mirror —
   * once the server exposes queue cancel (roadmap A0.3) the
   * controller will additionally issue a cancel call against the
   * active transport.
   *
   * @param id - Client-side queue entry id to remove.
   */
  cancelQueued(id: string): Promise<boolean>;
  /**
   * Drop every queued submission. Server-side cancel arrives with A0.3.
   */
  clearQueue(): Promise<void>;
  /**
   * Respond to a single pending protocol interrupt.
   *
   * When `options.interruptId` is omitted, resolution walks
   * {@link ThreadStream.interrupts `thread.interrupts`} from newest to
   * oldest and picks the first entry whose `interruptId` has not already
   * been resolved by a prior `respond()` call. That entry may be at the
   * root (`namespace: []`) or inside a subgraph (non-empty `namespace`).
   * {@link RootSnapshot.interrupts `rootStore.interrupts`} /
   * framework `stream.interrupts` mirrors the same pending interrupts
   * (including nested namespaces) for UI rendering.
   *
   * Omitting `interruptId` is fine when exactly one interrupt is pending.
   * When several can be active (parallel subagents, fan-out, nested
   * graphs), pass an explicit `interruptId` so you resume the interrupt
   * the user acted on. `namespace` is resolved automatically from
   * `thread.interrupts` / `stream.interrupts` when omitted.
   *
   * To resume several interrupts pending at the same checkpoint in one
   * command, use {@link respondAll} — sequential single `respond()` calls
   * would not work, since the first resume starts a run, leaving the
   * others with no interrupted run to respond to.
   *
   * The server validates `namespace` against the pending interrupt. Root
   * interrupts use `namespace: []`. Subgraph interrupts carry the exact
   * tuple on each {@link Interrupt} / {@link InterruptPayload} entry.
   *
   * @param response - Payload sent back to the interrupted namespace.
   * @param options - Optional target (`interruptId` / `namespace`) and
   *   run-level `config` / `metadata` folded into the run that services
   *   the resume (model/user config, trigger source, test flags, …).
   *   Equivalent to the same fields on {@link StreamSubmitOptions}.
   *
   * @example Single pending interrupt (safe to omit a target)
   * ```ts
   * await controller.respond({ approved: true });
   * ```
   *
   * @example Carry run config / metadata onto the resume
   * ```ts
   * await controller.respond(
   *   { approved: true },
   *   { config: { configurable: { model: "gpt-4o" } }, metadata: { source: "ui" } },
   * );
   * ```
   *
   * @example Multiple / nested interrupts — target by id
   * ```tsx
   * for (const intr of stream.interrupts) {
   *   await stream.respond(decide(intr.value), { interruptId: intr.id! });
   * }
   * ```
   *
   * Each {@link InterruptPayload} on `thread.interrupts` and each
   * {@link Interrupt} on `stream.interrupts` mirrors an
   * `input.requested` event with its protocol `namespace`.
   */
  respond(response: unknown, options?: StreamRespondOptions<ConfigurableType>): Promise<void>;
  /**
   * Resume several pending interrupts at the same checkpoint in a single
   * command.
   *
   * Required when a run pauses on multiple interrupts simultaneously
   * (e.g. parallel tool-authorization prompts): a single
   * `Command({ resume })` carrying every interrupt's payload resumes them
   * together. Sequential {@link respond} calls would fail because the
   * first resume starts a run, leaving the rest with no interrupted run to
   * respond to.
   *
   * `responsesById` maps each pending `interruptId` to the payload sent
   * back to it, so different interrupts can receive different responses
   * (approve one, deny another). To send the *same* payload to several
   * interrupts, build the map with that value for each id, e.g.
   * `Object.fromEntries(ids.map((id) => [id, response]))`.
   *
   * The server resumes by `interruptId`, so namespaces are resolved
   * internally from `getThread()?.interrupts` and need not be supplied.
   *
   * @param responsesById - Map of pending `interruptId` to its response
   *   payload. Must contain at least one entry.
   * @param options - Optional run-level `config` / `metadata` folded into
   *   the single run that services the batched resume. Equivalent to the
   *   same fields on {@link StreamSubmitOptions}.
   *
   * @example Distinct payloads per interrupt
   * ```tsx
   * await stream.respondAll({
   *   [interruptA.id]: { approved: true },
   *   [interruptB.id]: { approved: false },
   * });
   * ```
   *
   * @example Same payload to every pending interrupt
   * ```tsx
   * await stream.respondAll(
   *   Object.fromEntries(stream.interrupts.map((i) => [i.id!, { approved: true }])),
   * );
   * ```
   */
  respondAll(responsesById: Record<string, unknown>, options?: StreamRespondAllOptions<ConfigurableType>): Promise<void>;
  /**
   * Dispose the active thread, subscriptions, registry entries, and listeners.
   */
  dispose(): Promise<void>;
  /**
   * StrictMode-safe lifecycle hook for framework bindings.
   *
   * React 18+ `StrictMode` intentionally mounts → unmounts → remounts
   * components in dev to surface effect-cleanup bugs. A naive
   * `useEffect(() => () => controller.dispose())` would permanently
   * tear the controller down on that first synthetic unmount, leaving
   * every subsequent `submit()` a silent no-op.
   *
   * Call {@link activate} from the bind site's effect and return the
   * result as the effect's cleanup. The controller uses deferred
   * disposal: a `release()` only schedules a dispose on the next
   * microtask, which is cancelled if another `activate()` arrives
   * before it fires (the normal StrictMode remount path).
   */
  activate(): () => void;
  /**
   * Returns the bound {@link ThreadStream}, if one exists. Prefer
   * {@link StreamController.rootStore} and selector projections for
   * UI work; use this for low-level protocol access.
   */
  getThread(): ThreadStream | undefined;
  /**
   * Listen for `ThreadStream` lifecycle (swap on thread-id change,
   * detach on dispose). The listener fires immediately with the
   * current thread (may be `undefined`).
   *
   * @param listener - Callback invoked immediately and on every thread swap.
   */
  subscribeThread(listener: (thread: ThreadStream | undefined) => void): () => void;
}
//#endregion
export { ROOT_PUMP_CHANNELS, StreamController };
//# sourceMappingURL=controller.d.cts.map