{"version":3,"file":"dispatch_reasoner.cjs","names":[],"sources":["../../../src/batteries/orchestration/dispatch_reasoner.ts"],"sourcesContent":["/**\n * @module @nhtio/adk/batteries/orchestration/dispatch_reasoner\n */\n\nimport { Message, Tool } from '@nhtio/adk/common'\nimport { inMemoryMediaReader } from '@nhtio/adk/common'\nimport { DispatchRunner } from '@nhtio/adk/dispatch_runner'\nimport { isObject, isInstanceOf } from '../../lib/utils/guards'\nimport { InMemorySpoolStore } from '@nhtio/adk/batteries/storage/in_memory'\nimport { wrapInstruction, stripInstructionTags, validateReasonerOutput } from './reason'\nimport type { Schema } from '@nhtio/validation'\nimport type { ReasonerFn, EncodableValue } from './types'\nimport type { RawDispatchContext } from '@nhtio/adk/types'\nimport type { DispatchContext, DispatchExecutorFn } from '@nhtio/adk/types'\n\n/**\n * An `AbortController` pre-aborted when `signal` already is, so a cancelled reasoning request\n * does not start a dispatch. `RawDispatchContext` takes a controller rather than a signal.\n */\nconst signalController = (signal?: AbortSignal): AbortController => {\n  const controller = new AbortController()\n  if (signal?.aborted === true) controller.abort()\n  else signal?.addEventListener('abort', () => controller.abort(), { once: true })\n  return controller\n}\n\n/**\n * The persistence callbacks a standalone `DispatchContext` requires, as no-ops.\n *\n * A reasoning dispatch is a single question with a single structured answer: nothing it produces\n * outlives the call, so there is no store to write to and no history to fetch. Bytes are the one\n * exception — a large result needs somewhere to live for the duration — so those two callbacks\n * are backed by the supplied spool store. Kept internal deliberately: exporting this shape would\n * commit the battery to it.\n */\nconst noopPersistence = (\n  spoolStore: InMemorySpoolStore\n): Omit<RawDispatchContext, 'systemPrompt' | 'messages' | 'tools' | 'turnAbortController'> => ({\n  fetchMemories: async () => [],\n  fetchMessages: async () => [],\n  fetchThoughts: async () => [],\n  fetchToolCalls: async () => [],\n  fetchTools: async () => [],\n  fetchRetrievables: async () => [],\n  refreshStandingInstructions: async () => [],\n  storeStandingInstruction: async () => {},\n  mutateStandingInstruction: async () => {},\n  deleteStandingInstruction: async () => {},\n  storeMemory: async () => {},\n  mutateMemory: async () => {},\n  deleteMemory: async () => {},\n  storeRetrievable: async () => {},\n  mutateRetrievable: async () => {},\n  deleteRetrievable: async () => {},\n  storeMessage: async () => {},\n  mutateMessage: async () => {},\n  deleteMessage: async () => {},\n  storeThought: async () => {},\n  mutateThought: async () => {},\n  deleteThought: async () => {},\n  storeToolCall: async () => {},\n  mutateToolCall: async () => {},\n  deleteToolCall: async () => {},\n  storeMediaBytes: async (_ctx, _id, bytes) =>\n    inMemoryMediaReader(\n      isInstanceOf(bytes, 'Uint8Array', Uint8Array)\n        ? bytes\n        : new TextEncoder().encode(String(bytes))\n    ),\n  storeRetrievableBytes: (_ctx, id, bytes) => spoolStore.write(id, bytes as Uint8Array),\n})\n\n/**\n * Builds a reasoner that answers through a forced tool dispatch.\n *\n * A reasoning node must return a structured, schema-validated result, but\n * {@link DispatchRunner.dispatch} resolves to `void`. This helper bridges that\n * gap by wiring the node's output schema onto a forced tool's `inputSchema` and\n * capturing the arguments the model submits to that tool. Because the tool's\n * input schema is the node's output schema, validation rejects malformed\n * arguments before the handler ever runs, so the model cannot answer with\n * unstructured prose.\n *\n * The captured value is read from a closure after the dispatch resolves; the\n * void return value is irrelevant. If the tool is never called, the dispatch is\n * a halting failure and this function throws rather than fabricating or\n * returning a partial result.\n *\n * The boundary this draws is deliberate: the battery owns the forced-tool protocol and the retry\n * bound, and the CONSUMER owns the model. `DispatchRunner` takes an injected executor rather than\n * a model identifier, so a consumer wires whichever provider they already use and this helper\n * never grows credential handling.\n *\n * @param options - Configuration for the reasoner.\n * @param options.executor - The consumer's LLM call, invoked by the runner on every iteration.\n * @param options.spoolStore - Store used to spool large results. Defaults to a\n *   fresh {@link InMemorySpoolStore} so a large reasoning result always has\n *   somewhere to go.\n * @param options.toolName - Name of the forced capture tool. Defaults to\n *   `'submit_reasoning'`.\n * @returns A {@link ReasonerFn} that resolves to the validated reasoning\n *   result or throws if no valid result could be captured.\n */\nexport const createDispatchReasoner = (options: {\n  executor: DispatchExecutorFn\n  spoolStore?: InMemorySpoolStore\n  toolName?: string\n}): ReasonerFn => {\n  const spoolStore = options.spoolStore ?? new InMemorySpoolStore()\n  const toolName = options.toolName ?? 'submit_reasoning'\n\n  // Constructing Message and Tool here is a documented CONTRIBUTING §13\n  // exception, which is why this helper lives at its own subpath and not in the\n  // battery's environment-neutral barrel.\n  const reasoner: ReasonerFn = async (req) => {\n    const { prompt, outputSchema, maxAttempts, signal } = req\n    const schema: Schema = outputSchema\n\n    const corrective: string[] = []\n\n    for (let attempt = 0; attempt < maxAttempts; attempt += 1) {\n      let captured: Record<string, EncodableValue> | undefined\n\n      const tool = new Tool({\n        name: toolName,\n        description: 'Submit the structured reasoning result for this node.',\n        inputSchema: schema,\n        handler: async (args: unknown, ctx: DispatchContext) => {\n          if (isObject(args)) {\n            captured = args as Record<string, EncodableValue>\n          }\n          ctx.ack()\n          return 'Reasoning accepted.'\n        },\n      })\n\n      const body =\n        corrective.length > 0\n          ? wrapInstruction(prompt, `The previous attempt was rejected. ${corrective.join(' ')}`)\n          : prompt\n\n      const now = new Date()\n      const message = new Message({\n        id: `reason-${attempt}-${now.getTime()}`,\n        role: 'user',\n        content: body,\n        createdAt: now,\n        updatedAt: now,\n      })\n\n      await DispatchRunner.dispatch({\n        executor: options.executor,\n        raw: {\n          systemPrompt: 'Answer only by calling the submission tool.',\n          messages: [message],\n          tools: [tool],\n          turnAbortController: signalController(signal),\n          ...noopPersistence(spoolStore),\n        },\n      })\n\n      if (captured === undefined) {\n        corrective.push('The tool was never called; a structured result is required.')\n        continue\n      }\n\n      const validation = validateReasonerOutput(schema, captured)\n      if (validation === true) {\n        return captured\n      }\n\n      corrective.push(stripInstructionTags(validation))\n    }\n\n    throw new Error(\n      `Reasoning failed after ${maxAttempts} attempt(s): no valid structured result was captured.`\n    )\n  }\n\n  return reasoner\n}\n"],"mappings":";;;;;;;;;;;;;;;;;;AAmBA,IAAM,oBAAoB,WAA0C;CAClE,MAAM,aAAa,IAAI,gBAAgB;CACvC,IAAI,QAAQ,YAAY,MAAM,WAAW,MAAM;MAC1C,QAAQ,iBAAiB,eAAe,WAAW,MAAM,GAAG,EAAE,MAAM,KAAK,CAAC;CAC/E,OAAO;AACT;;;;;;;;;;AAWA,IAAM,mBACJ,gBAC6F;CAC7F,eAAe,YAAY,CAAC;CAC5B,eAAe,YAAY,CAAC;CAC5B,eAAe,YAAY,CAAC;CAC5B,gBAAgB,YAAY,CAAC;CAC7B,YAAY,YAAY,CAAC;CACzB,mBAAmB,YAAY,CAAC;CAChC,6BAA6B,YAAY,CAAC;CAC1C,0BAA0B,YAAY,CAAC;CACvC,2BAA2B,YAAY,CAAC;CACxC,2BAA2B,YAAY,CAAC;CACxC,aAAa,YAAY,CAAC;CAC1B,cAAc,YAAY,CAAC;CAC3B,cAAc,YAAY,CAAC;CAC3B,kBAAkB,YAAY,CAAC;CAC/B,mBAAmB,YAAY,CAAC;CAChC,mBAAmB,YAAY,CAAC;CAChC,cAAc,YAAY,CAAC;CAC3B,eAAe,YAAY,CAAC;CAC5B,eAAe,YAAY,CAAC;CAC5B,cAAc,YAAY,CAAC;CAC3B,eAAe,YAAY,CAAC;CAC5B,eAAe,YAAY,CAAC;CAC5B,eAAe,YAAY,CAAC;CAC5B,gBAAgB,YAAY,CAAC;CAC7B,gBAAgB,YAAY,CAAC;CAC7B,iBAAiB,OAAO,MAAM,KAAK,UACjC,eAAA,oBACE,eAAA,aAAa,OAAO,cAAc,UAAU,IACxC,QACA,IAAI,YAAY,EAAE,OAAO,OAAO,KAAK,CAAC,CAC5C;CACF,wBAAwB,MAAM,IAAI,UAAU,WAAW,MAAM,IAAI,KAAmB;AACtF;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;;AAiCA,IAAa,0BAA0B,YAIrB;CAChB,MAAM,aAAa,QAAQ,cAAc,IAAI,oCAAA,mBAAmB;CAChE,MAAM,WAAW,QAAQ,YAAY;CAKrC,MAAM,WAAuB,OAAO,QAAQ;EAC1C,MAAM,EAAE,QAAQ,cAAc,aAAa,WAAW;EACtD,MAAM,SAAiB;EAEvB,MAAM,aAAuB,CAAC;EAE9B,KAAK,IAAI,UAAU,GAAG,UAAU,aAAa,WAAW,GAAG;GACzD,IAAI;GAEJ,MAAM,OAAO,IAAI,yBAAA,KAAK;IACpB,MAAM;IACN,aAAa;IACb,aAAa;IACb,SAAS,OAAO,MAAe,QAAyB;KACtD,IAAI,eAAA,SAAS,IAAI,GACf,WAAW;KAEb,IAAI,IAAI;KACR,OAAO;IACT;GACF,CAAC;GAED,MAAM,OACJ,WAAW,SAAS,IAChB,uCAAA,gBAAgB,QAAQ,sCAAsC,WAAW,KAAK,GAAG,GAAG,IACpF;GAEN,MAAM,sBAAM,IAAI,KAAK;GACrB,MAAM,UAAU,IAAI,gBAAA,QAAQ;IAC1B,IAAI,UAAU,QAAQ,GAAG,IAAI,QAAQ;IACrC,MAAM;IACN,SAAS;IACT,WAAW;IACX,WAAW;GACb,CAAC;GAED,MAAM,wBAAA,eAAe,SAAS;IAC5B,UAAU,QAAQ;IAClB,KAAK;KACH,cAAc;KACd,UAAU,CAAC,OAAO;KAClB,OAAO,CAAC,IAAI;KACZ,qBAAqB,iBAAiB,MAAM;KAC5C,GAAG,gBAAgB,UAAU;IAC/B;GACF,CAAC;GAED,IAAI,aAAa,KAAA,GAAW;IAC1B,WAAW,KAAK,6DAA6D;IAC7E;GACF;GAEA,MAAM,aAAa,uCAAA,uBAAuB,QAAQ,QAAQ;GAC1D,IAAI,eAAe,MACjB,OAAO;GAGT,WAAW,KAAK,uCAAA,qBAAqB,UAAU,CAAC;EAClD;EAEA,MAAM,IAAI,MACR,0BAA0B,YAAY,sDACxC;CACF;CAEA,OAAO;AACT"}