Merge branch 'zustand/io/migration' of personal:logspace-ai/langflow into zustand/io/migration
This commit is contained in:
commit
5c5883f58f
4 changed files with 60 additions and 26 deletions
|
|
@ -185,15 +185,16 @@ async def build_vertex(
|
||||||
chat_service.clear_cache(flow_id)
|
chat_service.clear_cache(flow_id)
|
||||||
|
|
||||||
# Log the vertex build
|
# Log the vertex build
|
||||||
background_tasks.add_task(
|
if not vertex.will_stream:
|
||||||
log_vertex_build,
|
background_tasks.add_task(
|
||||||
flow_id=flow_id,
|
log_vertex_build,
|
||||||
vertex_id=vertex_id,
|
flow_id=flow_id,
|
||||||
valid=valid,
|
vertex_id=vertex_id,
|
||||||
params=params,
|
valid=valid,
|
||||||
data=result_data_response,
|
params=params,
|
||||||
artifacts=artifacts,
|
data=result_data_response,
|
||||||
)
|
artifacts=artifacts,
|
||||||
|
)
|
||||||
|
|
||||||
timedelta = time.perf_counter() - start_time
|
timedelta = time.perf_counter() - start_time
|
||||||
duration = format_elapsed_time(timedelta)
|
duration = format_elapsed_time(timedelta)
|
||||||
|
|
@ -243,22 +244,31 @@ async def build_vertex_stream(
|
||||||
vertex: "ChatVertex" = graph.get_vertex(vertex_id)
|
vertex: "ChatVertex" = graph.get_vertex(vertex_id)
|
||||||
if not hasattr(vertex, "stream"):
|
if not hasattr(vertex, "stream"):
|
||||||
raise ValueError(f"Vertex {vertex_id} does not support streaming")
|
raise ValueError(f"Vertex {vertex_id} does not support streaming")
|
||||||
if not vertex.pinned or not vertex._built:
|
if isinstance(vertex._built_result, str) and vertex._built_result:
|
||||||
|
stream_data = StreamData(
|
||||||
|
event="message",
|
||||||
|
data={"message": f"Streaming vertex {vertex_id}"},
|
||||||
|
)
|
||||||
|
yield str(stream_data)
|
||||||
|
stream_data = StreamData(
|
||||||
|
event="message",
|
||||||
|
data={"chunk": vertex._built_result},
|
||||||
|
)
|
||||||
|
yield str(stream_data)
|
||||||
|
|
||||||
|
elif not vertex.pinned or not vertex._built:
|
||||||
logger.debug(f"Streaming vertex {vertex_id}")
|
logger.debug(f"Streaming vertex {vertex_id}")
|
||||||
stream_data = StreamData(
|
stream_data = StreamData(
|
||||||
event="message",
|
event="message",
|
||||||
data={"message": f"Streaming vertex {vertex_id}"},
|
data={"message": f"Streaming vertex {vertex_id}"},
|
||||||
)
|
)
|
||||||
yield str(stream_data)
|
yield str(stream_data)
|
||||||
number_of_chunks = 0
|
|
||||||
async for chunk in vertex.stream():
|
async for chunk in vertex.stream():
|
||||||
stream_data = StreamData(
|
stream_data = StreamData(
|
||||||
event="message",
|
event="message",
|
||||||
data={"chunk": chunk},
|
data={"chunk": chunk},
|
||||||
)
|
)
|
||||||
number_of_chunks += 1
|
|
||||||
yield str(stream_data)
|
yield str(stream_data)
|
||||||
logger.debug(f"Number of chunks: {number_of_chunks}")
|
|
||||||
elif vertex.result is not None:
|
elif vertex.result is not None:
|
||||||
stream_data = StreamData(
|
stream_data = StreamData(
|
||||||
event="message",
|
event="message",
|
||||||
|
|
|
||||||
|
|
@ -2,16 +2,13 @@ import ast
|
||||||
import inspect
|
import inspect
|
||||||
import types
|
import types
|
||||||
from enum import Enum
|
from enum import Enum
|
||||||
from typing import TYPE_CHECKING, Any, Callable, Coroutine, Dict, List, Optional
|
from typing import (TYPE_CHECKING, Any, Callable, Coroutine, Dict, List,
|
||||||
|
Optional)
|
||||||
|
|
||||||
from loguru import logger
|
from loguru import logger
|
||||||
|
|
||||||
from langflow.graph.schema import (
|
from langflow.graph.schema import (INPUT_COMPONENTS, OUTPUT_COMPONENTS,
|
||||||
INPUT_COMPONENTS,
|
InterfaceComponentTypes, ResultData)
|
||||||
OUTPUT_COMPONENTS,
|
|
||||||
InterfaceComponentTypes,
|
|
||||||
ResultData,
|
|
||||||
)
|
|
||||||
from langflow.graph.utils import UnbuiltObject, UnbuiltResult
|
from langflow.graph.utils import UnbuiltObject, UnbuiltResult
|
||||||
from langflow.graph.vertex.utils import generate_result
|
from langflow.graph.vertex.utils import generate_result
|
||||||
from langflow.interface.initialize import loading
|
from langflow.interface.initialize import loading
|
||||||
|
|
@ -44,6 +41,7 @@ class Vertex:
|
||||||
) -> None:
|
) -> None:
|
||||||
# is_external means that the Vertex send or receives data from
|
# is_external means that the Vertex send or receives data from
|
||||||
# an external source (e.g the chat)
|
# an external source (e.g the chat)
|
||||||
|
self.will_stream = False
|
||||||
self.updated_raw_params = False
|
self.updated_raw_params = False
|
||||||
self.id: str = data["id"]
|
self.id: str = data["id"]
|
||||||
self.is_input = any(
|
self.is_input = any(
|
||||||
|
|
|
||||||
|
|
@ -11,7 +11,7 @@ from langflow.graph.utils import UnbuiltObject, flatten_list
|
||||||
from langflow.graph.vertex.base import StatefulVertex, StatelessVertex
|
from langflow.graph.vertex.base import StatefulVertex, StatelessVertex
|
||||||
from langflow.interface.utils import extract_input_variables_from_prompt
|
from langflow.interface.utils import extract_input_variables_from_prompt
|
||||||
from langflow.schema import Record
|
from langflow.schema import Record
|
||||||
from langflow.services.monitor.utils import log_message
|
from langflow.services.monitor.utils import log_vertex_build
|
||||||
from langflow.utils.schemas import ChatOutputResponse
|
from langflow.utils.schemas import ChatOutputResponse
|
||||||
|
|
||||||
|
|
||||||
|
|
@ -394,6 +394,8 @@ class ChatVertex(StatelessVertex):
|
||||||
sender_name=sender_name,
|
sender_name=sender_name,
|
||||||
stream_url=stream_url,
|
stream_url=stream_url,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
self.will_stream = stream_url is not None
|
||||||
if artifacts:
|
if artifacts:
|
||||||
self.artifacts = artifacts.model_dump()
|
self.artifacts = artifacts.model_dump()
|
||||||
if isinstance(self._built_object, (AsyncIterator, Iterator)):
|
if isinstance(self._built_object, (AsyncIterator, Iterator)):
|
||||||
|
|
@ -434,13 +436,15 @@ class ChatVertex(StatelessVertex):
|
||||||
self._built_result = complete_message
|
self._built_result = complete_message
|
||||||
# Update artifacts with the message
|
# Update artifacts with the message
|
||||||
# and remove the stream_url
|
# and remove the stream_url
|
||||||
|
self._finalize_build()
|
||||||
logger.debug(f"Streamed message: {complete_message}")
|
logger.debug(f"Streamed message: {complete_message}")
|
||||||
|
|
||||||
await log_message(
|
await log_vertex_build(
|
||||||
sender=self.params.get("sender", ""),
|
flow_id=self.graph.flow_id,
|
||||||
sender_name=self.params.get("sender_name", ""),
|
vertex_id=self.id,
|
||||||
message=complete_message,
|
valid=True,
|
||||||
session_id=self.params.get("session_id", ""),
|
params=self._built_object_repr(),
|
||||||
|
data=self.result,
|
||||||
artifacts=self.artifacts,
|
artifacts=self.artifacts,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -87,8 +87,30 @@ export default function Page({
|
||||||
const [lastSelection, setLastSelection] =
|
const [lastSelection, setLastSelection] =
|
||||||
useState<OnSelectionChangeParams | null>(null);
|
useState<OnSelectionChangeParams | null>(null);
|
||||||
|
|
||||||
|
const setNode = useFlowStore((state) => state.setNode);
|
||||||
useEffect(() => {
|
useEffect(() => {
|
||||||
const onKeyDown = (event: KeyboardEvent) => {
|
const onKeyDown = (event: KeyboardEvent) => {
|
||||||
|
const selectedNode = nodes.filter((obj) => obj.selected);
|
||||||
|
if ((event.ctrlKey || event.metaKey) && event.key === "p" && selectedNode.length > 0) {
|
||||||
|
event.preventDefault();
|
||||||
|
setNode(selectedNode[0].id, (old) => ({
|
||||||
|
...old,
|
||||||
|
data: {
|
||||||
|
...old.data,
|
||||||
|
node: {
|
||||||
|
...old.data.node,
|
||||||
|
pinned: old.data?.node?.pinned ? false : true,
|
||||||
|
},
|
||||||
|
},
|
||||||
|
}));
|
||||||
|
}
|
||||||
|
if ((event.ctrlKey || event.metaKey) && event.key === "d" && selectedNode.length > 0) {
|
||||||
|
event.preventDefault();
|
||||||
|
paste({nodes: selectedNode, edges: []}, {
|
||||||
|
x: position.current.x,
|
||||||
|
y: position.current.y,
|
||||||
|
});
|
||||||
|
}
|
||||||
if (!isWrappedWithClass(event, "noundo")) {
|
if (!isWrappedWithClass(event, "noundo")) {
|
||||||
if (
|
if (
|
||||||
(event.key === "y" || (event.key === "z" && event.shiftKey)) &&
|
(event.key === "y" || (event.key === "z" && event.shiftKey)) &&
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue