import asyncio
from typing import Any
from google.adk.events import Event
from google.adk.tools.long_running_tool import LongRunningFunctionTool
from google.genai.types import Content, FunctionCall, FunctionResponse, Part
from veadk import Agent, Runner
APP_NAME = "long_running_tool_app"
USER_ID = "long_running_tool_user"
SESSION_ID = "long_running_tool_session"
def big_data_processing(data_url: str) -> dict[str, Any]:
"""Start processing big data located at a URL.
Args:
data_url (str): The URL of the big data to process.
Returns:
dict[str, Any]: The initial task state, with "status" == "pending",
the "data-url", and a "task-id" to track the job.
"""
# Simulate a submitted job; replace this with a real job service.
return {
"status": "pending",
"data-url": data_url,
"task-id": "big-data-processing-1",
}
long_running_tool = LongRunningFunctionTool(func=big_data_processing)
agent = Agent(
name="long_running_tool_agent",
model_name="doubao-seed-2-1-pro-260628",
instruction="Use big_data_processing to process big data.",
tools=[long_running_tool],
)
runner = Runner(agent=agent, app_name=APP_NAME)
def get_long_running_call(event: Event) -> FunctionCall | None:
"""Return the long-running function call carried by an event, if any."""
if not event.long_running_tool_ids or not event.content or not event.content.parts:
return None
for part in event.content.parts:
if (
part.function_call
and part.function_call.id in event.long_running_tool_ids
):
return part.function_call
return None
def get_function_response(
event: Event, function_call_id: str
) -> FunctionResponse | None:
"""Return the function response matching a given call id, if any."""
if not event.content or not event.content.parts:
return None
for part in event.content.parts:
if (
part.function_response
and part.function_response.id == function_call_id
):
return part.function_response
return None
async def main():
# Create the session before running.
session = await runner.short_term_memory.create_session(
app_name=APP_NAME, user_id=USER_ID, session_id=SESSION_ID
)
query = "Process the big data from https://example.com/data.csv"
content = Content(role="user", parts=[Part(text=query)])
print("Running agent...")
long_running_call = None
long_running_response = None
async for event in runner.run_async(
session_id=session.id, user_id=USER_ID, new_message=content
):
if long_running_call is None:
long_running_call = get_long_running_call(event)
elif long_running_response is None:
long_running_response = get_function_response(event, long_running_call.id)
if event.content and event.content.parts:
if text := "".join(part.text or "" for part in event.content.parts):
print(f"[{event.author}]: {text}")
# Simulate completion and send the final result back to the agent.
if long_running_response is not None:
updated = long_running_response.model_copy(deep=True)
updated.response = {"status": "finish"}
async for event in runner.run_async(
session_id=session.id,
user_id=USER_ID,
new_message=Content(
role="user", parts=[Part(function_response=updated)]
),
):
if event.content and event.content.parts:
if text := "".join(part.text or "" for part in event.content.parts):
print(f"[{event.author}]: {text}")
if __name__ == "__main__":
asyncio.run(main())