Merge pull request #3343 from omChauhanDev/fix/auto-resolve-function-result
fix: keeping the Aggregator and Service states synchronized.
This commit is contained in:
@@ -595,10 +595,14 @@ class LLMService(AIService):
|
|||||||
cancel_on_interruption=item.cancel_on_interruption,
|
cancel_on_interruption=item.cancel_on_interruption,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
callback_executed = False
|
||||||
|
|
||||||
# Define a callback function that pushes a FunctionCallResultFrame upstream & downstream.
|
# Define a callback function that pushes a FunctionCallResultFrame upstream & downstream.
|
||||||
async def function_call_result_callback(
|
async def function_call_result_callback(
|
||||||
result: Any, *, properties: Optional[FunctionCallResultProperties] = None
|
result: Any, *, properties: Optional[FunctionCallResultProperties] = None
|
||||||
):
|
):
|
||||||
|
nonlocal callback_executed
|
||||||
|
callback_executed = True
|
||||||
await self.broadcast_frame(
|
await self.broadcast_frame(
|
||||||
FunctionCallResultFrame,
|
FunctionCallResultFrame,
|
||||||
function_name=runner_item.function_name,
|
function_name=runner_item.function_name,
|
||||||
@@ -609,40 +613,48 @@ class LLMService(AIService):
|
|||||||
properties=properties,
|
properties=properties,
|
||||||
)
|
)
|
||||||
|
|
||||||
if isinstance(item.handler, DirectFunctionWrapper):
|
try:
|
||||||
# Handler is a DirectFunctionWrapper
|
if isinstance(item.handler, DirectFunctionWrapper):
|
||||||
await item.handler.invoke(
|
# Handler is a DirectFunctionWrapper
|
||||||
args=runner_item.arguments,
|
await item.handler.invoke(
|
||||||
params=FunctionCallParams(
|
args=runner_item.arguments,
|
||||||
function_name=runner_item.function_name,
|
params=FunctionCallParams(
|
||||||
tool_call_id=runner_item.tool_call_id,
|
function_name=runner_item.function_name,
|
||||||
arguments=runner_item.arguments,
|
tool_call_id=runner_item.tool_call_id,
|
||||||
llm=self,
|
arguments=runner_item.arguments,
|
||||||
context=runner_item.context,
|
llm=self,
|
||||||
result_callback=function_call_result_callback,
|
context=runner_item.context,
|
||||||
),
|
result_callback=function_call_result_callback,
|
||||||
)
|
),
|
||||||
else:
|
|
||||||
# Handler is a FunctionCallHandler
|
|
||||||
if item.handler_deprecated:
|
|
||||||
await item.handler(
|
|
||||||
runner_item.function_name,
|
|
||||||
runner_item.tool_call_id,
|
|
||||||
runner_item.arguments,
|
|
||||||
self,
|
|
||||||
runner_item.context,
|
|
||||||
function_call_result_callback,
|
|
||||||
)
|
)
|
||||||
else:
|
else:
|
||||||
params = FunctionCallParams(
|
# Handler is a FunctionCallHandler
|
||||||
function_name=runner_item.function_name,
|
if item.handler_deprecated:
|
||||||
tool_call_id=runner_item.tool_call_id,
|
await item.handler(
|
||||||
arguments=runner_item.arguments,
|
runner_item.function_name,
|
||||||
llm=self,
|
runner_item.tool_call_id,
|
||||||
context=runner_item.context,
|
runner_item.arguments,
|
||||||
result_callback=function_call_result_callback,
|
self,
|
||||||
)
|
runner_item.context,
|
||||||
await item.handler(params)
|
function_call_result_callback,
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
params = FunctionCallParams(
|
||||||
|
function_name=runner_item.function_name,
|
||||||
|
tool_call_id=runner_item.tool_call_id,
|
||||||
|
arguments=runner_item.arguments,
|
||||||
|
llm=self,
|
||||||
|
context=runner_item.context,
|
||||||
|
result_callback=function_call_result_callback,
|
||||||
|
)
|
||||||
|
await item.handler(params)
|
||||||
|
except Exception as e:
|
||||||
|
error_message = f"Error executing function call [{runner_item.function_name}]: {e}"
|
||||||
|
logger.error(f"{self} {error_message}")
|
||||||
|
await self.push_error(error_msg=error_message, exception=e, fatal=False)
|
||||||
|
finally:
|
||||||
|
if not callback_executed:
|
||||||
|
await function_call_result_callback(None)
|
||||||
|
|
||||||
async def _cancel_function_call(self, function_name: Optional[str]):
|
async def _cancel_function_call(self, function_name: Optional[str]):
|
||||||
cancelled_tasks = set()
|
cancelled_tasks = set()
|
||||||
|
|||||||
Reference in New Issue
Block a user