diff --git a/src/backend/base/langflow/__main__.py b/src/backend/base/langflow/__main__.py index a1932baed..381628bc3 100644 --- a/src/backend/base/langflow/__main__.py +++ b/src/backend/base/langflow/__main__.py @@ -415,7 +415,7 @@ def print_banner(host: str, port: int, protocol: str) -> None: "To contribute, set: [bold]DO_NOT_TRACK=false[/bold] in your environment." ) ) - access_host = host if host != "0.0.0.0" else "localhost" + access_host = host if host != "0.0.0.0" else "localhost" # noqa: S104 access_link = f"[bold]🟢 Open Langflow →[/bold] [link={protocol}://{access_host}:{port}]{protocol}://{access_host}:{port}[/link]" message = f"{title}\n{info_text}\n\n{telemetry_text}\n\n{access_link}" diff --git a/src/backend/base/langflow/events/event_manager.py b/src/backend/base/langflow/events/event_manager.py index e499030d9..9d879809e 100644 --- a/src/backend/base/langflow/events/event_manager.py +++ b/src/backend/base/langflow/events/event_manager.py @@ -5,7 +5,7 @@ import json import time import uuid from functools import partial -from typing import TYPE_CHECKING, Literal +from typing import TYPE_CHECKING from fastapi.encoders import jsonable_encoder from loguru import logger @@ -50,7 +50,7 @@ class EventManager: def register_event( self, name: str, - event_type: Literal["message", "error", "warning", "info", "token"], + event_type: str, callback: EventCallback | None = None, ) -> None: if not name: @@ -65,7 +65,7 @@ class EventManager: callback_ = partial(callback, manager=self, event_type=event_type) self.events[name] = callback_ - def send_event(self, *, event_type: Literal["message", "error", "warning", "info", "token"], data: LoggableType): + def send_event(self, *, event_type: str, data: LoggableType): try: if isinstance(data, dict) and event_type in {"message", "error", "warning", "info", "token"}: data = create_event_by_type(event_type, **data) diff --git a/src/backend/base/langflow/initial_setup/starter_projects/Financial Report Parser.json b/src/backend/base/langflow/initial_setup/starter_projects/Financial Report Parser.json index 23d133b29..84cad8de9 100644 --- a/src/backend/base/langflow/initial_setup/starter_projects/Financial Report Parser.json +++ b/src/backend/base/langflow/initial_setup/starter_projects/Financial Report Parser.json @@ -512,7 +512,174 @@ }, { "data": { - "id": "ChatOutput-meCWe", + "id": "ParseData-hLbU0", + "node": { + "base_classes": [ + "Data", + "Message" + ], + "beta": false, + "category": "processing", + "conditional_paths": [], + "custom_fields": {}, + "description": "Convert Data objects into Messages using any {field_name} from input data.", + "display_name": "Data to Message", + "documentation": "", + "edited": false, + "field_order": [ + "data", + "template", + "sep" + ], + "frozen": false, + "icon": "message-square", + "key": "ParseData", + "legacy": true, + "lf_version": "1.1.5", + "metadata": { + "legacy_name": "Parse Data" + }, + "minimized": false, + "output_types": [], + "outputs": [ + { + "allows_loop": false, + "cache": true, + "display_name": "Message", + "group_outputs": false, + "method": "parse_data", + "name": "text", + "selected": "Message", + "tool_mode": true, + "types": [ + "Message" + ], + "value": "__UNDEFINED__" + }, + { + "allows_loop": false, + "cache": true, + "display_name": "Data List", + "group_outputs": false, + "method": "parse_data_as_list", + "name": "data_list", + "selected": "Data", + "tool_mode": true, + "types": [ + "Data" + ], + "value": "__UNDEFINED__" + } + ], + "pinned": false, + "score": 0.23285358167685585, + "template": { + "_type": "Component", + "code": { + "advanced": true, + "dynamic": true, + "fileTypes": [], + "file_path": "", + "info": "", + "list": false, + "load_from_db": false, + "multiline": true, + "name": "code", + "password": false, + "placeholder": "", + "required": true, + "show": true, + "title_case": false, + "type": "code", + "value": "from langflow.custom import Component\nfrom langflow.helpers.data import data_to_text, data_to_text_list\nfrom langflow.io import DataInput, MultilineInput, Output, StrInput\nfrom langflow.schema import Data\nfrom langflow.schema.message import Message\n\n\nclass ParseDataComponent(Component):\n display_name = \"Data to Message\"\n description = \"Convert Data objects into Messages using any {field_name} from input data.\"\n icon = \"message-square\"\n name = \"ParseData\"\n legacy = True\n metadata = {\n \"legacy_name\": \"Parse Data\",\n }\n\n inputs = [\n DataInput(\n name=\"data\",\n display_name=\"Data\",\n info=\"The data to convert to text.\",\n is_list=True,\n required=True,\n ),\n MultilineInput(\n name=\"template\",\n display_name=\"Template\",\n info=\"The template to use for formatting the data. \"\n \"It can contain the keys {text}, {data} or any other key in the Data.\",\n value=\"{text}\",\n required=True,\n ),\n StrInput(name=\"sep\", display_name=\"Separator\", advanced=True, value=\"\\n\"),\n ]\n\n outputs = [\n Output(\n display_name=\"Message\",\n name=\"text\",\n info=\"Data as a single Message, with each input Data separated by Separator\",\n method=\"parse_data\",\n ),\n Output(\n display_name=\"Data List\",\n name=\"data_list\",\n info=\"Data as a list of new Data, each having `text` formatted by Template\",\n method=\"parse_data_as_list\",\n ),\n ]\n\n def _clean_args(self) -> tuple[list[Data], str, str]:\n data = self.data if isinstance(self.data, list) else [self.data]\n template = self.template\n sep = self.sep\n return data, template, sep\n\n def parse_data(self) -> Message:\n data, template, sep = self._clean_args()\n result_string = data_to_text(template, data, sep)\n self.status = result_string\n return Message(text=result_string)\n\n def parse_data_as_list(self) -> list[Data]:\n data, template, _ = self._clean_args()\n text_list, data_list = data_to_text_list(template, data)\n for item, text in zip(data_list, text_list, strict=True):\n item.set_text(text)\n self.status = data_list\n return data_list\n" + }, + "data": { + "_input_type": "DataInput", + "advanced": false, + "display_name": "Data", + "dynamic": false, + "info": "The data to convert to text.", + "input_types": [ + "Data" + ], + "list": true, + "list_add_label": "Add More", + "name": "data", + "placeholder": "", + "required": true, + "show": true, + "title_case": false, + "tool_mode": false, + "trace_as_input": true, + "trace_as_metadata": true, + "type": "other", + "value": "" + }, + "sep": { + "_input_type": "StrInput", + "advanced": true, + "display_name": "Separator", + "dynamic": false, + "info": "", + "list": false, + "list_add_label": "Add More", + "load_from_db": false, + "name": "sep", + "placeholder": "", + "required": false, + "show": true, + "title_case": false, + "tool_mode": false, + "trace_as_metadata": true, + "type": "str", + "value": "\n" + }, + "template": { + "_input_type": "MultilineInput", + "advanced": false, + "display_name": "Template", + "dynamic": false, + "info": "The template to use for formatting the data. It can contain the keys {text}, {data} or any other key in the Data.", + "input_types": [ + "Message" + ], + "list": false, + "list_add_label": "Add More", + "load_from_db": false, + "multiline": true, + "name": "template", + "placeholder": "", + "required": true, + "show": true, + "title_case": false, + "tool_mode": false, + "trace_as_input": true, + "trace_as_metadata": true, + "type": "str", + "value": "EBITIDA: {EBITIDA} , \nNet Income: {NET_INCOME} ,\nGROSS_PROFIT: {GROSS_PROFIT}" + } + }, + "tool_mode": false + }, + "showNode": true, + "type": "ParseData" + }, + "dragging": false, + "id": "ParseData-hLbU0", + "measured": { + "height": 342, + "width": 320 + }, + "position": { + "x": 1788.8753184818502, + "y": 247.15776278878775 + }, + "selected": false, + "type": "genericNode" + }, + { + "data": { + "id": "ChatOutput-ZtLkt", "node": { "base_classes": [ "Message" diff --git a/src/backend/base/langflow/initial_setup/starter_projects/Meeting Summary.json b/src/backend/base/langflow/initial_setup/starter_projects/Meeting Summary.json index 676d739b5..3e5d29c40 100644 --- a/src/backend/base/langflow/initial_setup/starter_projects/Meeting Summary.json +++ b/src/backend/base/langflow/initial_setup/starter_projects/Meeting Summary.json @@ -432,7 +432,170 @@ }, { "data": { - "id": "OpenAIModel-ac6TO", + "id": "ParseData-LUfjb", + "node": { + "base_classes": [ + "Data", + "Message" + ], + "beta": false, + "conditional_paths": [], + "custom_fields": {}, + "description": "Convert Data objects into Messages using any {field_name} from input data.", + "display_name": "Data to Message", + "documentation": "", + "edited": false, + "field_order": [ + "data", + "template", + "sep" + ], + "frozen": false, + "icon": "message-square", + "legacy": true, + "lf_version": "1.1.5", + "metadata": { + "legacy_name": "Parse Data" + }, + "minimized": false, + "output_types": [], + "outputs": [ + { + "allows_loop": false, + "cache": true, + "display_name": "Message", + "group_outputs": false, + "method": "parse_data", + "name": "text", + "selected": "Message", + "tool_mode": true, + "types": [ + "Message" + ], + "value": "__UNDEFINED__" + }, + { + "allows_loop": false, + "cache": true, + "display_name": "Data List", + "group_outputs": false, + "method": "parse_data_as_list", + "name": "data_list", + "selected": "Data", + "tool_mode": true, + "types": [ + "Data" + ], + "value": "__UNDEFINED__" + } + ], + "pinned": false, + "template": { + "_type": "Component", + "code": { + "advanced": true, + "dynamic": true, + "fileTypes": [], + "file_path": "", + "info": "", + "list": false, + "load_from_db": false, + "multiline": true, + "name": "code", + "password": false, + "placeholder": "", + "required": true, + "show": true, + "title_case": false, + "type": "code", + "value": "from langflow.custom import Component\nfrom langflow.helpers.data import data_to_text, data_to_text_list\nfrom langflow.io import DataInput, MultilineInput, Output, StrInput\nfrom langflow.schema import Data\nfrom langflow.schema.message import Message\n\n\nclass ParseDataComponent(Component):\n display_name = \"Data to Message\"\n description = \"Convert Data objects into Messages using any {field_name} from input data.\"\n icon = \"message-square\"\n name = \"ParseData\"\n legacy = True\n metadata = {\n \"legacy_name\": \"Parse Data\",\n }\n\n inputs = [\n DataInput(\n name=\"data\",\n display_name=\"Data\",\n info=\"The data to convert to text.\",\n is_list=True,\n required=True,\n ),\n MultilineInput(\n name=\"template\",\n display_name=\"Template\",\n info=\"The template to use for formatting the data. \"\n \"It can contain the keys {text}, {data} or any other key in the Data.\",\n value=\"{text}\",\n required=True,\n ),\n StrInput(name=\"sep\", display_name=\"Separator\", advanced=True, value=\"\\n\"),\n ]\n\n outputs = [\n Output(\n display_name=\"Message\",\n name=\"text\",\n info=\"Data as a single Message, with each input Data separated by Separator\",\n method=\"parse_data\",\n ),\n Output(\n display_name=\"Data List\",\n name=\"data_list\",\n info=\"Data as a list of new Data, each having `text` formatted by Template\",\n method=\"parse_data_as_list\",\n ),\n ]\n\n def _clean_args(self) -> tuple[list[Data], str, str]:\n data = self.data if isinstance(self.data, list) else [self.data]\n template = self.template\n sep = self.sep\n return data, template, sep\n\n def parse_data(self) -> Message:\n data, template, sep = self._clean_args()\n result_string = data_to_text(template, data, sep)\n self.status = result_string\n return Message(text=result_string)\n\n def parse_data_as_list(self) -> list[Data]:\n data, template, _ = self._clean_args()\n text_list, data_list = data_to_text_list(template, data)\n for item, text in zip(data_list, text_list, strict=True):\n item.set_text(text)\n self.status = data_list\n return data_list\n" + }, + "data": { + "_input_type": "DataInput", + "advanced": false, + "display_name": "Data", + "dynamic": false, + "info": "The data to convert to text.", + "input_types": [ + "Data" + ], + "list": true, + "list_add_label": "Add More", + "name": "data", + "placeholder": "", + "required": true, + "show": true, + "title_case": false, + "tool_mode": false, + "trace_as_input": true, + "trace_as_metadata": true, + "type": "other", + "value": "" + }, + "sep": { + "_input_type": "StrInput", + "advanced": true, + "display_name": "Separator", + "dynamic": false, + "info": "", + "list": false, + "list_add_label": "Add More", + "load_from_db": false, + "name": "sep", + "placeholder": "", + "required": false, + "show": true, + "title_case": false, + "tool_mode": false, + "trace_as_metadata": true, + "type": "str", + "value": "\n" + }, + "template": { + "_input_type": "MultilineInput", + "advanced": false, + "display_name": "Template", + "dynamic": false, + "info": "The template to use for formatting the data. It can contain the keys {text}, {data} or any other key in the Data.", + "input_types": [ + "Message" + ], + "list": false, + "list_add_label": "Add More", + "load_from_db": false, + "multiline": true, + "name": "template", + "placeholder": "", + "required": true, + "show": true, + "title_case": false, + "tool_mode": false, + "trace_as_input": true, + "trace_as_metadata": true, + "type": "str", + "value": "{text}" + } + }, + "tool_mode": false + }, + "showNode": true, + "type": "ParseData" + }, + "id": "ParseData-LUfjb", + "measured": { + "height": 342, + "width": 320 + }, + "position": { + "x": 1330.927281184057, + "y": 382.3516758942169 + }, + "selected": false, + "type": "genericNode" + }, + { + "data": { + "id": "OpenAIModel-iudDZ", "node": { "base_classes": [ "LanguageModel", diff --git a/src/backend/base/langflow/schema/playground_events.py b/src/backend/base/langflow/schema/playground_events.py index d7db91ddb..de58af959 100644 --- a/src/backend/base/langflow/schema/playground_events.py +++ b/src/backend/base/langflow/schema/playground_events.py @@ -170,12 +170,13 @@ _EVENT_CREATORS: dict[str, tuple[Callable, inspect.Signature]] = { } -def create_event_by_type( - event_type: Literal["message", "error", "warning", "info", "token"], **kwargs -) -> PlaygroundEvent | dict: +def create_event_by_type(event_type: str, **kwargs) -> PlaygroundEvent | dict: if event_type not in _EVENT_CREATORS: return kwargs - - creator_func, signature = _EVENT_CREATORS[event_type] + try: + creator_func, signature = _EVENT_CREATORS[event_type] + except KeyError as e: + msg = f"Invalid event type: {event_type}" + raise ValueError(msg) from e valid_params = {k: v for k, v in kwargs.items() if k in signature.parameters} return creator_func(**valid_params) diff --git a/src/backend/base/langflow/services/job_queue/service.py b/src/backend/base/langflow/services/job_queue/service.py index a05a94a51..cee9cbb0b 100644 --- a/src/backend/base/langflow/services/job_queue/service.py +++ b/src/backend/base/langflow/services/job_queue/service.py @@ -4,7 +4,7 @@ import asyncio from loguru import logger -from langflow.events.event_manager import EventManager, create_default_event_manager +from langflow.events.event_manager import EventManager from langflow.services.base import Service @@ -133,18 +133,17 @@ class JobQueueService(Service): - The asyncio.Queue instance for handling the job's tasks or messages. - The EventManager instance for event handling tied to the queue. """ - if job_id in self._queues: - msg = f"Queue for job_id {job_id} already exists" - logger.error(msg) - raise ValueError(msg) - if self._closed: msg = "Queue service is closed" - logger.error(msg) raise RuntimeError(msg) + existing_queue = self._queues.get(job_id) + if existing_queue: + msg = f"Queue for job_id {job_id} already exists" + raise ValueError(msg) + main_queue: asyncio.Queue = asyncio.Queue() - event_manager = create_default_event_manager(main_queue) + event_manager: EventManager = self._create_default_event_manager(main_queue) # Register the queue without an active task. self._queues[job_id] = (main_queue, event_manager, None, None) @@ -300,3 +299,29 @@ class JobQueueService(Service): # Enough time has passed, perform the actual cleanup logger.debug(f"Cleaning up job_id {job_id} after grace period") await self.cleanup_job(job_id) + + def _create_default_event_manager(self, queue: asyncio.Queue) -> EventManager: + """Creates the default event manager with predefined events. + + Args: + queue (asyncio.Queue): The queue to be associated with the event manager. + + Returns: + EventManager: The configured EventManager instance. + """ + manager = EventManager(queue) + # Registering predefined events + event_names_types = [ + ("on_token", "token"), + ("on_vertices_sorted", "vertices_sorted"), + ("on_error", "error"), + ("on_end", "end"), + ("on_message", "add_message"), + ("on_remove_message", "remove_message"), + ("on_end_vertex", "end_vertex"), + ("on_build_start", "build_start"), + ("on_build_end", "build_end"), + ] + for name, event_type in event_names_types: + manager.register_event(name, event_type) + return manager