diff --git a/src/backend/base/langflow/api/v2/agui_translator.py b/src/backend/base/langflow/api/v2/agui_translator.py index 727e71e56e..0fa76a39f2 100644 --- a/src/backend/base/langflow/api/v2/agui_translator.py +++ b/src/backend/base/langflow/api/v2/agui_translator.py @@ -12,6 +12,7 @@ list. The translator is stateful: one instance per run. from __future__ import annotations import json +from dataclasses import dataclass, field from ag_ui.core import ( BaseEvent, @@ -37,6 +38,28 @@ from ag_ui.core import ( _CUSTOM_CONTENT_TYPES = frozenset({"json", "code", "media", "error"}) +@dataclass +class _BufferedMessage: + """A text message waiting for the wire while another message is open. + + AG-UI allows one open text message at a time, but parallel components + stream tokens for different message ids interleaved. Tokens for a message + that cannot open yet accumulate here; ``final_text`` is set when its + ``add_message`` finalizer arrives before it ever reached the wire. + """ + + chunks: list[str] = field(default_factory=list) + final_text: str | None = None + + @property + def complete(self) -> bool: + return self.final_text is not None + + @property + def text(self) -> str: + return self.final_text if self.final_text is not None else "".join(self.chunks) + + class AGUITranslator: """Translates Langflow ``EventManager`` events into AG-UI protocol events. @@ -47,11 +70,16 @@ class AGUITranslator: def __init__(self, run_id: str, thread_id: str) -> None: self.run_id = run_id self.thread_id = thread_id - # Id of the text message currently being streamed by ``token`` events, - # or ``None`` when no message is open. + # Id of the text message currently open on the wire, or ``None``. self._open_message_id: str | None = None - # Message ids already emitted as a complete (non-streamed) text message. + # Message ids whose TEXT_MESSAGE_END has been emitted; the protocol + # considers them closed, so nothing may reopen them. self._emitted_text_message_ids: set[str] = set() + # Messages from parallel components waiting for the wire, in arrival + # order (dict preserves insertion order). Invariant: non-empty only + # while another message holds the wire — closing the open message + # immediately promotes the next buffered one. + self._buffered_messages: dict[str, _BufferedMessage] = {} # Tool-call ids already emitted as TOOL_CALL_START / already resolved # with a TOOL_CALL_RESULT. ``add_message`` can re-fire with the same # (append-only) content_blocks, so emissions must be deduplicated. @@ -95,11 +123,11 @@ class AGUITranslator: # streamed message and must stay transparent, or the message would be # split into multiple START/END pairs reusing an already-ended id. if event_type == "end": - events = self._close_open_message() + events = self._drain_messages() events.append(RunFinishedEvent(run_id=self.run_id, thread_id=self.thread_id)) return events if event_type == "error": - events = self._close_open_message() + events = self._drain_messages() # The ``error`` payload varies by emission path: a full ErrorMessage # dump carries the reason in ``text``; the minimal path sends # ``{"error": str}``. @@ -111,11 +139,13 @@ class AGUITranslator: def _translate_token(self, data: dict) -> list[BaseEvent]: """Map a ``token`` event to text-message events. - The first token of a message opens it with ``TEXT_MESSAGE_START``; a - token for a different message id closes the previous one first. A - token whose id was already ended via ``add_message`` (or a prior - token boundary) is dropped: re-opening it would emit a second - ``TEXT_MESSAGE_START`` for an id the protocol considers closed. + The first token of a message opens it with ``TEXT_MESSAGE_START``. + While a message holds the wire, tokens for other message ids (parallel + components stream interleaved) buffer instead of closing it; the + buffered message flushes when the open one genuinely ends. A token + whose id was already ended via ``add_message`` is dropped: re-opening + it would emit a second ``TEXT_MESSAGE_START`` for an id the protocol + considers closed. """ message_id = str(data.get("id") or "") if not message_id: @@ -123,16 +153,22 @@ class AGUITranslator: # be correlated. Dropping the event is preferable to emitting a # malformed stream with empty message_ids. return [] - if message_id in self._emitted_text_message_ids and self._open_message_id != message_id: + if message_id in self._emitted_text_message_ids: return [] chunk = data.get("chunk", "") - events: list[BaseEvent] = [] - if self._open_message_id != message_id: - events.extend(self._close_open_message()) - events.append(TextMessageStartEvent(message_id=message_id, role="assistant")) + if self._open_message_id == message_id: + return [TextMessageContentEvent(message_id=message_id, delta=chunk)] + if self._open_message_id is None: self._open_message_id = message_id - events.append(TextMessageContentEvent(message_id=message_id, delta=chunk)) - return events + return [ + TextMessageStartEvent(message_id=message_id, role="assistant"), + TextMessageContentEvent(message_id=message_id, delta=chunk), + ] + # Another message holds the wire: buffer this one until it is free. + buffered = self._buffered_messages.setdefault(message_id, _BufferedMessage()) + if not buffered.complete: + buffered.chunks.append(chunk) + return [] def _translate_vertices_sorted(self, data: dict) -> list[BaseEvent]: """Map ``vertices_sorted`` to a ``STATE_SNAPSHOT`` of the node graph. @@ -198,21 +234,32 @@ class AGUITranslator: # Message text. if message_id and message_id == self._open_message_id: - # Finalizer of a token-streamed message: close it. The text was - # already streamed token by token, so it must not be re-emitted now - # or by any later add_message that re-fires for the same id. - self._emitted_text_message_ids.add(message_id) + # Finalizer of the wire-open streamed message: close it. The text + # was already streamed token by token, so it must not be re-emitted + # now or by any later add_message that re-fires for the same id. events.extend(self._close_open_message()) + elif message_id and message_id in self._buffered_messages: + # Finalizer for a message still waiting for the wire: record the + # authoritative full text; the trio is emitted on promotion. + self._buffered_messages[message_id].final_text = ( + data.get("text") or self._buffered_messages[message_id].text + ) else: text = data.get("text") or "" # Skip text-message lifecycle emission without a stable message_id; # tool-call events above are namespaced by block/content index so # they can still ride a missing id, but TEXT_MESSAGE_* cannot. if text and message_id and message_id not in self._emitted_text_message_ids: - self._emitted_text_message_ids.add(message_id) - events.append(TextMessageStartEvent(message_id=message_id, role="assistant")) - events.append(TextMessageContentEvent(message_id=message_id, delta=text)) - events.append(TextMessageEndEvent(message_id=message_id)) + if self._open_message_id is None: + self._emitted_text_message_ids.add(message_id) + events.append(TextMessageStartEvent(message_id=message_id, role="assistant")) + events.append(TextMessageContentEvent(message_id=message_id, delta=text)) + events.append(TextMessageEndEvent(message_id=message_id)) + else: + # A parallel component finished while another message holds + # the wire: buffer the complete message instead of opening + # a second text message mid-stream. + self._buffered_messages[message_id] = _BufferedMessage(final_text=text) return events def _translate_tool_use( @@ -286,17 +333,57 @@ class AGUITranslator: return {"op": "add", "path": f"/nodes/{node_id}", "value": {"status": status, "output": output}} def _close_open_message(self) -> list[BaseEvent]: - """Emit ``TEXT_MESSAGE_END`` for the open message, if any. + """Emit ``TEXT_MESSAGE_END`` for the open message and free the wire. The closed id is recorded in ``_emitted_text_message_ids`` so a later - token (e.g. an interleaved ``A, B, A`` sequence) cannot re-open it - and emit a second ``TEXT_MESSAGE_START`` for an id the protocol - already considers closed. + token cannot re-open it and emit a second ``TEXT_MESSAGE_START`` for + an id the protocol already considers closed. With the wire free, the + next buffered parallel message (if any) is promoted onto it. """ if self._open_message_id is None: return [] closed_id = self._open_message_id - end = TextMessageEndEvent(message_id=closed_id) self._emitted_text_message_ids.add(closed_id) self._open_message_id = None - return [end] + events: list[BaseEvent] = [TextMessageEndEvent(message_id=closed_id)] + events.extend(self._promote_next_buffered()) + return events + + def _promote_next_buffered(self) -> list[BaseEvent]: + """Move buffered parallel messages onto the freed wire, in arrival order. + + Already-complete messages emit their full START/CONTENT/END trio and the + promotion continues; the first still-streaming message replays its + buffered chunks, takes the wire, and stays open for its live tokens. + """ + events: list[BaseEvent] = [] + while self._buffered_messages and self._open_message_id is None: + message_id, buffered = next(iter(self._buffered_messages.items())) + del self._buffered_messages[message_id] + if message_id in self._emitted_text_message_ids: + continue + text = buffered.text + if buffered.complete: + self._emitted_text_message_ids.add(message_id) + if text: + events.append(TextMessageStartEvent(message_id=message_id, role="assistant")) + events.append(TextMessageContentEvent(message_id=message_id, delta=text)) + events.append(TextMessageEndEvent(message_id=message_id)) + continue + self._open_message_id = message_id + events.append(TextMessageStartEvent(message_id=message_id, role="assistant")) + if text: + events.append(TextMessageContentEvent(message_id=message_id, delta=text)) + return events + + def _drain_messages(self) -> list[BaseEvent]: + """Close the open message and flush every buffered one (run boundary). + + At ``end``/``error`` nothing else will free the wire, so buffered + parallel messages flush now — each promoted, emitted, and closed — + rather than being silently lost. + """ + events = self._close_open_message() + while self._open_message_id is not None: + events.extend(self._close_open_message()) + return events diff --git a/src/backend/tests/unit/api/v2/test_agui_translator.py b/src/backend/tests/unit/api/v2/test_agui_translator.py index 0797105bc7..995cea3d86 100644 --- a/src/backend/tests/unit/api/v2/test_agui_translator.py +++ b/src/backend/tests/unit/api/v2/test_agui_translator.py @@ -90,13 +90,31 @@ def test_token_sequence_emits_start_contents_then_end_on_boundary(): assert isinstance(ended[1], RunFinishedEvent) -def test_new_message_id_closes_previous_message_and_opens_new(): +def test_token_for_second_message_is_buffered_until_first_closes(): + """A token for a different message id must not close the streaming one. + + Parallel components stream tokens for different message ids interleaved. + Closing the open message on the first foreign token burned its id, so all + its later tokens were dropped. Instead the foreign message buffers until + the open one genuinely ends (its ``add_message`` finalizer), then flushes. + """ t = AGUITranslator(run_id="r1", thread_id="t1") t.start() t.translate("token", {"chunk": "a", "id": "m1"}) - out = t.translate("token", {"chunk": "b", "id": "m2"}) + buffered = t.translate("token", {"chunk": "b", "id": "m2"}) + # m2 buffers silently; m1 stays open. + assert buffered == [] + + # m1 keeps streaming: its tokens still flow. + more = t.translate("token", {"chunk": "a2", "id": "m1"}) + assert len(more) == 1 + assert isinstance(more[0], TextMessageContentEvent) + assert more[0].message_id == "m1" + + # m1's finalizer closes it and promotes m2 with its buffered content. + out = t.translate("add_message", {"id": "m1", "text": "aa2"}) assert isinstance(out[0], TextMessageEndEvent) assert out[0].message_id == "m1" assert isinstance(out[1], TextMessageStartEvent) @@ -166,32 +184,83 @@ def test_token_for_already_ended_message_id_is_dropped(): ) -def test_token_after_boundary_close_is_dropped(): - """An interleaved token sequence ``A, B, A`` must not re-open id A. +def test_parallel_interleaved_tokens_preserve_all_content(): + """Two components streaming in parallel must not lose either stream. - Switching from id A to id B closes A through ``_close_open_message``. A - later token for A used to slip past the dedup guard because - ``_close_open_message`` was the only finalizer that did not record the - closed id in ``_emitted_text_message_ids``. The translator must treat a - token-boundary close the same way it treats an ``add_message`` close. + Reproduces the reported parallel-components drop: with a single-slot + tracker, START(m2), CONTENT(m2), START(m3) burned m2, so CONTENT(m2) + was dropped and m2's remaining text never reached the client. With + buffering, m2 holds the wire until it genuinely ends; m3's tokens + buffer and flush afterwards. Nothing is dropped, and the wire carries + at most one open text message at a time. + """ + t = AGUITranslator(run_id="r1", thread_id="t1") + sequence = [ + ("token", {"chunk": "Result ", "id": "m2"}), + ("token", {"chunk": "Side ", "id": "m3"}), + ("token", {"chunk": "for X", "id": "m2"}), + ("token", {"chunk": "for Y", "id": "m3"}), + ("add_message", {"id": "m2", "text": "Result for X"}), + ("add_message", {"id": "m3", "text": "Side for Y"}), + ("end", {}), + ] + + out = _run_sequence(t, sequence) + + _assert_well_formed(out) + text_by_message: dict[str, str] = {} + for event in out: + if isinstance(event, TextMessageContentEvent): + text_by_message[event.message_id] = text_by_message.get(event.message_id, "") + event.delta + assert text_by_message["m2"] == "Result for X" + assert text_by_message["m3"] == "Side for Y" + + +def test_add_message_for_other_message_while_streaming_is_buffered(): + """A complete message landing mid-stream must not interleave its trio. + + A parallel non-streaming component can finish (add_message with full + text) while another component holds the wire with an open streamed + message. Emitting START/CONTENT/END for the finished one immediately + would open a second text message mid-stream; instead it buffers and + flushes when the open message closes. """ t = AGUITranslator(run_id="r1", thread_id="t1") t.start() - a1 = t.translate("token", {"chunk": "hi", "id": "m1"}) - assert any(isinstance(e, TextMessageStartEvent) and e.message_id == "m1" for e in a1) + t.translate("token", {"chunk": "streaming...", "id": "m1"}) + parallel_done = t.translate("add_message", {"id": "m9", "text": "finished early"}) + assert all( + not isinstance(e, (TextMessageStartEvent, TextMessageContentEvent, TextMessageEndEvent)) for e in parallel_done + ), f"buffered message leaked text events mid-stream: {parallel_done}" - # Switching to a new id closes m1 via _close_open_message. - b = t.translate("token", {"chunk": "yo", "id": "m2"}) - assert any(isinstance(e, TextMessageEndEvent) and e.message_id == "m1" for e in b) - assert any(isinstance(e, TextMessageStartEvent) and e.message_id == "m2" for e in b) + out = t.translate("add_message", {"id": "m1", "text": "streaming..."}) + assert [type(e) for e in out] == [ + TextMessageEndEvent, + TextMessageStartEvent, + TextMessageContentEvent, + TextMessageEndEvent, + ] + assert out[0].message_id == "m1" + assert out[1].message_id == "m9" + assert out[2].delta == "finished early" - # A late token for m1 must be dropped, not re-open the ended message. - late = t.translate("token", {"chunk": "again", "id": "m1"}) - assert all(not isinstance(e, TextMessageStartEvent) for e in late), ( - f"Token boundary did not mark m1 as ended; emitted: {late}" - ) - assert late == [], f"Expected no events for late token after boundary close; got {late}" + +def test_end_drains_buffered_messages_before_run_finished(): + """A run ending while messages are still buffered must flush them all.""" + t = AGUITranslator(run_id="r1", thread_id="t1") + t.start() + + t.translate("token", {"chunk": "open", "id": "m1"}) + t.translate("token", {"chunk": "waiting", "id": "m2"}) + out = t.translate("end", {}) + + assert isinstance(out[-1], RunFinishedEvent) + m2_content = [e for e in out if isinstance(e, TextMessageContentEvent) and e.message_id == "m2"] + assert len(m2_content) == 1 + assert m2_content[0].delta == "waiting" + ends = [e.message_id for e in out if isinstance(e, TextMessageEndEvent)] + assert ends == ["m1", "m2"] def test_vertices_sorted_emits_state_snapshot_of_all_nodes(): @@ -610,20 +679,22 @@ def _assert_well_formed(events: list) -> None: assert isinstance(events[0], RunStartedEvent) assert isinstance(events[-1], (RunFinishedEvent, RunErrorEvent)) - open_messages: set[str] = set() + open_message: str | None = None seen_messages: set[str] = set() for event in events: if isinstance(event, TextMessageStartEvent): - assert event.message_id not in open_messages, "text message started while already open" + # AG-UI's reference verifier allows at most one open text message; + # interleaved STARTs are rejected by conforming clients. + assert open_message is None, f"text message {event.message_id} started while {open_message} is open" assert event.message_id not in seen_messages, "text message id reused after it ended" - open_messages.add(event.message_id) + open_message = event.message_id seen_messages.add(event.message_id) elif isinstance(event, TextMessageContentEvent): - assert event.message_id in open_messages, "text content for a message that is not open" + assert event.message_id == open_message, "text content for a message that is not open" elif isinstance(event, TextMessageEndEvent): - assert event.message_id in open_messages, "text message ended without being open" - open_messages.discard(event.message_id) - assert not open_messages, "text messages left unclosed" + assert event.message_id == open_message, "text message ended without being open" + open_message = None + assert open_message is None, "text message left unclosed" started_tools: set[str] = set() ended_tools: set[str] = set() diff --git a/src/lfx/src/lfx/_assets/component_index.json b/src/lfx/src/lfx/_assets/component_index.json index 2a0ed39d02..9fc99281be 100644 --- a/src/lfx/src/lfx/_assets/component_index.json +++ b/src/lfx/src/lfx/_assets/component_index.json @@ -91135,7 +91135,7 @@ "icon": "bot", "legacy": false, "metadata": { - "code_hash": "ce7ae1d8e196", + "code_hash": "d08c17dbf863", "dependencies": { "dependencies": [ { @@ -91299,7 +91299,7 @@ "show": true, "title_case": false, "type": "code", - "value": "from __future__ import annotations\n\nimport uuid\nfrom contextlib import contextmanager\nfrom datetime import datetime, timezone\nfrom typing import TYPE_CHECKING, Any, cast\n\nfrom langchain.agents import create_agent\nfrom langchain.agents.middleware import ModelCallLimitMiddleware, ToolRetryMiddleware\n\nfrom lfx.components.models_and_agents.agent_helpers.graph_event_adapter import (\n adapt_graph_events_to_executor_shape,\n)\nfrom lfx.components.models_and_agents.agent_helpers.messages_input_builder import (\n build_initial_messages,\n)\nfrom lfx.components.models_and_agents.agent_helpers.placeholder_corrective_middleware import (\n WatsonXPlaceholderMiddleware,\n)\nfrom lfx.components.models_and_agents.agent_helpers.single_tool_call_middleware import (\n SingleToolCallMiddleware,\n)\nfrom lfx.components.models_and_agents.memory import MemoryComponent, aget_agent_chat_history\n\nif TYPE_CHECKING:\n from langchain_core.tools import Tool\n\n from lfx.schema.log import OnTokenFunctionType, SendMessageFunctionType\n\nfrom lfx.base.agents.agent import LCToolsAgentComponent\nfrom lfx.base.agents.callback import AgentAsyncHandler\nfrom lfx.base.agents.default_system_prompt import DEFAULT_SYSTEM_PROMPT_TEMPLATE\nfrom lfx.base.agents.events import ExceptionWithMessageError, process_agent_events\nfrom lfx.base.agents.token_callback import TokenUsageCallbackHandler\nfrom lfx.base.agents.utils import get_chat_output_sender_name\nfrom lfx.base.constants import STREAM_INFO_TEXT\nfrom lfx.base.models.unified_models import (\n get_language_model_options,\n get_llm,\n handle_model_input_update,\n)\nfrom lfx.base.models.watsonx_constants import IBM_WATSONX_URLS\nfrom lfx.components.agentics.helpers.model_config import validate_model_selection\nfrom lfx.components.helpers import CalculatorComponent, CurrentDateComponent\nfrom lfx.components.langchain_utilities.ibm_granite_handler import is_watsonx_model\nfrom lfx.components.langchain_utilities.tool_calling import ToolCallingAgentComponent\nfrom lfx.custom.custom_component.component import get_component_toolkit\nfrom lfx.field_typing.range_spec import RangeSpec\nfrom lfx.inputs.inputs import BoolInput, DropdownInput, ModelInput, StrInput\nfrom lfx.io import IntInput, MessageTextInput, MultilineInput, Output, SecretStrInput, TableInput\nfrom lfx.log.logger import logger\nfrom lfx.memory import delete_message\nfrom lfx.schema.data import Data\nfrom lfx.schema.dotdict import dotdict\nfrom lfx.schema.message import Message\nfrom lfx.schema.table import EditMode\nfrom lfx.utils.constants import MESSAGE_SENDER_AI\n\n\ndef set_advanced_true(component_input):\n component_input.advanced = True\n return component_input\n\n\ndef _agent_base_inputs():\n \"\"\"Return base inputs tailored to AgentComponent's create_agent path.\n\n `get_base_inputs()` returns a shared list — replace, don't mutate. We drop\n inputs that are no-ops here and override info text on the inputs whose\n semantics shifted under create_agent.\n\n `verbose` is dropped because the create_agent event stream already surfaces\n every agent step via the \"Agent Steps\" content blocks; the legacy boolean\n has nothing to toggle. Saved flows that still carry a `verbose` value just\n ignore it on load (the schema no longer declares it).\n \"\"\"\n drop = {\"verbose\"}\n overrides = {\n \"handle_parsing_errors\": BoolInput(\n name=\"handle_parsing_errors\",\n display_name=\"Handle Parse Errors\",\n value=True,\n advanced=True,\n info=(\n \"Adds tool-execution retry as a safety net. `create_agent` already \"\n \"feeds tool-call validation errors back to the LLM automatically; \"\n \"this flag layers `ToolRetryMiddleware` on top so transient tool \"\n \"runtime failures are retried (max 2 retries).\"\n ),\n ),\n \"max_iterations\": IntInput(\n name=\"max_iterations\",\n display_name=\"Max Iterations\",\n value=15,\n advanced=True,\n range_spec=RangeSpec(min=1, max=128000, step=1, step_type=\"int\"),\n info=(\n \"Maximum number of model calls the agent can make before stopping \"\n \"(maps to `ModelCallLimitMiddleware.run_limit` on the create_agent \"\n \"path). Must be at least 1 — it is a safety cap, never 'unlimited'.\"\n ),\n ),\n }\n return [overrides.get(inp.name, inp) for inp in LCToolsAgentComponent.get_base_inputs() if inp.name not in drop]\n\n\ndef _extract_text_content(value) -> str:\n \"\"\"Pull a string payload from a Message-like, AIMessage-like, or string value.\"\"\"\n if isinstance(value, str):\n return value\n text = getattr(value, \"text\", None)\n if isinstance(text, str):\n return text\n content = getattr(value, \"content\", None)\n if isinstance(content, str):\n return content\n return str(value) if value is not None else \"\"\n\n\n@contextmanager\ndef _suppress_send_message(component: Any):\n \"\"\"Temporarily replace component.send_message with a no-op for the duration of the block.\n\n Used during the structured-output prompt fallback: run_agent streams the agent's\n final answer through self.send_message (correct for message_response), but in\n json_response the orchestrator parses that text into structured Data which the\n downstream Chat Output emits — leaving the original emission in place produces a\n duplicate message in the playground. The original method is always restored on exit,\n even when the wrapped call raises.\n \"\"\"\n original = component.send_message\n\n async def _noop(message, *_args, **_kwargs):\n return message\n\n component.send_message = _noop\n try:\n yield\n finally:\n component.send_message = original\n\n\nclass AgentComponent(ToolCallingAgentComponent):\n display_name: str = \"Agent\"\n description: str = \"Define the agent's instructions, then enter a task to complete using tools.\"\n documentation: str = \"https://docs.langflow.org/agents\"\n icon = \"bot\"\n beta = False\n name = \"Agent\"\n\n memory_inputs = [set_advanced_true(component_input) for component_input in MemoryComponent().inputs]\n\n inputs = [\n ModelInput(\n name=\"model\",\n display_name=\"Language Model\",\n info=\"Select your model provider\",\n real_time_refresh=True,\n required=True,\n # Agents require tool calling — the filter is honored by\n # ``handle_model_input_update`` so models that can't run with\n # tools never reach the picker (and any saved selection that\n # no longer satisfies the constraint is auto-replaced).\n filters={\"tool_calling\": True},\n ),\n SecretStrInput(\n name=\"api_key\",\n display_name=\"API Key\",\n info=\"Overrides global provider settings. Leave blank to use your pre-configured API Key.\",\n real_time_refresh=True,\n advanced=True,\n ),\n DropdownInput(\n name=\"base_url_ibm_watsonx\",\n display_name=\"watsonx API Endpoint\",\n info=\"The base URL of the API (IBM watsonx.ai only)\",\n options=IBM_WATSONX_URLS,\n value=IBM_WATSONX_URLS[0],\n combobox=True,\n show=False,\n real_time_refresh=True,\n ),\n StrInput(\n name=\"project_id\",\n display_name=\"watsonx Project ID\",\n info=\"The project ID associated with the foundation model (IBM watsonx.ai only)\",\n show=False,\n required=False,\n ),\n MultilineInput(\n name=\"system_prompt\",\n display_name=\"Agent Instructions\",\n info=(\n \"System Prompt: Initial instructions and context provided to guide the agent's behavior. \"\n \"Supports dynamic placeholders: {current_date}, {model_name}, {optional_user_context}.\"\n ),\n value=DEFAULT_SYSTEM_PROMPT_TEMPLATE,\n advanced=False,\n ),\n MessageTextInput(\n name=\"context_id\",\n display_name=\"Context ID\",\n info=\"The context ID of the chat. Adds an extra layer to the local memory.\",\n value=\"\",\n advanced=True,\n ),\n IntInput(\n name=\"n_messages\",\n display_name=\"Number of Chat History Messages\",\n value=100,\n info=\"Number of chat history messages to retrieve.\",\n advanced=True,\n show=True,\n ),\n IntInput(\n name=\"max_tokens\",\n display_name=\"Max Tokens\",\n info=\"Maximum number of tokens to generate. Field name varies by provider.\",\n advanced=True,\n range_spec=RangeSpec(min=1, max=128000, step=1, step_type=\"int\"),\n ),\n MultilineInput(\n name=\"format_instructions\",\n display_name=\"Output Format Instructions\",\n info=\"Generic Template for structured output formatting. Valid only with Structured response.\",\n value=(\n \"You are an AI that extracts structured JSON objects from unstructured text. \"\n \"Use a predefined schema with expected types (str, int, float, bool, dict). \"\n \"Extract ALL relevant instances that match the schema - if multiple patterns exist, capture them all. \"\n \"Fill missing or ambiguous values with defaults: null for missing values. \"\n \"Remove exact duplicates but keep variations that have different field values. \"\n \"Always return valid JSON in the expected format, never throw errors. \"\n \"If multiple objects can be extracted, return them all in the structured format.\"\n ),\n advanced=True,\n ),\n TableInput(\n name=\"output_schema\",\n display_name=\"Output Schema\",\n info=(\n \"Schema Validation: Define the structure and data types for structured output. \"\n \"No validation if no output schema.\"\n ),\n advanced=True,\n required=False,\n value=[],\n table_schema=[\n {\n \"name\": \"name\",\n \"display_name\": \"Name\",\n \"type\": \"str\",\n \"description\": \"Specify the name of the output field.\",\n \"default\": \"field\",\n \"edit_mode\": EditMode.INLINE,\n },\n {\n \"name\": \"description\",\n \"display_name\": \"Description\",\n \"type\": \"str\",\n \"description\": \"Describe the purpose of the output field.\",\n \"default\": \"description of field\",\n \"edit_mode\": EditMode.POPOVER,\n },\n {\n \"name\": \"type\",\n \"display_name\": \"Type\",\n \"type\": \"str\",\n \"edit_mode\": EditMode.INLINE,\n \"description\": (\"Indicate the data type of the output field (e.g., str, int, float, bool, dict).\"),\n \"options\": [\"str\", \"int\", \"float\", \"bool\", \"dict\"],\n \"default\": \"str\",\n },\n {\n \"name\": \"multiple\",\n \"display_name\": \"As List\",\n \"type\": \"boolean\",\n \"description\": \"Set to True if this output field should be a list of the specified type.\",\n \"default\": \"False\",\n \"edit_mode\": EditMode.INLINE,\n },\n ],\n ),\n *_agent_base_inputs(),\n # removed memory inputs from agent component\n # *memory_inputs,\n BoolInput(\n name=\"stream\",\n display_name=\"Stream\",\n info=STREAM_INFO_TEXT,\n value=True,\n advanced=True,\n ),\n BoolInput(\n name=\"add_current_date_tool\",\n display_name=\"Current Date\",\n advanced=True,\n info=\"If true, will add a tool to the agent that returns the current date.\",\n value=True,\n ),\n BoolInput(\n name=\"add_calculator_tool\",\n display_name=\"Calculator\",\n advanced=True,\n info=(\n \"If true, adds a zero-config arithmetic calculator tool to the agent \"\n \"(safe: only +, -, *, /, ** operators via AST).\"\n ),\n value=True,\n ),\n ]\n outputs = [\n Output(name=\"response\", display_name=\"Response\", method=\"message_response\"),\n Output(\n name=\"structured_response\",\n display_name=\"Structured Response\",\n method=\"json_response\",\n types=[\"Data\"],\n ),\n ]\n\n def _resolve_selected_model(self):\n \"\"\"Resolve the selected model, including legacy agent_llm/model_name inputs.\"\"\"\n try:\n from langchain_core.language_models import BaseLanguageModel\n\n if isinstance(self.model, BaseLanguageModel):\n return self.model\n except ImportError:\n pass\n\n if isinstance(self.model, list) and self.model:\n return self.model\n\n legacy_provider = getattr(self, \"agent_llm\", None)\n legacy_model_name = getattr(self, \"model_name\", None)\n if not legacy_provider or not legacy_model_name:\n return self.model\n\n options = get_language_model_options(user_id=self.user_id)\n for option in options:\n if option.get(\"provider\") == legacy_provider and option.get(\"name\") == legacy_model_name:\n return [option]\n\n return [\n {\n \"name\": legacy_model_name,\n \"provider\": legacy_provider,\n \"metadata\": {},\n }\n ]\n\n def _get_max_tokens_value(self):\n \"\"\"Return the user-supplied max_tokens or None when unset/zero.\"\"\"\n val = getattr(self, \"max_tokens\", None)\n if val in {\"\", 0}:\n return None\n return val\n\n def _get_llm(self):\n \"\"\"Override parent to include max_tokens from the Agent's input field.\n\n Streaming is mandatory for AgentComponent: ``runnable.astream_events(v2)`` only\n emits ``on_chat_model_stream`` chunks when the underlying chat model is\n instantiated with ``streaming=True``. Unlike the LanguageModel component (where\n ``stream`` is a user-facing toggle), the Agent has no opt-out — the toggle is\n kept in the UI for backwards compatibility but is intentionally ignored here.\n Without ``stream=True``, the chat model accumulates the whole response and\n only emits ``on_chat_model_end``, silently disabling the Playground's live-\n typing view and breaking the streaming contract on the /events surface.\n \"\"\"\n return get_llm(\n model=self.model,\n user_id=self.user_id,\n api_key=getattr(self, \"api_key\", None),\n stream=True,\n max_tokens=self._get_max_tokens_value(),\n watsonx_url=getattr(self, \"base_url_ibm_watsonx\", None),\n watsonx_project_id=getattr(self, \"project_id\", None),\n )\n\n async def get_agent_requirements(self):\n \"\"\"Get the agent requirements for the agent.\"\"\"\n from langchain_core.tools import StructuredTool\n\n selected_model = self._resolve_selected_model()\n try:\n from langchain_core.language_models import BaseLanguageModel\n\n is_connected_model = isinstance(selected_model, BaseLanguageModel)\n except ImportError:\n is_connected_model = False\n\n if not is_connected_model:\n validate_model_selection(selected_model)\n\n # Ensure _get_llm() uses the resolved model (e.g. from legacy agent_llm/model_name)\n self.model = selected_model\n llm_model = self._get_llm()\n if llm_model is None:\n msg = \"No language model selected. Please choose a model to proceed.\"\n raise ValueError(msg)\n\n # Get memory data\n self.chat_history = await self.get_memory_data()\n await logger.adebug(f\"Retrieved {len(self.chat_history)} chat history messages\")\n if isinstance(self.chat_history, Message):\n self.chat_history = [self.chat_history]\n\n # Add current date tool if enabled\n if self.add_current_date_tool:\n if not isinstance(self.tools, list): # type: ignore[has-type]\n self.tools = []\n current_date_tool = (await CurrentDateComponent(**self.get_base_args()).to_toolkit()).pop(0)\n\n if not isinstance(current_date_tool, StructuredTool):\n msg = \"CurrentDateComponent must be converted to a StructuredTool\"\n raise TypeError(msg)\n # Skip if an externally-connected tool already provides the same name.\n # Duplicate tool names are rejected by Anthropic/Gemini with HTTP 400.\n if not any(getattr(t, \"name\", None) == current_date_tool.name for t in self.tools):\n self.tools.append(current_date_tool)\n\n # Add calculator tool if enabled (zero-config arithmetic)\n if getattr(self, \"add_calculator_tool\", False):\n if not isinstance(self.tools, list): # type: ignore[has-type]\n self.tools = []\n calculator_tool = (await CalculatorComponent(**self.get_base_args()).to_toolkit()).pop(0)\n\n if not isinstance(calculator_tool, StructuredTool):\n msg = \"CalculatorComponent must be converted to a StructuredTool\"\n raise TypeError(msg)\n # Skip if an externally-connected tool already provides the same name.\n # Duplicate tool names are rejected by Anthropic/Gemini with HTTP 400.\n if not any(getattr(t, \"name\", None) == calculator_tool.name for t in self.tools):\n self.tools.append(calculator_tool)\n\n # Set shared callbacks for tracing the tools used by the agent\n self.set_tools_callbacks(self.tools, self._get_shared_callbacks())\n\n return llm_model, self.chat_history, self.tools\n\n def _get_resolved_model_name(self) -> str:\n \"\"\"Best-effort human-readable model name for {model_name} injection.\"\"\"\n try:\n from langchain_core.language_models import BaseLanguageModel\n\n if isinstance(self.model, BaseLanguageModel):\n return type(self.model).__name__\n except ImportError:\n pass\n\n if isinstance(self.model, list) and self.model:\n first = self.model[0]\n if isinstance(first, dict):\n name = first.get(\"name\")\n if isinstance(name, str) and name:\n return name\n\n legacy_model_name = getattr(self, \"model_name\", None)\n if isinstance(legacy_model_name, str) and legacy_model_name:\n return legacy_model_name\n return \"\"\n\n def _inject_dynamic_prompt_values(self, prompt: Any | None) -> str | None:\n \"\"\"Replace known env placeholders in the system prompt.\n\n Handles {current_date}, {model_name}, and {optional_user_context} (the\n last one ships with the structured DEFAULT_SYSTEM_PROMPT_TEMPLATE and\n is currently unused at the AgentComponent layer, so it resolves to \"\").\n Uses str.replace (not str.format) so user prompts containing literal\n braces such as JSON examples ({\"key\": 1}) never break the agent.\n\n `system_prompt` is a connectable MultilineInput, so the value can arrive\n as a Message (e.g. a Prompt node wired in). Normalize it to text first —\n a raw Message has no `.replace` and used to crash the agent build.\n \"\"\"\n if prompt is None:\n return None\n prompt = _extract_text_content(prompt)\n if not prompt:\n return prompt\n replacements = {\n \"{current_date}\": datetime.now(tz=timezone.utc).strftime(\"%Y-%m-%d %H:%M:%S UTC\"),\n \"{model_name}\": self._get_resolved_model_name(),\n \"{optional_user_context}\": \"\",\n }\n for placeholder, value in replacements.items():\n prompt = prompt.replace(placeholder, value)\n return prompt\n\n def create_agent_runnable(self):\n \"\"\"Build the LangGraph `CompiledStateGraph` via `langchain.agents.create_agent`.\n\n Replaces the legacy `AgentExecutor` runnable inherited from\n `ToolCallingAgentComponent`. Other agent components (tool_calling, csv, json,\n openapi, sql*, vector_store_router) keep the legacy path — only AgentComponent\n runs on the new graph API.\n\n `max_iterations` and `handle_parsing_errors` (legacy AgentExecutor knobs) are\n translated to LangGraph middleware. Without that translation those user inputs\n would silently become no-ops on the new API.\n\n Provider notes:\n - WatsonX/Granite work natively with create_agent — `ChatWatsonx.bind_tools`\n handles tool_choice correctly. The legacy `create_granite_agent` path was\n dropped because it hardcoded `tool_choice='required'`, which the WatsonX\n API now rejects.\n - Ollama and other small/local models often emit malformed tool args. The\n ToolRetryMiddleware (default `retry_on=(Exception,)`, `on_failure='continue'`)\n catches Pydantic ValidationErrors from bad args and feeds the error back\n to the LLM as a retry signal, so the agent recovers gracefully.\n \"\"\"\n llm = self._get_llm()\n tools = self.tools or []\n\n # Eager bind_tools validation. `create_agent(...)` is lazy — without this,\n # an LLM that doesn't support tool calling fails on the first user message\n # instead of when the user wires up the component, which is a much worse UX.\n # Gated on a non-empty tools list so a no-tool Agent on a plain chat model\n # (which legitimately has no `bind_tools`) isn't shut out at flow-build time.\n # Providers signal \"no tool calling\" inconsistently — `NotImplementedError`\n # (langchain default), `AttributeError` (no `bind_tools` attr), or `TypeError`\n # (signature mismatch). Treat all three as the same UX failure.\n if tools:\n try:\n llm.bind_tools(tools)\n except (NotImplementedError, AttributeError, TypeError) as exc:\n # Include the underlying error so a broken tool schema or a\n # provider implementation bug is not silently disguised as a\n # \"model can't call tools\" UX error.\n msg = (\n f\"{self.display_name} does not support tool calling, \"\n \"or one of the connected tools failed to bind. \"\n \"Please connect a tool-calling capable language model and \"\n f\"verify your tools are well-formed. Underlying error: {exc!s}\"\n )\n raise NotImplementedError(msg) from exc\n\n middleware = self._build_middleware(llm)\n return create_agent(\n model=llm,\n tools=tools,\n system_prompt=self.system_prompt or \"\",\n middleware=middleware or None,\n )\n\n def _compute_recursion_limit(self) -> int:\n \"\"\"Derive the LangGraph recursion_limit from the user-set max_iterations.\n\n Mirrors the clamp in `_build_middleware` (max(1, max_iterations)) so a\n saved 0 or negative value cannot under-cap the graph below one full\n iteration. The +5 buffer covers start/end/router overhead.\n \"\"\"\n raw = getattr(self, \"max_iterations\", None)\n run_limit = max(1, int(raw)) if raw is not None else 15\n return run_limit * 2 + 5\n\n def _build_middleware(self, llm: Any) -> list:\n # `llm` is passed in (rather than re-fetched via `self._get_llm()`)\n # because some providers do credential resolution / client instantiation\n # lazily on each call. The caller — `create_agent_runnable` — already\n # resolved it once for `bind_tools`, so reuse that instance here.\n middleware: list = []\n max_iterations = getattr(self, \"max_iterations\", None)\n if max_iterations is not None:\n # `max_iterations` is a safety cap, not an \"unlimited\" toggle. A saved\n # 0 or negative value (falsy) must NOT silently drop the limiter and\n # allow an unbounded model/tool loop — clamp it to a real minimum.\n run_limit = max(1, int(max_iterations))\n middleware.append(ModelCallLimitMiddleware(run_limit=run_limit))\n # ToolRetryMiddleware only matters when there ARE tools to retry. Attaching\n # it on a no-tools agent inflates the compiled graph and adds per-invocation\n # middleware overhead for nothing, which is a measurable contributor to\n # trivial-prompt latency (QA UI-003).\n if getattr(self, \"handle_parsing_errors\", False) and self.tools:\n middleware.append(ToolRetryMiddleware(max_retries=2))\n # WatsonX models have two known platform quirks; both still reproduce on\n # the current API, so we keep the protections from the legacy\n # `create_granite_agent` path.\n # 1. Multi-tool-call assistant turns are rejected (\"This model only\n # supports single tool-calls at once!\"). Clamp to one per turn.\n # 2. Tool args occasionally come back as literal placeholder strings\n # (e.g. ``). Re-invoke once with a corrective\n # SystemMessage.\n # Order: SingleToolCallMiddleware first (outermost) so the clamp is\n # applied to the final response, including any corrective re-invoke\n # produced by WatsonXPlaceholderMiddleware.\n if is_watsonx_model(llm):\n middleware.append(SingleToolCallMiddleware())\n middleware.append(WatsonXPlaceholderMiddleware())\n return middleware\n\n async def run_agent(self, agent) -> Message:\n \"\"\"Run the LangGraph `CompiledStateGraph` and return the final agent Message.\n\n Overrides the legacy `LCAgentComponent.run_agent` (which builds an\n `{\"input\": str, \"chat_history\": [...]}` dict for `AgentExecutor`). The graph\n wants `{\"messages\": [BaseMessage, ...]}`. The event stream is wrapped with\n `adapt_graph_events_to_executor_shape` so the legacy `process_agent_events`\n (in `lfx.base.agents.events`) can be reused unchanged.\n \"\"\"\n messages = build_initial_messages(\n input_value=self.input_value,\n chat_history=getattr(self, \"chat_history\", None),\n )\n input_dict = {\"messages\": messages}\n\n agent_message = self._build_initial_agent_message()\n token_usage_handler = TokenUsageCallbackHandler()\n\n # Stream tokens to the event manager when running inside the Langflow runtime.\n # This is what powers the live-typing view in the chat UI.\n on_token_callback: OnTokenFunctionType | None = None\n if getattr(self, \"_event_manager\", None):\n on_token_callback = cast(\"OnTokenFunctionType\", self._event_manager.on_token)\n\n # Align LangGraph's `recursion_limit` with `max_iterations` so the\n # middleware cap (ModelCallLimitMiddleware) is what bounds the loop —\n # not LangGraph's default 25-step guard, which fires at ~12 model+tool\n # iterations and raises a raw GraphRecursionError (QA UI-009/UI-010).\n # Each iteration is ~2 graph steps (model node + tools node); add 5\n # for start/end overhead.\n recursion_limit = self._compute_recursion_limit()\n\n stream = adapt_graph_events_to_executor_shape(\n agent.astream_events(\n input_dict,\n config={\n \"callbacks\": [\n AgentAsyncHandler(self.log),\n token_usage_handler,\n *self._get_shared_callbacks(),\n ],\n \"recursion_limit\": recursion_limit,\n },\n version=\"v2\",\n )\n )\n try:\n result = await process_agent_events(\n stream,\n agent_message,\n cast(\"SendMessageFunctionType\", self.send_message),\n on_token_callback,\n )\n except ExceptionWithMessageError as e:\n # Drop the half-stored partial message from the DB (only if it was\n # actually persisted) and tell the frontend to remove the stale bubble.\n if hasattr(e, \"agent_message\"):\n msg_id = e.agent_message.get_id()\n if msg_id:\n await delete_message(id_=msg_id)\n await self._send_message_event(e.agent_message, category=\"remove_message\")\n logger.error(f\"ExceptionWithMessageError: {e}\")\n raise\n\n usage_data = token_usage_handler.get_usage()\n if usage_data:\n self._token_usage = usage_data\n result.properties.usage = usage_data\n # Only round-trip the DB when the message was stored (Chat Output wired).\n # `_should_skip_message=True` leaves `result.get_id()` empty; persisting\n # then would create a phantom row.\n if result.get_id():\n stored_result = await self._update_stored_message(result)\n await self._send_message_event(stored_result)\n result = stored_result\n\n self.status = result\n return result\n\n def _build_initial_agent_message(self) -> Message:\n \"\"\"Construct the placeholder agent Message that `process_agent_events` mutates.\"\"\"\n if hasattr(self, \"graph\"):\n session_id = self.graph.session_id\n elif hasattr(self, \"_session_id\"):\n session_id = self._session_id\n else:\n session_id = None\n\n sender_name = get_chat_output_sender_name(self) or self.display_name or \"AI\"\n return Message(\n sender=MESSAGE_SENDER_AI,\n sender_name=sender_name,\n properties={\"icon\": \"Bot\", \"state\": \"partial\"},\n # `text=\"\"` sentinel so MessageTable's no_content check accepts\n # an in-flight agent message whose content_blocks haven't been\n # populated yet. Mirrors ChatInput's convention.\n text=\"\",\n # Flat chronological event log; see lfx.base.agents.events.\n content_blocks=[],\n session_id=session_id or uuid.uuid4(),\n )\n\n async def message_response(self) -> Message:\n try:\n llm_model, self.chat_history, self.tools = await self.get_agent_requirements()\n # Set up and run agent\n self.set(\n llm=llm_model,\n tools=self.tools or [],\n chat_history=self.chat_history,\n input_value=self.input_value,\n system_prompt=self._inject_dynamic_prompt_values(self.system_prompt),\n )\n agent = self.create_agent_runnable()\n result = await self.run_agent(agent)\n\n # Store result for potential JSON output\n self._agent_result = result\n\n except (ValueError, TypeError, KeyError) as e:\n await logger.aerror(f\"{type(e).__name__}: {e!s}\")\n raise\n except ExceptionWithMessageError as e:\n await logger.aerror(f\"ExceptionWithMessageError occurred: {e}\")\n raise\n # Avoid catching blind Exception; let truly unexpected exceptions propagate\n except Exception as e:\n await logger.aerror(f\"Unexpected error: {e!s}\")\n raise\n else:\n return result\n\n async def json_response(self) -> Data:\n \"\"\"Produce structured Data via native LLM structured output, with prompt-based fallback.\n\n Native path (no tools, llm has with_structured_output) bypasses the agent loop and\n returns provider-validated JSON. When tools are attached, falls back to running the\n agent with a schema-augmented system prompt and parsing the final message content.\n \"\"\"\n from lfx.components.models_and_agents.structured_output.structured_output_orchestrator import (\n orchestrate_structured_output,\n )\n\n try:\n llm_model, self.chat_history, self.tools = await self.get_agent_requirements()\n except (ValueError, TypeError) as exc:\n await logger.aerror(f\"json_response.requirements_failed: {exc}\")\n return Data(data={\"content\": \"\", \"error\": str(exc)})\n\n injected_system_prompt = self._inject_dynamic_prompt_values(getattr(self, \"system_prompt\", \"\") or \"\") or \"\"\n format_instructions = getattr(self, \"format_instructions\", \"\") or \"\"\n output_schema = getattr(self, \"output_schema\", None) or []\n has_tools = bool(self.tools)\n\n async def _run_agent_for_fallback(augmented_prompt: str) -> str:\n self.set(\n llm=llm_model,\n tools=self.tools or [],\n chat_history=self.chat_history,\n input_value=self.input_value,\n system_prompt=augmented_prompt,\n )\n agent_runnable = self.create_agent_runnable()\n with _suppress_send_message(self):\n result = await self.run_agent(agent_runnable)\n return _extract_text_content(result)\n\n try:\n return await orchestrate_structured_output(\n llm=llm_model,\n output_schema=output_schema,\n system_prompt=injected_system_prompt,\n format_instructions=format_instructions,\n input_value=_extract_text_content(self.input_value),\n run_prompt_fallback=_run_agent_for_fallback,\n prefer_native=not has_tools,\n )\n except (\n ExceptionWithMessageError,\n ValueError,\n TypeError,\n NotImplementedError,\n AttributeError,\n ) as exc:\n await logger.aerror(f\"json_response.orchestration_failed: {exc}\")\n return Data(data={\"content\": \"\", \"error\": str(exc)})\n\n async def get_memory_data(self):\n # Scope by flow_id so default playground session names (e.g. \"New Session 0\")\n # cannot leak chat history across unrelated flows. See issue #13059.\n # The helper also returns [] when n_messages == 0, preserving the\n # explicit \"memory disabled\" contract from MemoryComponent.retrieve_messages.\n messages = await aget_agent_chat_history(\n session_id=self.graph.session_id,\n flow_id=getattr(self.graph, \"flow_id\", None),\n context_id=self.context_id,\n n_messages=self.n_messages,\n )\n return [\n message for message in messages if getattr(message, \"id\", None) != getattr(self.input_value, \"id\", None)\n ]\n\n def update_input_types(self, build_config: dotdict) -> dotdict:\n \"\"\"Update input types for all fields in build_config.\"\"\"\n for key, value in build_config.items():\n if isinstance(value, dict):\n if value.get(\"input_types\") is None:\n build_config[key][\"input_types\"] = []\n elif hasattr(value, \"input_types\") and value.input_types is None:\n value.input_types = []\n return build_config\n\n async def update_build_config(\n self,\n build_config: dotdict,\n field_value: list[dict],\n field_name: str | None = None,\n ) -> dotdict:\n # Update model options with caching (for all field changes).\n # The tool-calling constraint lives on the ModelInput's ``filters``\n # field (declared above); ``handle_model_input_update`` reads it\n # and applies the filter to both the dropdown options and the\n # sticky-default re-injection path.\n build_config = handle_model_input_update(\n component=self,\n build_config=dict(build_config),\n field_value=field_value,\n field_name=field_name,\n )\n build_config = dotdict(build_config)\n\n if field_name == \"model\":\n build_config = self.update_input_types(build_config)\n\n # Validate required keys. `verbose` was dropped from the input set\n # (see `_agent_base_inputs` — the create_agent event stream already\n # surfaces every step), so it is intentionally NOT required here.\n # Saved flows that still carry a `verbose` value just ignore it on\n # load.\n default_keys = [\n \"code\",\n \"_type\",\n \"model\",\n \"tools\",\n \"input_value\",\n \"add_current_date_tool\",\n \"add_calculator_tool\",\n \"system_prompt\",\n \"max_iterations\",\n \"handle_parsing_errors\",\n ]\n missing_keys = [key for key in default_keys if key not in build_config]\n if missing_keys:\n msg = f\"Missing required keys in build_config: {missing_keys}\"\n raise ValueError(msg)\n return dotdict({k: v.to_dict() if hasattr(v, \"to_dict\") else v for k, v in build_config.items()})\n\n async def _get_tools(self) -> list[Tool]:\n component_toolkit = get_component_toolkit()\n\n tools = component_toolkit(component=self).get_tools(\n tool_name=\"Call_Agent\",\n # here we do not use the shared callbacks as we are exposing the agent as a tool\n callbacks=self.get_langchain_callbacks(),\n )\n if hasattr(self, \"tools_metadata\"):\n tools = component_toolkit(component=self, metadata=self.tools_metadata).update_tools_metadata(tools=tools)\n\n return tools\n" + "value": "from __future__ import annotations\n\nimport uuid\nfrom contextlib import contextmanager\nfrom datetime import datetime, timezone\nfrom typing import TYPE_CHECKING, Any, cast\n\nfrom langchain.agents import create_agent\nfrom langchain.agents.middleware import ModelCallLimitMiddleware, ToolRetryMiddleware\n\nfrom lfx.components.models_and_agents.agent_helpers.graph_event_adapter import (\n adapt_graph_events_to_executor_shape,\n)\nfrom lfx.components.models_and_agents.agent_helpers.messages_input_builder import (\n build_initial_messages,\n)\nfrom lfx.components.models_and_agents.agent_helpers.placeholder_corrective_middleware import (\n WatsonXPlaceholderMiddleware,\n)\nfrom lfx.components.models_and_agents.agent_helpers.single_tool_call_middleware import (\n SingleToolCallMiddleware,\n)\nfrom lfx.components.models_and_agents.memory import MemoryComponent, aget_agent_chat_history\n\nif TYPE_CHECKING:\n from langchain_core.tools import Tool\n\n from lfx.schema.log import OnTokenFunctionType, SendMessageFunctionType\n\nfrom lfx.base.agents.agent import LCToolsAgentComponent\nfrom lfx.base.agents.callback import AgentAsyncHandler\nfrom lfx.base.agents.default_system_prompt import DEFAULT_SYSTEM_PROMPT_TEMPLATE\nfrom lfx.base.agents.events import ExceptionWithMessageError, process_agent_events\nfrom lfx.base.agents.token_callback import TokenUsageCallbackHandler\nfrom lfx.base.agents.utils import get_chat_output_sender_name\nfrom lfx.base.constants import STREAM_INFO_TEXT\nfrom lfx.base.models.unified_models import (\n get_language_model_options,\n get_llm,\n handle_model_input_update,\n)\nfrom lfx.base.models.watsonx_constants import IBM_WATSONX_URLS\nfrom lfx.components.agentics.helpers.model_config import validate_model_selection\nfrom lfx.components.helpers import CalculatorComponent, CurrentDateComponent\nfrom lfx.components.langchain_utilities.ibm_granite_handler import is_watsonx_model\nfrom lfx.components.langchain_utilities.tool_calling import ToolCallingAgentComponent\nfrom lfx.custom.custom_component.component import get_component_toolkit\nfrom lfx.field_typing.range_spec import RangeSpec\nfrom lfx.inputs.inputs import BoolInput, DropdownInput, ModelInput, StrInput\nfrom lfx.io import IntInput, MessageTextInput, MultilineInput, Output, SecretStrInput, TableInput\nfrom lfx.log.logger import logger\nfrom lfx.memory import delete_message\nfrom lfx.schema.content_block import ContentBlock\nfrom lfx.schema.data import Data\nfrom lfx.schema.dotdict import dotdict\nfrom lfx.schema.message import Message\nfrom lfx.schema.table import EditMode\nfrom lfx.utils.constants import MESSAGE_SENDER_AI\n\n\ndef set_advanced_true(component_input):\n component_input.advanced = True\n return component_input\n\n\ndef _agent_base_inputs():\n \"\"\"Return base inputs tailored to AgentComponent's create_agent path.\n\n `get_base_inputs()` returns a shared list — replace, don't mutate. We drop\n inputs that are no-ops here and override info text on the inputs whose\n semantics shifted under create_agent.\n\n `verbose` is dropped because the create_agent event stream already surfaces\n every agent step via the \"Agent Steps\" content blocks; the legacy boolean\n has nothing to toggle. Saved flows that still carry a `verbose` value just\n ignore it on load (the schema no longer declares it).\n \"\"\"\n drop = {\"verbose\"}\n overrides = {\n \"handle_parsing_errors\": BoolInput(\n name=\"handle_parsing_errors\",\n display_name=\"Handle Parse Errors\",\n value=True,\n advanced=True,\n info=(\n \"Adds tool-execution retry as a safety net. `create_agent` already \"\n \"feeds tool-call validation errors back to the LLM automatically; \"\n \"this flag layers `ToolRetryMiddleware` on top so transient tool \"\n \"runtime failures are retried (max 2 retries).\"\n ),\n ),\n \"max_iterations\": IntInput(\n name=\"max_iterations\",\n display_name=\"Max Iterations\",\n value=15,\n advanced=True,\n range_spec=RangeSpec(min=1, max=128000, step=1, step_type=\"int\"),\n info=(\n \"Maximum number of model calls the agent can make before stopping \"\n \"(maps to `ModelCallLimitMiddleware.run_limit` on the create_agent \"\n \"path). Must be at least 1 — it is a safety cap, never 'unlimited'.\"\n ),\n ),\n }\n return [overrides.get(inp.name, inp) for inp in LCToolsAgentComponent.get_base_inputs() if inp.name not in drop]\n\n\ndef _extract_text_content(value) -> str:\n \"\"\"Pull a string payload from a Message-like, AIMessage-like, or string value.\"\"\"\n if isinstance(value, str):\n return value\n text = getattr(value, \"text\", None)\n if isinstance(text, str):\n return text\n content = getattr(value, \"content\", None)\n if isinstance(content, str):\n return content\n return str(value) if value is not None else \"\"\n\n\n@contextmanager\ndef _suppress_send_message(component: Any):\n \"\"\"Temporarily replace component.send_message with a no-op for the duration of the block.\n\n Used during the structured-output prompt fallback: run_agent streams the agent's\n final answer through self.send_message (correct for message_response), but in\n json_response the orchestrator parses that text into structured Data which the\n downstream Chat Output emits — leaving the original emission in place produces a\n duplicate message in the playground. The original method is always restored on exit,\n even when the wrapped call raises.\n \"\"\"\n original = component.send_message\n\n async def _noop(message, *_args, **_kwargs):\n return message\n\n component.send_message = _noop\n try:\n yield\n finally:\n component.send_message = original\n\n\nclass AgentComponent(ToolCallingAgentComponent):\n display_name: str = \"Agent\"\n description: str = \"Define the agent's instructions, then enter a task to complete using tools.\"\n documentation: str = \"https://docs.langflow.org/agents\"\n icon = \"bot\"\n beta = False\n name = \"Agent\"\n\n memory_inputs = [set_advanced_true(component_input) for component_input in MemoryComponent().inputs]\n\n inputs = [\n ModelInput(\n name=\"model\",\n display_name=\"Language Model\",\n info=\"Select your model provider\",\n real_time_refresh=True,\n required=True,\n # Agents require tool calling — the filter is honored by\n # ``handle_model_input_update`` so models that can't run with\n # tools never reach the picker (and any saved selection that\n # no longer satisfies the constraint is auto-replaced).\n filters={\"tool_calling\": True},\n ),\n SecretStrInput(\n name=\"api_key\",\n display_name=\"API Key\",\n info=\"Overrides global provider settings. Leave blank to use your pre-configured API Key.\",\n real_time_refresh=True,\n advanced=True,\n ),\n DropdownInput(\n name=\"base_url_ibm_watsonx\",\n display_name=\"watsonx API Endpoint\",\n info=\"The base URL of the API (IBM watsonx.ai only)\",\n options=IBM_WATSONX_URLS,\n value=IBM_WATSONX_URLS[0],\n combobox=True,\n show=False,\n real_time_refresh=True,\n ),\n StrInput(\n name=\"project_id\",\n display_name=\"watsonx Project ID\",\n info=\"The project ID associated with the foundation model (IBM watsonx.ai only)\",\n show=False,\n required=False,\n ),\n MultilineInput(\n name=\"system_prompt\",\n display_name=\"Agent Instructions\",\n info=(\n \"System Prompt: Initial instructions and context provided to guide the agent's behavior. \"\n \"Supports dynamic placeholders: {current_date}, {model_name}, {optional_user_context}.\"\n ),\n value=DEFAULT_SYSTEM_PROMPT_TEMPLATE,\n advanced=False,\n ),\n MessageTextInput(\n name=\"context_id\",\n display_name=\"Context ID\",\n info=\"The context ID of the chat. Adds an extra layer to the local memory.\",\n value=\"\",\n advanced=True,\n ),\n IntInput(\n name=\"n_messages\",\n display_name=\"Number of Chat History Messages\",\n value=100,\n info=\"Number of chat history messages to retrieve.\",\n advanced=True,\n show=True,\n ),\n IntInput(\n name=\"max_tokens\",\n display_name=\"Max Tokens\",\n info=\"Maximum number of tokens to generate. Field name varies by provider.\",\n advanced=True,\n range_spec=RangeSpec(min=1, max=128000, step=1, step_type=\"int\"),\n ),\n MultilineInput(\n name=\"format_instructions\",\n display_name=\"Output Format Instructions\",\n info=\"Generic Template for structured output formatting. Valid only with Structured response.\",\n value=(\n \"You are an AI that extracts structured JSON objects from unstructured text. \"\n \"Use a predefined schema with expected types (str, int, float, bool, dict). \"\n \"Extract ALL relevant instances that match the schema - if multiple patterns exist, capture them all. \"\n \"Fill missing or ambiguous values with defaults: null for missing values. \"\n \"Remove exact duplicates but keep variations that have different field values. \"\n \"Always return valid JSON in the expected format, never throw errors. \"\n \"If multiple objects can be extracted, return them all in the structured format.\"\n ),\n advanced=True,\n ),\n TableInput(\n name=\"output_schema\",\n display_name=\"Output Schema\",\n info=(\n \"Schema Validation: Define the structure and data types for structured output. \"\n \"No validation if no output schema.\"\n ),\n advanced=True,\n required=False,\n value=[],\n table_schema=[\n {\n \"name\": \"name\",\n \"display_name\": \"Name\",\n \"type\": \"str\",\n \"description\": \"Specify the name of the output field.\",\n \"default\": \"field\",\n \"edit_mode\": EditMode.INLINE,\n },\n {\n \"name\": \"description\",\n \"display_name\": \"Description\",\n \"type\": \"str\",\n \"description\": \"Describe the purpose of the output field.\",\n \"default\": \"description of field\",\n \"edit_mode\": EditMode.POPOVER,\n },\n {\n \"name\": \"type\",\n \"display_name\": \"Type\",\n \"type\": \"str\",\n \"edit_mode\": EditMode.INLINE,\n \"description\": (\"Indicate the data type of the output field (e.g., str, int, float, bool, dict).\"),\n \"options\": [\"str\", \"int\", \"float\", \"bool\", \"dict\"],\n \"default\": \"str\",\n },\n {\n \"name\": \"multiple\",\n \"display_name\": \"As List\",\n \"type\": \"boolean\",\n \"description\": \"Set to True if this output field should be a list of the specified type.\",\n \"default\": \"False\",\n \"edit_mode\": EditMode.INLINE,\n },\n ],\n ),\n *_agent_base_inputs(),\n # removed memory inputs from agent component\n # *memory_inputs,\n BoolInput(\n name=\"stream\",\n display_name=\"Stream\",\n info=STREAM_INFO_TEXT,\n value=True,\n advanced=True,\n ),\n BoolInput(\n name=\"add_current_date_tool\",\n display_name=\"Current Date\",\n advanced=True,\n info=\"If true, will add a tool to the agent that returns the current date.\",\n value=True,\n ),\n BoolInput(\n name=\"add_calculator_tool\",\n display_name=\"Calculator\",\n advanced=True,\n info=(\n \"If true, adds a zero-config arithmetic calculator tool to the agent \"\n \"(safe: only +, -, *, /, ** operators via AST).\"\n ),\n value=True,\n ),\n ]\n outputs = [\n Output(name=\"response\", display_name=\"Response\", method=\"message_response\"),\n Output(\n name=\"structured_response\",\n display_name=\"Structured Response\",\n method=\"json_response\",\n types=[\"Data\"],\n ),\n ]\n\n def _resolve_selected_model(self):\n \"\"\"Resolve the selected model, including legacy agent_llm/model_name inputs.\"\"\"\n try:\n from langchain_core.language_models import BaseLanguageModel\n\n if isinstance(self.model, BaseLanguageModel):\n return self.model\n except ImportError:\n pass\n\n if isinstance(self.model, list) and self.model:\n return self.model\n\n legacy_provider = getattr(self, \"agent_llm\", None)\n legacy_model_name = getattr(self, \"model_name\", None)\n if not legacy_provider or not legacy_model_name:\n return self.model\n\n options = get_language_model_options(user_id=self.user_id)\n for option in options:\n if option.get(\"provider\") == legacy_provider and option.get(\"name\") == legacy_model_name:\n return [option]\n\n return [\n {\n \"name\": legacy_model_name,\n \"provider\": legacy_provider,\n \"metadata\": {},\n }\n ]\n\n def _get_max_tokens_value(self):\n \"\"\"Return the user-supplied max_tokens or None when unset/zero.\"\"\"\n val = getattr(self, \"max_tokens\", None)\n if val in {\"\", 0}:\n return None\n return val\n\n def _get_llm(self):\n \"\"\"Override parent to include max_tokens from the Agent's input field.\n\n Streaming is mandatory for AgentComponent: ``runnable.astream_events(v2)`` only\n emits ``on_chat_model_stream`` chunks when the underlying chat model is\n instantiated with ``streaming=True``. Unlike the LanguageModel component (where\n ``stream`` is a user-facing toggle), the Agent has no opt-out — the toggle is\n kept in the UI for backwards compatibility but is intentionally ignored here.\n Without ``stream=True``, the chat model accumulates the whole response and\n only emits ``on_chat_model_end``, silently disabling the Playground's live-\n typing view and breaking the streaming contract on the /events surface.\n \"\"\"\n return get_llm(\n model=self.model,\n user_id=self.user_id,\n api_key=getattr(self, \"api_key\", None),\n stream=True,\n max_tokens=self._get_max_tokens_value(),\n watsonx_url=getattr(self, \"base_url_ibm_watsonx\", None),\n watsonx_project_id=getattr(self, \"project_id\", None),\n )\n\n async def get_agent_requirements(self):\n \"\"\"Get the agent requirements for the agent.\"\"\"\n from langchain_core.tools import StructuredTool\n\n selected_model = self._resolve_selected_model()\n try:\n from langchain_core.language_models import BaseLanguageModel\n\n is_connected_model = isinstance(selected_model, BaseLanguageModel)\n except ImportError:\n is_connected_model = False\n\n if not is_connected_model:\n validate_model_selection(selected_model)\n\n # Ensure _get_llm() uses the resolved model (e.g. from legacy agent_llm/model_name)\n self.model = selected_model\n llm_model = self._get_llm()\n if llm_model is None:\n msg = \"No language model selected. Please choose a model to proceed.\"\n raise ValueError(msg)\n\n # Get memory data\n self.chat_history = await self.get_memory_data()\n await logger.adebug(f\"Retrieved {len(self.chat_history)} chat history messages\")\n if isinstance(self.chat_history, Message):\n self.chat_history = [self.chat_history]\n\n # Add current date tool if enabled\n if self.add_current_date_tool:\n if not isinstance(self.tools, list): # type: ignore[has-type]\n self.tools = []\n current_date_tool = (await CurrentDateComponent(**self.get_base_args()).to_toolkit()).pop(0)\n\n if not isinstance(current_date_tool, StructuredTool):\n msg = \"CurrentDateComponent must be converted to a StructuredTool\"\n raise TypeError(msg)\n # Skip if an externally-connected tool already provides the same name.\n # Duplicate tool names are rejected by Anthropic/Gemini with HTTP 400.\n if not any(getattr(t, \"name\", None) == current_date_tool.name for t in self.tools):\n self.tools.append(current_date_tool)\n\n # Add calculator tool if enabled (zero-config arithmetic)\n if getattr(self, \"add_calculator_tool\", False):\n if not isinstance(self.tools, list): # type: ignore[has-type]\n self.tools = []\n calculator_tool = (await CalculatorComponent(**self.get_base_args()).to_toolkit()).pop(0)\n\n if not isinstance(calculator_tool, StructuredTool):\n msg = \"CalculatorComponent must be converted to a StructuredTool\"\n raise TypeError(msg)\n # Skip if an externally-connected tool already provides the same name.\n # Duplicate tool names are rejected by Anthropic/Gemini with HTTP 400.\n if not any(getattr(t, \"name\", None) == calculator_tool.name for t in self.tools):\n self.tools.append(calculator_tool)\n\n # Set shared callbacks for tracing the tools used by the agent\n self.set_tools_callbacks(self.tools, self._get_shared_callbacks())\n\n return llm_model, self.chat_history, self.tools\n\n def _get_resolved_model_name(self) -> str:\n \"\"\"Best-effort human-readable model name for {model_name} injection.\"\"\"\n try:\n from langchain_core.language_models import BaseLanguageModel\n\n if isinstance(self.model, BaseLanguageModel):\n return type(self.model).__name__\n except ImportError:\n pass\n\n if isinstance(self.model, list) and self.model:\n first = self.model[0]\n if isinstance(first, dict):\n name = first.get(\"name\")\n if isinstance(name, str) and name:\n return name\n\n legacy_model_name = getattr(self, \"model_name\", None)\n if isinstance(legacy_model_name, str) and legacy_model_name:\n return legacy_model_name\n return \"\"\n\n def _inject_dynamic_prompt_values(self, prompt: Any | None) -> str | None:\n \"\"\"Replace known env placeholders in the system prompt.\n\n Handles {current_date}, {model_name}, and {optional_user_context} (the\n last one ships with the structured DEFAULT_SYSTEM_PROMPT_TEMPLATE and\n is currently unused at the AgentComponent layer, so it resolves to \"\").\n Uses str.replace (not str.format) so user prompts containing literal\n braces such as JSON examples ({\"key\": 1}) never break the agent.\n\n `system_prompt` is a connectable MultilineInput, so the value can arrive\n as a Message (e.g. a Prompt node wired in). Normalize it to text first —\n a raw Message has no `.replace` and used to crash the agent build.\n \"\"\"\n if prompt is None:\n return None\n prompt = _extract_text_content(prompt)\n if not prompt:\n return prompt\n replacements = {\n \"{current_date}\": datetime.now(tz=timezone.utc).strftime(\"%Y-%m-%d %H:%M:%S UTC\"),\n \"{model_name}\": self._get_resolved_model_name(),\n \"{optional_user_context}\": \"\",\n }\n for placeholder, value in replacements.items():\n prompt = prompt.replace(placeholder, value)\n return prompt\n\n def create_agent_runnable(self):\n \"\"\"Build the LangGraph `CompiledStateGraph` via `langchain.agents.create_agent`.\n\n Replaces the legacy `AgentExecutor` runnable inherited from\n `ToolCallingAgentComponent`. Other agent components (tool_calling, csv, json,\n openapi, sql*, vector_store_router) keep the legacy path — only AgentComponent\n runs on the new graph API.\n\n `max_iterations` and `handle_parsing_errors` (legacy AgentExecutor knobs) are\n translated to LangGraph middleware. Without that translation those user inputs\n would silently become no-ops on the new API.\n\n Provider notes:\n - WatsonX/Granite work natively with create_agent — `ChatWatsonx.bind_tools`\n handles tool_choice correctly. The legacy `create_granite_agent` path was\n dropped because it hardcoded `tool_choice='required'`, which the WatsonX\n API now rejects.\n - Ollama and other small/local models often emit malformed tool args. The\n ToolRetryMiddleware (default `retry_on=(Exception,)`, `on_failure='continue'`)\n catches Pydantic ValidationErrors from bad args and feeds the error back\n to the LLM as a retry signal, so the agent recovers gracefully.\n \"\"\"\n llm = self._get_llm()\n tools = self.tools or []\n\n # Eager bind_tools validation. `create_agent(...)` is lazy — without this,\n # an LLM that doesn't support tool calling fails on the first user message\n # instead of when the user wires up the component, which is a much worse UX.\n # Gated on a non-empty tools list so a no-tool Agent on a plain chat model\n # (which legitimately has no `bind_tools`) isn't shut out at flow-build time.\n # Providers signal \"no tool calling\" inconsistently — `NotImplementedError`\n # (langchain default), `AttributeError` (no `bind_tools` attr), or `TypeError`\n # (signature mismatch). Treat all three as the same UX failure.\n if tools:\n try:\n llm.bind_tools(tools)\n except (NotImplementedError, AttributeError, TypeError) as exc:\n # Include the underlying error so a broken tool schema or a\n # provider implementation bug is not silently disguised as a\n # \"model can't call tools\" UX error.\n msg = (\n f\"{self.display_name} does not support tool calling, \"\n \"or one of the connected tools failed to bind. \"\n \"Please connect a tool-calling capable language model and \"\n f\"verify your tools are well-formed. Underlying error: {exc!s}\"\n )\n raise NotImplementedError(msg) from exc\n\n middleware = self._build_middleware(llm)\n return create_agent(\n model=llm,\n tools=tools,\n system_prompt=self.system_prompt or \"\",\n middleware=middleware or None,\n )\n\n def _compute_recursion_limit(self) -> int:\n \"\"\"Derive the LangGraph recursion_limit from the user-set max_iterations.\n\n Mirrors the clamp in `_build_middleware` (max(1, max_iterations)) so a\n saved 0 or negative value cannot under-cap the graph below one full\n iteration. The +5 buffer covers start/end/router overhead.\n \"\"\"\n raw = getattr(self, \"max_iterations\", None)\n run_limit = max(1, int(raw)) if raw is not None else 15\n return run_limit * 2 + 5\n\n def _build_middleware(self, llm: Any) -> list:\n # `llm` is passed in (rather than re-fetched via `self._get_llm()`)\n # because some providers do credential resolution / client instantiation\n # lazily on each call. The caller — `create_agent_runnable` — already\n # resolved it once for `bind_tools`, so reuse that instance here.\n middleware: list = []\n max_iterations = getattr(self, \"max_iterations\", None)\n if max_iterations is not None:\n # `max_iterations` is a safety cap, not an \"unlimited\" toggle. A saved\n # 0 or negative value (falsy) must NOT silently drop the limiter and\n # allow an unbounded model/tool loop — clamp it to a real minimum.\n run_limit = max(1, int(max_iterations))\n middleware.append(ModelCallLimitMiddleware(run_limit=run_limit))\n # ToolRetryMiddleware only matters when there ARE tools to retry. Attaching\n # it on a no-tools agent inflates the compiled graph and adds per-invocation\n # middleware overhead for nothing, which is a measurable contributor to\n # trivial-prompt latency (QA UI-003).\n if getattr(self, \"handle_parsing_errors\", False) and self.tools:\n middleware.append(ToolRetryMiddleware(max_retries=2))\n # WatsonX models have two known platform quirks; both still reproduce on\n # the current API, so we keep the protections from the legacy\n # `create_granite_agent` path.\n # 1. Multi-tool-call assistant turns are rejected (\"This model only\n # supports single tool-calls at once!\"). Clamp to one per turn.\n # 2. Tool args occasionally come back as literal placeholder strings\n # (e.g. ``). Re-invoke once with a corrective\n # SystemMessage.\n # Order: SingleToolCallMiddleware first (outermost) so the clamp is\n # applied to the final response, including any corrective re-invoke\n # produced by WatsonXPlaceholderMiddleware.\n if is_watsonx_model(llm):\n middleware.append(SingleToolCallMiddleware())\n middleware.append(WatsonXPlaceholderMiddleware())\n return middleware\n\n async def run_agent(self, agent) -> Message:\n \"\"\"Run the LangGraph `CompiledStateGraph` and return the final agent Message.\n\n Overrides the legacy `LCAgentComponent.run_agent` (which builds an\n `{\"input\": str, \"chat_history\": [...]}` dict for `AgentExecutor`). The graph\n wants `{\"messages\": [BaseMessage, ...]}`. The event stream is wrapped with\n `adapt_graph_events_to_executor_shape` so the legacy `process_agent_events`\n (in `lfx.base.agents.events`) can be reused unchanged.\n \"\"\"\n messages = build_initial_messages(\n input_value=self.input_value,\n chat_history=getattr(self, \"chat_history\", None),\n )\n input_dict = {\"messages\": messages}\n\n agent_message = self._build_initial_agent_message()\n token_usage_handler = TokenUsageCallbackHandler()\n\n # Stream tokens to the event manager when running inside the Langflow runtime.\n # This is what powers the live-typing view in the chat UI.\n on_token_callback: OnTokenFunctionType | None = None\n if getattr(self, \"_event_manager\", None):\n on_token_callback = cast(\"OnTokenFunctionType\", self._event_manager.on_token)\n\n # Align LangGraph's `recursion_limit` with `max_iterations` so the\n # middleware cap (ModelCallLimitMiddleware) is what bounds the loop —\n # not LangGraph's default 25-step guard, which fires at ~12 model+tool\n # iterations and raises a raw GraphRecursionError (QA UI-009/UI-010).\n # Each iteration is ~2 graph steps (model node + tools node); add 5\n # for start/end overhead.\n recursion_limit = self._compute_recursion_limit()\n\n stream = adapt_graph_events_to_executor_shape(\n agent.astream_events(\n input_dict,\n config={\n \"callbacks\": [\n AgentAsyncHandler(self.log),\n token_usage_handler,\n *self._get_shared_callbacks(),\n ],\n \"recursion_limit\": recursion_limit,\n },\n version=\"v2\",\n )\n )\n try:\n result = await process_agent_events(\n stream,\n agent_message,\n cast(\"SendMessageFunctionType\", self.send_message),\n on_token_callback,\n )\n except ExceptionWithMessageError as e:\n # Drop the half-stored partial message from the DB (only if it was\n # actually persisted) and tell the frontend to remove the stale bubble.\n if hasattr(e, \"agent_message\"):\n msg_id = e.agent_message.get_id()\n if msg_id:\n await delete_message(id_=msg_id)\n await self._send_message_event(e.agent_message, category=\"remove_message\")\n logger.error(f\"ExceptionWithMessageError: {e}\")\n raise\n\n usage_data = token_usage_handler.get_usage()\n if usage_data:\n self._token_usage = usage_data\n result.properties.usage = usage_data\n # Only round-trip the DB when the message was stored (Chat Output wired).\n # `_should_skip_message=True` leaves `result.get_id()` empty; persisting\n # then would create a phantom row.\n if result.get_id():\n stored_result = await self._update_stored_message(result)\n await self._send_message_event(stored_result)\n result = stored_result\n\n self.status = result\n return result\n\n def _build_initial_agent_message(self) -> Message:\n \"\"\"Construct the placeholder agent Message that `process_agent_events` mutates.\"\"\"\n if hasattr(self, \"graph\"):\n session_id = self.graph.session_id\n elif hasattr(self, \"_session_id\"):\n session_id = self._session_id\n else:\n session_id = None\n\n sender_name = get_chat_output_sender_name(self) or self.display_name or \"AI\"\n return Message(\n sender=MESSAGE_SENDER_AI,\n sender_name=sender_name,\n properties={\"icon\": \"Bot\", \"state\": \"partial\"},\n content_blocks=[ContentBlock(title=\"Agent Steps\", contents=[])],\n session_id=session_id or uuid.uuid4(),\n )\n\n async def message_response(self) -> Message:\n try:\n llm_model, self.chat_history, self.tools = await self.get_agent_requirements()\n # Set up and run agent\n self.set(\n llm=llm_model,\n tools=self.tools or [],\n chat_history=self.chat_history,\n input_value=self.input_value,\n system_prompt=self._inject_dynamic_prompt_values(self.system_prompt),\n )\n agent = self.create_agent_runnable()\n result = await self.run_agent(agent)\n\n # Store result for potential JSON output\n self._agent_result = result\n\n except (ValueError, TypeError, KeyError) as e:\n await logger.aerror(f\"{type(e).__name__}: {e!s}\")\n raise\n except ExceptionWithMessageError as e:\n await logger.aerror(f\"ExceptionWithMessageError occurred: {e}\")\n raise\n # Avoid catching blind Exception; let truly unexpected exceptions propagate\n except Exception as e:\n await logger.aerror(f\"Unexpected error: {e!s}\")\n raise\n else:\n return result\n\n async def json_response(self) -> Data:\n \"\"\"Produce structured Data via native LLM structured output, with prompt-based fallback.\n\n Native path (no tools, llm has with_structured_output) bypasses the agent loop and\n returns provider-validated JSON. When tools are attached, falls back to running the\n agent with a schema-augmented system prompt and parsing the final message content.\n \"\"\"\n from lfx.components.models_and_agents.structured_output.structured_output_orchestrator import (\n orchestrate_structured_output,\n )\n\n try:\n llm_model, self.chat_history, self.tools = await self.get_agent_requirements()\n except (ValueError, TypeError) as exc:\n await logger.aerror(f\"json_response.requirements_failed: {exc}\")\n return Data(data={\"content\": \"\", \"error\": str(exc)})\n\n injected_system_prompt = self._inject_dynamic_prompt_values(getattr(self, \"system_prompt\", \"\") or \"\") or \"\"\n format_instructions = getattr(self, \"format_instructions\", \"\") or \"\"\n output_schema = getattr(self, \"output_schema\", None) or []\n has_tools = bool(self.tools)\n\n async def _run_agent_for_fallback(augmented_prompt: str) -> str:\n self.set(\n llm=llm_model,\n tools=self.tools or [],\n chat_history=self.chat_history,\n input_value=self.input_value,\n system_prompt=augmented_prompt,\n )\n agent_runnable = self.create_agent_runnable()\n with _suppress_send_message(self):\n result = await self.run_agent(agent_runnable)\n return _extract_text_content(result)\n\n try:\n return await orchestrate_structured_output(\n llm=llm_model,\n output_schema=output_schema,\n system_prompt=injected_system_prompt,\n format_instructions=format_instructions,\n input_value=_extract_text_content(self.input_value),\n run_prompt_fallback=_run_agent_for_fallback,\n prefer_native=not has_tools,\n )\n except (\n ExceptionWithMessageError,\n ValueError,\n TypeError,\n NotImplementedError,\n AttributeError,\n ) as exc:\n await logger.aerror(f\"json_response.orchestration_failed: {exc}\")\n return Data(data={\"content\": \"\", \"error\": str(exc)})\n\n async def get_memory_data(self):\n # Scope by flow_id so default playground session names (e.g. \"New Session 0\")\n # cannot leak chat history across unrelated flows. See issue #13059.\n # The helper also returns [] when n_messages == 0, preserving the\n # explicit \"memory disabled\" contract from MemoryComponent.retrieve_messages.\n messages = await aget_agent_chat_history(\n session_id=self.graph.session_id,\n flow_id=getattr(self.graph, \"flow_id\", None),\n context_id=self.context_id,\n n_messages=self.n_messages,\n )\n return [\n message for message in messages if getattr(message, \"id\", None) != getattr(self.input_value, \"id\", None)\n ]\n\n def update_input_types(self, build_config: dotdict) -> dotdict:\n \"\"\"Update input types for all fields in build_config.\"\"\"\n for key, value in build_config.items():\n if isinstance(value, dict):\n if value.get(\"input_types\") is None:\n build_config[key][\"input_types\"] = []\n elif hasattr(value, \"input_types\") and value.input_types is None:\n value.input_types = []\n return build_config\n\n async def update_build_config(\n self,\n build_config: dotdict,\n field_value: list[dict],\n field_name: str | None = None,\n ) -> dotdict:\n # Update model options with caching (for all field changes).\n # The tool-calling constraint lives on the ModelInput's ``filters``\n # field (declared above); ``handle_model_input_update`` reads it\n # and applies the filter to both the dropdown options and the\n # sticky-default re-injection path.\n build_config = handle_model_input_update(\n component=self,\n build_config=dict(build_config),\n field_value=field_value,\n field_name=field_name,\n )\n build_config = dotdict(build_config)\n\n if field_name == \"model\":\n build_config = self.update_input_types(build_config)\n\n # Validate required keys. `verbose` was dropped from the input set\n # (see `_agent_base_inputs` — the create_agent event stream already\n # surfaces every step), so it is intentionally NOT required here.\n # Saved flows that still carry a `verbose` value just ignore it on\n # load.\n default_keys = [\n \"code\",\n \"_type\",\n \"model\",\n \"tools\",\n \"input_value\",\n \"add_current_date_tool\",\n \"add_calculator_tool\",\n \"system_prompt\",\n \"max_iterations\",\n \"handle_parsing_errors\",\n ]\n missing_keys = [key for key in default_keys if key not in build_config]\n if missing_keys:\n msg = f\"Missing required keys in build_config: {missing_keys}\"\n raise ValueError(msg)\n return dotdict({k: v.to_dict() if hasattr(v, \"to_dict\") else v for k, v in build_config.items()})\n\n async def _get_tools(self) -> list[Tool]:\n component_toolkit = get_component_toolkit()\n\n tools = component_toolkit(component=self).get_tools(\n tool_name=\"Call_Agent\",\n # here we do not use the shared callbacks as we are exposing the agent as a tool\n callbacks=self.get_langchain_callbacks(),\n )\n if hasattr(self, \"tools_metadata\"):\n tools = component_toolkit(component=self, metadata=self.tools_metadata).update_tools_metadata(tools=tools)\n\n return tools\n" }, "context_id": { "_input_type": "MessageTextInput", @@ -118454,6 +118454,6 @@ "num_components": 354, "num_modules": 95 }, - "sha256": "544fc8ca7b56f49e13681bf74be7acdc2cbbd4aa5db57f8206026d88ccfb03f0", + "sha256": "57420a9e8822a72ba8e8f8fd0d597f5ff7b9d4f85d965ed40713d108512cc7fe", "version": "1.11.0" }