Refactor vertex streaming logic in build_vertex_stream function
This commit is contained in:
parent
9af9647fb3
commit
116518202c
1 changed files with 23 additions and 13 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",
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue