{"version":3,"file":"channel.cjs","names":["namespaceKey","isRootNamespace","inferChannel","openProjectionSubscription"],"sources":["../../../src/stream/projections/channel.ts"],"sourcesContent":["/**\n * Raw channel escape hatch.\n *\n * Subscribes to an arbitrary list of channels at an arbitrary\n * namespace and retains a bounded buffer of events. The subscription\n * resumes across serial runs, so the buffer keeps accumulating events\n * for the lifetime of the thread (use `extensionProjection` instead when\n * you only need the most-recent payload of a single `custom:<name>`\n * channel). Consumers that need assembly semantics should use\n * `messagesProjection`, `toolCallsProjection`, etc. instead; this one is\n * for inspection, custom reducers, or niche use-cases.\n */\nimport type { Channel, Event } from \"@langchain/protocol\";\nimport { inferChannel } from \"../../client/stream/subscription.js\";\nimport type { ProjectionSpec, ProjectionRuntime } from \"../types.js\";\nimport { isRootNamespace, namespaceKey } from \"../namespace.js\";\nimport { openProjectionSubscription } from \"./runtime.js\";\n\n/** Max events retained per raw subscription. Older events are dropped. */\nconst DEFAULT_BUFFER = 4096;\n\nexport interface ChannelProjectionOptions {\n  /**\n   * Maximum number of events retained in the projection snapshot.\n   * Defaults to 4096 so replay-backed discovery hooks can tolerate\n   * bursty token/media streams without dropping early lifecycle events.\n   */\n  bufferSize?: number;\n  /**\n   * Whether to open a real subscription and receive replayed history.\n   * Defaults to true. Set false for live-only root-bus inspection when\n   * replay is unnecessary.\n   */\n  replay?: boolean;\n}\n\nexport function channelProjection(\n  channels: readonly Channel[],\n  namespace: readonly string[],\n  options: ChannelProjectionOptions = {}\n): ProjectionSpec<Event[]> {\n  const chs = [...channels].sort();\n  const ns = [...namespace];\n  const bufferSize = options.bufferSize ?? DEFAULT_BUFFER;\n  const replay = options.replay ?? true;\n  const key = `channel|${bufferSize}|${replay ? \"replay\" : \"live\"}|${chs.join(\",\")}|${namespaceKey(ns)}`;\n\n  return {\n    key,\n    namespace: ns,\n    initial: [],\n    open({ thread, store, rootBus }): ProjectionRuntime {\n      // If this projection is scoped to the root namespace AND every\n      // requested channel is already covered by the controller's root\n      // pump, attach to the shared fan-out instead of opening a\n      // second server subscription. This is the common case for\n      // lightweight event-trace / debug panels.\n      const covered =\n        !replay &&\n        isRootNamespace(ns) &&\n        chs.every((c) => rootBus.channels.includes(c));\n\n      if (covered) {\n        const requestedSet = new Set(chs as Channel[]);\n        // Classify via `inferChannel` so namespaced methods (e.g.\n        // `input.requested` → `input`) and named custom events\n        // (`custom` + `data.name` → `custom:<name>`) match the same\n        // way as the slow path's `matchesSubscription`. Comparing\n        // `event.method` to channel names would silently drop `input`.\n        const matches = (event: Event): boolean => {\n          const channel = inferChannel(event);\n          if (channel === undefined) return false;\n          return (\n            requestedSet.has(channel) ||\n            (channel.startsWith(\"custom:\") && requestedSet.has(\"custom\"))\n          );\n        };\n        const push = (event: Event): void => {\n          if (!matches(event)) return;\n          const current = store.getSnapshot();\n          const next =\n            current.length >= bufferSize\n              ? [...current.slice(current.length - bufferSize + 1), event]\n              : [...current, event];\n          store.setValue(next);\n        };\n        const unsubscribe = rootBus.subscribe(push);\n        return {\n          dispose() {\n            unsubscribe();\n          },\n        };\n      }\n\n      return openProjectionSubscription({\n        thread,\n        channels: chs as Channel[],\n        namespace: ns,\n        // Keep the buffer alive across serial runs. Some transports pause a\n        // subscription on each run's terminal lifecycle event; without\n        // resuming, the channel would go silent after the first run and a\n        // second prompt on the same thread would emit no further events.\n        resumeOnPause: true,\n        onEvent(event) {\n          const current = store.getSnapshot();\n          const next =\n            current.length >= bufferSize\n              ? [...current.slice(current.length - bufferSize + 1), event]\n              : [...current, event];\n          store.setValue(next);\n        },\n      });\n    },\n  };\n}\n"],"mappings":";;;;;AAmBA,MAAM,iBAAiB;AAiBvB,SAAgB,kBACd,UACA,WACA,UAAoC,EAAE,EACb;CACzB,MAAM,MAAM,CAAC,GAAG,SAAS,CAAC,MAAM;CAChC,MAAM,KAAK,CAAC,GAAG,UAAU;CACzB,MAAM,aAAa,QAAQ,cAAc;CACzC,MAAM,SAAS,QAAQ,UAAU;AAGjC,QAAO;EACL,KAHU,WAAW,WAAW,GAAG,SAAS,WAAW,OAAO,GAAG,IAAI,KAAK,IAAI,CAAC,GAAGA,kBAAAA,aAAa,GAAG;EAIlG,WAAW;EACX,SAAS,EAAE;EACX,KAAK,EAAE,QAAQ,OAAO,WAA8B;AAWlD,OAJE,CAAC,UACDC,kBAAAA,gBAAgB,GAAG,IACnB,IAAI,OAAO,MAAM,QAAQ,SAAS,SAAS,EAAE,CAAC,EAEnC;IACX,MAAM,eAAe,IAAI,IAAI,IAAiB;IAM9C,MAAM,WAAW,UAA0B;KACzC,MAAM,UAAUC,qBAAAA,aAAa,MAAM;AACnC,SAAI,YAAY,KAAA,EAAW,QAAO;AAClC,YACE,aAAa,IAAI,QAAQ,IACxB,QAAQ,WAAW,UAAU,IAAI,aAAa,IAAI,SAAS;;IAGhE,MAAM,QAAQ,UAAuB;AACnC,SAAI,CAAC,QAAQ,MAAM,CAAE;KACrB,MAAM,UAAU,MAAM,aAAa;KACnC,MAAM,OACJ,QAAQ,UAAU,aACd,CAAC,GAAG,QAAQ,MAAM,QAAQ,SAAS,aAAa,EAAE,EAAE,MAAM,GAC1D,CAAC,GAAG,SAAS,MAAM;AACzB,WAAM,SAAS,KAAK;;IAEtB,MAAM,cAAc,QAAQ,UAAU,KAAK;AAC3C,WAAO,EACL,UAAU;AACR,kBAAa;OAEhB;;AAGH,UAAOC,gBAAAA,2BAA2B;IAChC;IACA,UAAU;IACV,WAAW;IAKX,eAAe;IACf,QAAQ,OAAO;KACb,MAAM,UAAU,MAAM,aAAa;KACnC,MAAM,OACJ,QAAQ,UAAU,aACd,CAAC,GAAG,QAAQ,MAAM,QAAQ,SAAS,aAAa,EAAE,EAAE,MAAM,GAC1D,CAAC,GAAG,SAAS,MAAM;AACzB,WAAM,SAAS,KAAK;;IAEvB,CAAC;;EAEL"}