🐛 fix(chat.py): add support for running celery tasks for building vertices and fallback to local build if celery task fails
This commit is contained in:
parent
24344ee0fa
commit
2dd43d3f99
1 changed files with 21 additions and 1 deletions
|
|
@ -123,6 +123,9 @@ async def stream_build(flow_id: str):
|
||||||
"log": f"Building node {vertex.vertex_type}",
|
"log": f"Building node {vertex.vertex_type}",
|
||||||
}
|
}
|
||||||
yield str(StreamData(event="log", data=log_dict))
|
yield str(StreamData(event="log", data=log_dict))
|
||||||
|
if vertex.is_task:
|
||||||
|
vertex = try_running_celery_task(vertex)
|
||||||
|
else:
|
||||||
vertex.build()
|
vertex.build()
|
||||||
params = vertex._built_object_repr()
|
params = vertex._built_object_repr()
|
||||||
valid = True
|
valid = True
|
||||||
|
|
@ -179,3 +182,20 @@ async def stream_build(flow_id: str):
|
||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
logger.error(f"Error streaming build: {exc}")
|
logger.error(f"Error streaming build: {exc}")
|
||||||
raise HTTPException(status_code=500, detail=str(exc))
|
raise HTTPException(status_code=500, detail=str(exc))
|
||||||
|
|
||||||
|
|
||||||
|
def try_running_celery_task(vertex):
|
||||||
|
# Try running the task in celery
|
||||||
|
# and set the task_id to the local vertex
|
||||||
|
# if it fails, run the task locally
|
||||||
|
try:
|
||||||
|
from langflow.worker import build_vertex
|
||||||
|
|
||||||
|
task = build_vertex.delay(vertex)
|
||||||
|
vertex.task_id = task.id
|
||||||
|
except Exception as exc:
|
||||||
|
logger.exception(exc)
|
||||||
|
logger.error("Error running task in celery, running locally")
|
||||||
|
vertex.task_id = None
|
||||||
|
vertex.build()
|
||||||
|
return vertex
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue