diff --git a/src/lfx/src/lfx/graph/schema.py b/src/lfx/src/lfx/graph/schema.py index 38c6610715..f45d117034 100644 --- a/src/lfx/src/lfx/graph/schema.py +++ b/src/lfx/src/lfx/graph/schema.py @@ -39,6 +39,9 @@ class ResultData(BaseModel): if message is None: continue + if not isinstance(message, dict): + continue + if "stream_url" in message and "type" in message: stream_url = StreamURL(location=message["stream_url"]) values["outputs"].update({key: OutputValue(message=stream_url, type=message["type"])}) diff --git a/src/lfx/src/lfx/schema/schema.py b/src/lfx/src/lfx/schema/schema.py index b1a2e19e20..f1d9e02232 100644 --- a/src/lfx/src/lfx/schema/schema.py +++ b/src/lfx/src/lfx/schema/schema.py @@ -110,7 +110,7 @@ def build_output_logs(vertex, result) -> dict: type_ = get_type(output_result) match type_: - case LogType.STREAM if "stream_url" in message: + case LogType.STREAM if isinstance(message, dict) and "stream_url" in message: message = StreamURL(location=message["stream_url"]) case LogType.STREAM: @@ -123,9 +123,13 @@ def build_output_logs(vertex, result) -> dict: message = "" case LogType.ARRAY: - if isinstance(message, DataFrame): + if message is None: + message = [] + elif isinstance(message, DataFrame): message = message.to_dict(orient="records") - message = [serialize(item) for item in message] + message = [serialize(item) for item in message] + else: + message = [serialize(item) for item in message] name = output.get("name", f"output_{index}") outputs |= {name: OutputValue(message=message, type=type_).model_dump()}