Skip to content
New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

fix: serialize array output and logs #3040

Merged
merged 3 commits into from
Jul 29, 2024
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
22 changes: 1 addition & 21 deletions src/backend/base/langflow/custom/custom_component/component.py
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
import inspect
from typing import Any, AsyncIterator, Callable, ClassVar, Generator, Iterator, List, Optional, Union
from typing import Any, Callable, ClassVar, List, Optional, Union
from uuid import UUID

import yaml
Expand All @@ -15,26 +15,6 @@
from .custom_component import CustomComponent


def recursive_serialize_or_str(obj):
try:
if isinstance(obj, dict):
return {k: recursive_serialize_or_str(v) for k, v in obj.items()}
elif isinstance(obj, list):
return [recursive_serialize_or_str(v) for v in obj]
elif isinstance(obj, BaseModel):
return {k: recursive_serialize_or_str(v) for k, v in obj.model_dump().items()}
elif isinstance(obj, (AsyncIterator, Generator, Iterator)):
# contain memory addresses
# without consuming the iterator
# return list(obj) consumes the iterator
# return f"{obj}" this generates '<generator object BaseChatModel.stream at 0x33e9ec770>'
# it is not useful
return "Unconsumed Stream"
return str(obj)
except Exception:
return str(obj)


class Component(CustomComponent):
inputs: List[InputTypes] = []
outputs: List[Output] = []
Expand Down
11 changes: 11 additions & 0 deletions src/backend/base/langflow/schema/artifact.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@

from langflow.schema import Data
from langflow.schema.message import Message
from langflow.schema.schema import recursive_serialize_or_str


class ArtifactType(str, Enum):
Expand Down Expand Up @@ -52,6 +53,16 @@ def get_artifact_type(value, build_result=None) -> str:
def post_process_raw(raw, artifact_type: str):
if artifact_type == ArtifactType.STREAM.value:
raw = ""
elif artifact_type == ArtifactType.ARRAY.value:
_raw = []
for item in raw:
if hasattr(item, "dict"):
_raw.append(recursive_serialize_or_str(item))
elif hasattr(item, "model_dump"):
_raw.append(recursive_serialize_or_str(item))
else:
_raw.append(str(item))
raw = _raw
elif artifact_type == ArtifactType.UNKNOWN.value and raw is not None:
if isinstance(raw, (BaseModel, dict)):
try:
Expand Down
37 changes: 36 additions & 1 deletion src/backend/base/langflow/schema/schema.py
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
from enum import Enum
from typing import Generator, Literal, Union
from typing import AsyncIterator, Generator, Iterator, Literal, Union

from pydantic import BaseModel
from typing_extensions import TypedDict
Expand Down Expand Up @@ -104,7 +104,42 @@ def build_output_logs(vertex, result) -> dict:

case LogType.UNKNOWN:
message = ""

case LogType.ARRAY:
message = [recursive_serialize_or_str(item) for item in message]
name = output.get("name", f"output_{index}")
outputs |= {name: OutputValue(message=message, type=_type).model_dump()}

return outputs


def recursive_serialize_or_str(obj):
try:
if isinstance(obj, dict):
return {k: recursive_serialize_or_str(v) for k, v in obj.items()}
elif isinstance(obj, list):
return [recursive_serialize_or_str(v) for v in obj]
elif isinstance(obj, BaseModel):
if hasattr(obj, "model_dump"):
obj_dict = obj.model_dump()
elif hasattr(obj, "dict"):
obj_dict = obj.dict() # type: ignore
return {k: recursive_serialize_or_str(v) for k, v in obj_dict.items()}

elif isinstance(obj, (AsyncIterator, Generator, Iterator)):
# contain memory addresses
# without consuming the iterator
# return list(obj) consumes the iterator
# return f"{obj}" this generates '<generator object BaseChatModel.stream at 0x33e9ec770>'
# it is not useful
return "Unconsumed Stream"
elif hasattr(obj, "dict"):
return {k: recursive_serialize_or_str(v) for k, v in obj.dict().items()}
elif hasattr(obj, "model_dump"):
return {k: recursive_serialize_or_str(v) for k, v in obj.model_dump().items()}
elif issubclass(obj, BaseModel):
# This a type BaseModel and not an instance of it
return repr(obj)
return str(obj)
except Exception:
return str(obj)
Loading