services.py 54 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346
  1. import json
  2. from datetime import datetime, timedelta
  3. from typing import cast
  4. from sqlalchemy.orm import Session
  5. from core_events import EventPublishContract, EventServiceClient, EventServiceClientError
  6. from uuid import uuid4
  7. from core_domain import (
  8. AgentSkillRefContract,
  9. AgentToolRefContract,
  10. ChatCompletionRequestContract,
  11. ChatCompletionResponseContract,
  12. ChatMessageContract,
  13. MemoryCreateContract,
  14. MemoryScopeType,
  15. MemorySearchRequestContract,
  16. MemorySearchResultContract)
  17. from core_shared import JSONValue, try_build_redis_client
  18. from core_shared.task_queue import TaskQueuePublisher
  19. from app.bootstrap.settings import AgentServiceSettings
  20. from app.db.models import AgentDefinition, AgentRun, AgentToolInvocation, AgentVersion
  21. from app.domain.repositories import (
  22. AgentDefinitionRepository,
  23. AgentRunRepository,
  24. AgentToolInvocationRepository,
  25. AgentVersionRepository)
  26. from app.infrastructure.model_gateway_client import ModelGatewayClient, ModelGatewayClientError
  27. from app.infrastructure.memory_client import MemoryClient, MemoryClientError
  28. from app.infrastructure.skill_client import SkillServiceClient, SkillServiceClientError
  29. from app.infrastructure.tool_client import ToolServiceClient, ToolServiceClientError
  30. from app.schemas.agent import (
  31. AgentCreateRequest,
  32. AgentConfigCreateRequest,
  33. AgentConfigListRequest,
  34. AgentRunCreateRequest,
  35. AgentRunDetailRequest,
  36. AgentRunExecuteRequest,
  37. AgentRunStatusUpdateRequest,
  38. AgentStatusUpdateRequest,
  39. AgentUpdateRequest,
  40. AgentVersionCreateRequest)
  41. def generate_agent_code() -> str:
  42. return f"agent_{uuid4().hex[:16]}"
  43. class AgentApplicationService:
  44. def __init__(
  45. self,
  46. *,
  47. agent_repository: AgentDefinitionRepository,
  48. agent_version_repository: AgentVersionRepository,
  49. agent_run_repository: AgentRunRepository,
  50. agent_tool_invocation_repository: AgentToolInvocationRepository,
  51. model_gateway_client: ModelGatewayClient | None = None,
  52. memory_client: MemoryClient | None = None,
  53. tool_client: ToolServiceClient | None = None,
  54. skill_client: SkillServiceClient | None = None,
  55. event_client: EventServiceClient | None = None,
  56. task_queue_publisher: TaskQueuePublisher | None = None,
  57. react_max_steps: int = 5,
  58. react_max_tool_calls: int = 10,
  59. react_tool_retry_count: int = 1) -> None:
  60. self.agent_repository = agent_repository
  61. self.agent_version_repository = agent_version_repository
  62. self.agent_run_repository = agent_run_repository
  63. self.agent_tool_invocation_repository = agent_tool_invocation_repository
  64. self.model_gateway_client = model_gateway_client
  65. self.memory_client = memory_client
  66. self.tool_client = tool_client
  67. self.skill_client = skill_client
  68. self.event_client = event_client
  69. self.task_queue_publisher = task_queue_publisher
  70. self.react_max_steps = react_max_steps
  71. self.react_max_tool_calls = react_max_tool_calls
  72. self.react_tool_retry_count = react_tool_retry_count
  73. def create_agent(self, payload: AgentCreateRequest) -> AgentDefinition:
  74. return self.agent_repository.create(
  75. code=payload.code or generate_agent_code(),
  76. name=payload.name,
  77. description=payload.description,
  78. agent_type=payload.agent_type,
  79. owner_user_id=payload.owner_user_id,
  80. metadata_json=payload.metadata_json)
  81. def list_agents(self) -> list[AgentDefinition]:
  82. return self.agent_repository.list_all()
  83. def update_agent(self, payload: AgentUpdateRequest) -> AgentDefinition | None:
  84. return self.agent_repository.update(
  85. agent_id=payload.agent_id,
  86. name=payload.name,
  87. description=payload.description,
  88. metadata_json=payload.metadata_json)
  89. def update_agent_status(
  90. self,
  91. *,
  92. agent_id: str,
  93. payload: AgentStatusUpdateRequest) -> AgentDefinition | None:
  94. return self.agent_repository.update_status(
  95. agent_id=agent_id,
  96. status=payload.status)
  97. def create_agent_config(self, payload: AgentConfigCreateRequest) -> AgentVersion:
  98. return self.create_agent_version(
  99. AgentVersionCreateRequest(
  100. agent_id=payload.agent_id,
  101. status="draft",
  102. role=payload.role,
  103. goal=payload.goal,
  104. system_prompt=payload.system_prompt,
  105. model_config=payload.model_config_data,
  106. memory_policy=payload.memory_policy,
  107. tool_refs=payload.tool_refs,
  108. skill_refs=payload.skill_refs))
  109. def list_agent_configs(self, payload: AgentConfigListRequest) -> list[AgentVersion]:
  110. return self.list_agent_versions(agent_id=payload.agent_id)
  111. def create_agent_version(self, payload: AgentVersionCreateRequest) -> AgentVersion:
  112. agent = self.agent_repository.get_by_id(
  113. agent_id=payload.agent_id)
  114. if agent is None:
  115. raise ValueError(f"agent not found: {payload.agent_id}")
  116. return self.agent_version_repository.create(
  117. agent_id=payload.agent_id,
  118. status=payload.status,
  119. role=payload.role,
  120. goal=payload.goal,
  121. system_prompt=payload.system_prompt,
  122. model_config_json=payload.model_config_data.model_dump(mode="json"),
  123. memory_policy_json=payload.memory_policy.model_dump(mode="json"),
  124. tool_refs_json=[item.model_dump(mode="json") for item in payload.tool_refs],
  125. skill_refs_json=[item.model_dump(mode="json") for item in payload.skill_refs])
  126. def list_agent_versions(self, *, agent_id: str) -> list[AgentVersion]:
  127. return self.agent_version_repository.list_by_agent(agent_id=agent_id)
  128. def create_agent_run(self, payload: AgentRunCreateRequest) -> AgentRun:
  129. agent_version = self._resolve_agent_version(
  130. agent_id=payload.agent_id,
  131. agent_version_id=payload.agent_config_id or payload.agent_version_id)
  132. if agent_version is None:
  133. raise ValueError("agent config not found")
  134. agent_run = self.agent_run_repository.create(
  135. agent_id=payload.agent_id,
  136. agent_version_id=agent_version.id,
  137. session_id=payload.session_id,
  138. input_text=payload.input_text,
  139. input_json=payload.input_json)
  140. self._publish_event(
  141. event_type="agent.run.created",
  142. agent_run=agent_run,
  143. payload_json={"agent_run_id": agent_run.id, "status": agent_run.status})
  144. if self.task_queue_publisher is not None:
  145. self.task_queue_publisher.publish_agent_run(
  146. agent_run_id=agent_run.id)
  147. return agent_run
  148. def list_agent_runs(
  149. self,
  150. *,
  151. agent_id: str | None = None,
  152. session_id: str | None = None) -> list[AgentRun]:
  153. return self.agent_run_repository.list_by_scope(
  154. agent_id=agent_id,
  155. session_id=session_id)
  156. def get_agent_run(self, payload: AgentRunDetailRequest) -> AgentRun | None:
  157. return self.agent_run_repository.get_by_id(
  158. agent_run_id=payload.agent_run_id)
  159. def list_agent_tool_invocations(
  160. self,
  161. *,
  162. agent_run_id: str) -> list[AgentToolInvocation]:
  163. return self.agent_tool_invocation_repository.list_by_run(
  164. agent_run_id=agent_run_id)
  165. def update_agent_run_status(
  166. self,
  167. *,
  168. agent_run_id: str,
  169. payload: AgentRunStatusUpdateRequest) -> AgentRun | None:
  170. entity = self.agent_run_repository.get_by_id(
  171. agent_run_id=agent_run_id)
  172. if entity is None:
  173. return None
  174. return self.agent_run_repository.update_status(
  175. agent_run_id=agent_run_id,
  176. status=payload.status,
  177. worker_key=payload.worker_key,
  178. output_text=payload.output_text,
  179. output_json=payload.output_json,
  180. error_code=payload.error_code,
  181. error_message=payload.error_message)
  182. def execute_agent_run(
  183. self,
  184. *,
  185. agent_run_id: str,
  186. payload: AgentRunExecuteRequest) -> AgentRun | None:
  187. agent_run = self.agent_run_repository.get_by_id(
  188. agent_run_id=agent_run_id)
  189. if agent_run is None:
  190. return None
  191. agent_version = self.agent_version_repository.get_by_id(
  192. agent_version_id=agent_run.agent_version_id)
  193. if agent_version is None:
  194. return self.agent_run_repository.update_status(
  195. agent_run_id=agent_run.id,
  196. status="failed",
  197. worker_key=payload.worker_key,
  198. error_code="agent_version_missing",
  199. error_message=f"agent version not found: {agent_run.agent_version_id}")
  200. self.agent_run_repository.update_status(
  201. agent_run_id=agent_run.id,
  202. status="running",
  203. worker_key=payload.worker_key)
  204. memory_results, memory_metadata = self._read_relevant_memories(
  205. agent_run=agent_run,
  206. agent_version=agent_version)
  207. selected_tools = self._select_tool_refs(agent_run=agent_run, agent_version=agent_version)
  208. selected_skills = self._select_skill_refs(agent_run=agent_run, agent_version=agent_version)
  209. if payload.dry_run:
  210. messages = self._build_chat_messages(
  211. agent_run=agent_run,
  212. agent_version=agent_version,
  213. memory_results=memory_results,
  214. capability_context=self._format_capability_plan(
  215. selected_tools=selected_tools,
  216. selected_skills=selected_skills))
  217. completed_run = self.agent_run_repository.update_status(
  218. agent_run_id=agent_run.id,
  219. status="completed",
  220. worker_key=payload.worker_key,
  221. output_text=self._build_dry_run_output(
  222. agent_run=agent_run,
  223. agent_version=agent_version),
  224. output_json={
  225. "dry_run": True,
  226. "agent_version_id": agent_version.id,
  227. "message_count": len(messages),
  228. "messages": [message.model_dump(mode="json") for message in messages],
  229. "selected_tool_refs": [
  230. tool_ref.model_dump(mode="json") for tool_ref in selected_tools
  231. ],
  232. "selected_skill_refs": [
  233. skill_ref.model_dump(mode="json") for skill_ref in selected_skills
  234. ],
  235. **memory_metadata,
  236. })
  237. if completed_run is not None:
  238. self._publish_event(
  239. event_type="agent.run.completed",
  240. agent_run=completed_run,
  241. payload_json={
  242. "agent_run_id": completed_run.id,
  243. "dry_run": True,
  244. "status": completed_run.status,
  245. })
  246. return completed_run
  247. if self._read_bool(agent_version.model_config_json, "react_enabled", default=False):
  248. return self._execute_react_agent_run(
  249. agent_run=agent_run,
  250. agent_version=agent_version,
  251. payload=payload,
  252. memory_results=memory_results,
  253. memory_metadata=memory_metadata,
  254. selected_tools=selected_tools,
  255. selected_skills=selected_skills)
  256. tool_invocations = self._invoke_selected_tools(
  257. agent_run=agent_run,
  258. agent_version=agent_version,
  259. selected_tools=selected_tools)
  260. skill_invocations = self._invoke_selected_skills(
  261. agent_run=agent_run,
  262. selected_skills=selected_skills,
  263. worker_key=payload.worker_key)
  264. messages = self._build_chat_messages(
  265. agent_run=agent_run,
  266. agent_version=agent_version,
  267. memory_results=memory_results,
  268. capability_context=self._format_capability_results(
  269. tool_invocations=tool_invocations,
  270. skill_invocations=skill_invocations))
  271. if self.model_gateway_client is None:
  272. return self.agent_run_repository.update_status(
  273. agent_run_id=agent_run.id,
  274. status="failed",
  275. worker_key=payload.worker_key,
  276. error_code="model_gateway_missing",
  277. error_message="model gateway client is not configured",
  278. output_json={
  279. "tool_invocations": tool_invocations,
  280. "skill_invocations": skill_invocations,
  281. **memory_metadata,
  282. })
  283. try:
  284. response = self.model_gateway_client.create_chat_completion(
  285. ChatCompletionRequestContract(
  286. model=self._read_optional_string(agent_version.model_config_json, "model"),
  287. temperature=self._read_optional_float(
  288. agent_version.model_config_json,
  289. "temperature"),
  290. max_tokens=self._read_optional_int(
  291. agent_version.model_config_json,
  292. "max_tokens"),
  293. messages=messages,
  294. metadata_json={
  295. "agent_id": agent_run.agent_id,
  296. "agent_version_id": agent_version.id,
  297. "agent_run_id": agent_run.id,
  298. })
  299. )
  300. except ModelGatewayClientError as exc:
  301. return self.agent_run_repository.update_status(
  302. agent_run_id=agent_run.id,
  303. status="failed",
  304. worker_key=payload.worker_key,
  305. error_code="model_gateway_error",
  306. error_message=str(exc))
  307. memory_write_metadata = self._write_interaction_memory(
  308. agent_run=agent_run,
  309. agent_version=agent_version,
  310. output_text=response.content)
  311. completed_run = self.agent_run_repository.update_status(
  312. agent_run_id=agent_run.id,
  313. status="completed",
  314. worker_key=payload.worker_key,
  315. output_text=response.content,
  316. output_json={
  317. "dry_run": False,
  318. "agent_version_id": agent_version.id,
  319. "model": response.model,
  320. "finish_reason": response.finish_reason,
  321. "usage_json": response.usage_json,
  322. "raw_response_json": response.raw_response_json,
  323. "tool_invocations": tool_invocations,
  324. "skill_invocations": skill_invocations,
  325. **memory_metadata,
  326. **memory_write_metadata,
  327. })
  328. if completed_run is not None:
  329. self._publish_event(
  330. event_type="agent.run.completed",
  331. agent_run=completed_run,
  332. payload_json={
  333. "agent_run_id": completed_run.id,
  334. "dry_run": False,
  335. "status": completed_run.status,
  336. })
  337. return completed_run
  338. def _publish_event(
  339. self,
  340. *,
  341. event_type: str,
  342. agent_run: AgentRun,
  343. payload_json: dict[str, JSONValue]) -> None:
  344. if self.event_client is None:
  345. return
  346. try:
  347. self.event_client.publish_event(
  348. EventPublishContract(
  349. event_type=event_type,
  350. source_service="agent-service",
  351. aggregate_type="agent_run",
  352. aggregate_id=agent_run.id,
  353. correlation_id=agent_run.session_id,
  354. payload_json={
  355. **payload_json,
  356. "agent_id": agent_run.agent_id,
  357. "agent_version_id": agent_run.agent_version_id,
  358. })
  359. )
  360. except EventServiceClientError:
  361. return
  362. def _execute_react_agent_run(
  363. self,
  364. *,
  365. agent_run: AgentRun,
  366. agent_version: AgentVersion,
  367. payload: AgentRunExecuteRequest,
  368. memory_results: list[MemorySearchResultContract],
  369. memory_metadata: dict[str, JSONValue],
  370. selected_tools: list[AgentToolRefContract],
  371. selected_skills: list[AgentSkillRefContract]) -> AgentRun | None:
  372. if self.model_gateway_client is None:
  373. return self.agent_run_repository.update_status(
  374. agent_run_id=agent_run.id,
  375. status="failed",
  376. worker_key=payload.worker_key,
  377. error_code="model_gateway_missing",
  378. error_message="model gateway client is not configured")
  379. skill_invocations = self._invoke_selected_skills(
  380. agent_run=agent_run,
  381. selected_skills=selected_skills,
  382. worker_key=payload.worker_key)
  383. messages = self._build_chat_messages(
  384. agent_run=agent_run,
  385. agent_version=agent_version,
  386. memory_results=memory_results,
  387. capability_context=self._format_react_instruction(
  388. agent_run=agent_run,
  389. selected_tools=selected_tools,
  390. skill_invocations=skill_invocations))
  391. react_steps: list[dict[str, JSONValue]] = []
  392. tool_invocations: list[dict[str, JSONValue]] = []
  393. final_answer: str | None = None
  394. tool_call_count = 0
  395. max_steps = self._read_int(
  396. agent_version.model_config_json,
  397. "react_max_steps",
  398. default=self.react_max_steps)
  399. for step_index in range(max(max_steps, 1)):
  400. try:
  401. response = self.model_gateway_client.create_chat_completion(
  402. self._build_chat_completion_request(
  403. agent_run=agent_run,
  404. agent_version=agent_version,
  405. messages=messages,
  406. selected_tools=selected_tools)
  407. )
  408. except ModelGatewayClientError as exc:
  409. return self.agent_run_repository.update_status(
  410. agent_run_id=agent_run.id,
  411. status="failed",
  412. worker_key=payload.worker_key,
  413. error_code="model_gateway_error",
  414. error_message=str(exc),
  415. output_json={
  416. "react_steps": react_steps,
  417. "tool_invocations": tool_invocations,
  418. "skill_invocations": skill_invocations,
  419. **memory_metadata,
  420. })
  421. action = self._parse_react_action_from_response(response)
  422. react_step: dict[str, JSONValue] = {
  423. "step_index": step_index,
  424. "model_content": response.content,
  425. "action": action,
  426. }
  427. react_steps.append(react_step)
  428. if action.get("action") == "finish":
  429. answer_value = action.get("answer")
  430. final_answer = answer_value if isinstance(answer_value, str) else response.content
  431. break
  432. if action.get("action") != "tool":
  433. final_answer = response.content
  434. break
  435. max_tool_calls = self._read_int(
  436. agent_version.model_config_json,
  437. "react_max_tool_calls",
  438. default=self.react_max_tool_calls)
  439. if tool_call_count >= max(max_tool_calls, 0):
  440. final_answer = "Tool call budget exhausted."
  441. react_step["observation"] = final_answer
  442. break
  443. tool_code = action.get("tool_code")
  444. matching_tools = [
  445. item for item in selected_tools if item.tool_code == tool_code
  446. ]
  447. if not matching_tools:
  448. observation = f"tool not available: {tool_code}"
  449. react_step["observation"] = observation
  450. messages.append(ChatMessageContract(role="assistant", content=response.content))
  451. messages.append(ChatMessageContract(role="user", content=observation))
  452. continue
  453. tool_input = action.get("input_json")
  454. original_input_json = agent_run.input_json
  455. if isinstance(tool_input, dict):
  456. agent_run.input_json = {
  457. str(item_key): item_value for item_key, item_value in tool_input.items()
  458. }
  459. current_invocations = self._invoke_react_tool_with_retry(
  460. agent_run=agent_run,
  461. agent_version=agent_version,
  462. tool_ref=matching_tools[0])
  463. tool_call_count += len(current_invocations)
  464. agent_run.input_json = original_input_json
  465. tool_invocations.extend(current_invocations)
  466. observation = self._format_react_observation(current_invocations)
  467. react_step["observation"] = observation
  468. messages.append(ChatMessageContract(role="assistant", content=response.content))
  469. messages.append(ChatMessageContract(role="user", content=observation))
  470. if final_answer is None:
  471. final_answer = "ReAct loop reached max steps without a final answer."
  472. memory_write_metadata = self._write_interaction_memory(
  473. agent_run=agent_run,
  474. agent_version=agent_version,
  475. output_text=final_answer)
  476. completed_run = self.agent_run_repository.update_status(
  477. agent_run_id=agent_run.id,
  478. status="completed",
  479. worker_key=payload.worker_key,
  480. output_text=final_answer,
  481. output_json={
  482. "dry_run": False,
  483. "agent_version_id": agent_version.id,
  484. "react_enabled": True,
  485. "react_steps": react_steps,
  486. "react_tool_call_count": tool_call_count,
  487. "tool_invocations": tool_invocations,
  488. "skill_invocations": skill_invocations,
  489. **memory_metadata,
  490. **memory_write_metadata,
  491. })
  492. if completed_run is not None:
  493. self._publish_event(
  494. event_type="agent.run.completed",
  495. agent_run=completed_run,
  496. payload_json={
  497. "agent_run_id": completed_run.id,
  498. "react_enabled": True,
  499. "status": completed_run.status,
  500. })
  501. return completed_run
  502. def execute_next_claimed_agent_run(
  503. self,
  504. *,
  505. worker_key: str,
  506. lease_seconds: int,
  507. dry_run: bool,
  508. redis_client: object | None = None) -> tuple[AgentRun, int] | None:
  509. released_lease_count = self.agent_run_repository.release_expired_leases(
  510. now_time=datetime.utcnow())
  511. claimed_agent_run = self.agent_run_repository.claim_next_queued(
  512. worker_key=worker_key,
  513. lease_expire_time=datetime.utcnow() + timedelta(seconds=lease_seconds))
  514. if claimed_agent_run is None:
  515. return None
  516. if redis_client is not None:
  517. from core_shared.redis_primitives import DistributedLock, IdempotencyStore
  518. lock = DistributedLock(
  519. client=redis_client,
  520. name=f"agent-run:{claimed_agent_run.id}:lock",
  521. ttl_seconds=lease_seconds)
  522. if not lock.acquire():
  523. return None
  524. idempotency_store = IdempotencyStore(
  525. client=redis_client,
  526. prefix="agent-run-idempotency")
  527. if not idempotency_store.begin(key=claimed_agent_run.id):
  528. lock.release()
  529. return None
  530. else:
  531. lock = None
  532. idempotency_store = None
  533. try:
  534. result = self.execute_agent_run(
  535. agent_run_id=claimed_agent_run.id,
  536. payload=AgentRunExecuteRequest(
  537. worker_key=worker_key,
  538. dry_run=dry_run))
  539. if idempotency_store is not None and result is not None:
  540. idempotency_store.complete(
  541. key=claimed_agent_run.id,
  542. result={"status": result.status, "agent_run_id": result.id})
  543. finally:
  544. if lock is not None:
  545. lock.release()
  546. if result is None:
  547. return None
  548. return result, released_lease_count
  549. def _resolve_agent_version(
  550. self,
  551. *,
  552. agent_id: str,
  553. agent_version_id: str | None) -> AgentVersion | None:
  554. if agent_version_id is not None:
  555. return self.agent_version_repository.get_by_id(
  556. agent_version_id=agent_version_id)
  557. return self.agent_version_repository.get_latest_by_agent(
  558. agent_id=agent_id)
  559. def _build_chat_messages(
  560. self,
  561. *,
  562. agent_run: AgentRun,
  563. agent_version: AgentVersion,
  564. memory_results: list[MemorySearchResultContract] | None = None,
  565. capability_context: str | None = None) -> list[ChatMessageContract]:
  566. messages = [
  567. ChatMessageContract(role="system", content=agent_version.system_prompt),
  568. ]
  569. if agent_version.goal:
  570. messages.append(
  571. ChatMessageContract(role="system", content=f"Goal: {agent_version.goal}")
  572. )
  573. if memory_results:
  574. messages.append(
  575. ChatMessageContract(
  576. role="system",
  577. content=self._format_memory_context(memory_results))
  578. )
  579. if capability_context:
  580. messages.append(ChatMessageContract(role="system", content=capability_context))
  581. if agent_run.input_text:
  582. messages.append(ChatMessageContract(role="user", content=agent_run.input_text))
  583. if agent_run.input_json:
  584. messages.append(
  585. ChatMessageContract(
  586. role="user",
  587. content=f"Structured input: {agent_run.input_json}")
  588. )
  589. return messages
  590. def _select_tool_refs(
  591. self,
  592. *,
  593. agent_run: AgentRun,
  594. agent_version: AgentVersion) -> list[AgentToolRefContract]:
  595. input_preview = self._build_input_preview(agent_run)
  596. selected: list[AgentToolRefContract] = []
  597. for item in agent_version.tool_refs_json:
  598. ref = AgentToolRefContract.model_validate(item)
  599. if (
  600. ref.required
  601. or self._read_bool(ref.config_json, "auto_invoke", default=False)
  602. or self._matches_selection_keywords(ref.config_json, input_preview)
  603. ):
  604. selected.append(ref)
  605. return selected
  606. def _select_skill_refs(
  607. self,
  608. *,
  609. agent_run: AgentRun,
  610. agent_version: AgentVersion) -> list[AgentSkillRefContract]:
  611. input_preview = self._build_input_preview(agent_run)
  612. selected: list[AgentSkillRefContract] = []
  613. for item in agent_version.skill_refs_json:
  614. ref = AgentSkillRefContract.model_validate(item)
  615. auto_invoke = self._read_bool(ref.config_json, "auto_invoke", default=True)
  616. if auto_invoke or self._matches_selection_keywords(ref.config_json, input_preview):
  617. selected.append(ref)
  618. return selected
  619. def _invoke_selected_tools(
  620. self,
  621. *,
  622. agent_run: AgentRun,
  623. agent_version: AgentVersion,
  624. selected_tools: list[AgentToolRefContract]) -> list[dict[str, JSONValue]]:
  625. invocations: list[dict[str, JSONValue]] = []
  626. for ref in selected_tools:
  627. invocation = self.agent_tool_invocation_repository.create(
  628. agent_run_id=agent_run.id,
  629. agent_id=agent_run.agent_id,
  630. agent_version_id=agent_version.id,
  631. tool_code=ref.tool_code,
  632. tool_binding_id=ref.tool_binding_id,
  633. status="selected",
  634. input_json=agent_run.input_json or {})
  635. if ref.tool_binding_id is None:
  636. self.agent_tool_invocation_repository.update_status(
  637. invocation_id=invocation.id,
  638. status="skipped",
  639. reason="tool_binding_id_missing")
  640. invocations.append(
  641. {
  642. "status": "skipped",
  643. "reason": "tool_binding_id_missing",
  644. "tool_code": ref.tool_code,
  645. }
  646. )
  647. continue
  648. if self.tool_client is None:
  649. self.agent_tool_invocation_repository.update_status(
  650. invocation_id=invocation.id,
  651. status="failed",
  652. reason="tool_client_missing",
  653. error_message="tool client is not configured")
  654. invocations.append(
  655. {
  656. "status": "failed",
  657. "reason": "tool_client_missing",
  658. "tool_binding_id": ref.tool_binding_id,
  659. }
  660. )
  661. continue
  662. try:
  663. self.agent_tool_invocation_repository.update_status(
  664. invocation_id=invocation.id,
  665. status="running")
  666. detail = self.tool_client.get_tool_binding_detail(
  667. binding_id=ref.tool_binding_id)
  668. if not detail.binding.enabled:
  669. self.agent_tool_invocation_repository.update_status(
  670. invocation_id=invocation.id,
  671. status="failed",
  672. reason="tool_binding_disabled",
  673. error_message="tool binding is disabled")
  674. invocations.append(
  675. {
  676. "status": "failed",
  677. "reason": "tool_binding_disabled",
  678. "tool_binding_id": ref.tool_binding_id,
  679. }
  680. )
  681. continue
  682. if detail.tool_definition.tool_type != "http":
  683. self.agent_tool_invocation_repository.update_status(
  684. invocation_id=invocation.id,
  685. status="skipped",
  686. reason="unsupported_tool_type",
  687. error_message=detail.tool_definition.tool_type)
  688. invocations.append(
  689. {
  690. "status": "skipped",
  691. "reason": "unsupported_tool_type",
  692. "tool_type": detail.tool_definition.tool_type,
  693. "tool_binding_id": ref.tool_binding_id,
  694. }
  695. )
  696. continue
  697. output_text, output_json = self.tool_client.invoke_http_tool(
  698. detail=detail,
  699. input_json=agent_run.input_json or {},
  700. config_json=ref.config_json)
  701. except ToolServiceClientError as exc:
  702. self.agent_tool_invocation_repository.update_status(
  703. invocation_id=invocation.id,
  704. status="failed",
  705. reason="tool_service_error",
  706. error_message=str(exc))
  707. invocations.append(
  708. {
  709. "status": "failed",
  710. "reason": str(exc),
  711. "tool_binding_id": ref.tool_binding_id,
  712. }
  713. )
  714. continue
  715. self.agent_tool_invocation_repository.update_status(
  716. invocation_id=invocation.id,
  717. status="completed",
  718. output_text=output_text,
  719. output_json=output_json)
  720. invocations.append(
  721. {
  722. "status": "completed",
  723. "tool_binding_id": ref.tool_binding_id,
  724. "tool_code": detail.tool_definition.code,
  725. "output_text": output_text,
  726. "output_json": output_json,
  727. }
  728. )
  729. return invocations
  730. def _invoke_selected_skills(
  731. self,
  732. *,
  733. agent_run: AgentRun,
  734. selected_skills: list[AgentSkillRefContract],
  735. worker_key: str | None) -> list[dict[str, JSONValue]]:
  736. invocations: list[dict[str, JSONValue]] = []
  737. for ref in selected_skills:
  738. if self.skill_client is None:
  739. invocations.append(
  740. {
  741. "status": "failed",
  742. "reason": "skill_client_missing",
  743. "skill_id": ref.skill_id,
  744. "skill_code": ref.skill_code,
  745. }
  746. )
  747. continue
  748. skill_id = ref.skill_id or self._resolve_skill_id_by_code(
  749. skill_code=ref.skill_code)
  750. if skill_id is None:
  751. invocations.append(
  752. {
  753. "status": "failed",
  754. "reason": "skill_id_missing",
  755. "skill_code": ref.skill_code,
  756. }
  757. )
  758. continue
  759. try:
  760. created_run = self.skill_client.create_skill_run(
  761. skill_id=skill_id,
  762. skill_version_id=self._read_optional_string(
  763. ref.config_json,
  764. "skill_version_id"),
  765. installation_id=self._read_optional_string(
  766. ref.config_json,
  767. "installation_id"),
  768. input_json=self._build_skill_input_json(agent_run=agent_run, ref=ref))
  769. executed_run = self.skill_client.execute_skill_run(
  770. skill_run_id=created_run.id,
  771. worker_key=worker_key)
  772. except SkillServiceClientError as exc:
  773. invocations.append(
  774. {
  775. "status": "failed",
  776. "reason": str(exc),
  777. "skill_id": skill_id,
  778. "skill_code": ref.skill_code,
  779. }
  780. )
  781. continue
  782. invocations.append(
  783. {
  784. "status": executed_run.status,
  785. "skill_id": skill_id,
  786. "skill_code": ref.skill_code,
  787. "skill_run_id": executed_run.id,
  788. "output_text": executed_run.output_text,
  789. "output_json": executed_run.output_json,
  790. "error_code": executed_run.error_code,
  791. "error_message": executed_run.error_message,
  792. }
  793. )
  794. return invocations
  795. def _format_capability_plan(
  796. self,
  797. *,
  798. selected_tools: list[AgentToolRefContract],
  799. selected_skills: list[AgentSkillRefContract]) -> str:
  800. return (
  801. "Selected capability plan before model call:\n"
  802. f"Tools: {[item.model_dump(mode='json') for item in selected_tools]}\n"
  803. f"Skills: {[item.model_dump(mode='json') for item in selected_skills]}"
  804. )
  805. def _format_capability_results(
  806. self,
  807. *,
  808. tool_invocations: list[dict[str, JSONValue]],
  809. skill_invocations: list[dict[str, JSONValue]]) -> str | None:
  810. if not tool_invocations and not skill_invocations:
  811. return None
  812. return (
  813. "Capability invocation results before model call:\n"
  814. f"Tools: {tool_invocations}\n"
  815. f"Skills: {skill_invocations}"
  816. )
  817. def _format_react_instruction(
  818. self,
  819. *,
  820. agent_run: AgentRun,
  821. selected_tools: list[AgentToolRefContract],
  822. skill_invocations: list[dict[str, JSONValue]]) -> str:
  823. tool_schemas = self._build_react_tool_schemas(
  824. agent_run=agent_run,
  825. selected_tools=selected_tools)
  826. return (
  827. "Use ReAct JSON only. Respond with one JSON object per turn.\n"
  828. "To call a tool: "
  829. '{"action":"tool","tool_code":"code","input_json":{...}}\n'
  830. "To finish: "
  831. '{"action":"finish","answer":"final answer"}\n'
  832. f"Available tools: {tool_schemas}\n"
  833. f"Pre-run skill results: {skill_invocations}"
  834. )
  835. def _build_react_tool_schemas(
  836. self,
  837. *,
  838. agent_run: AgentRun,
  839. selected_tools: list[AgentToolRefContract]) -> list[dict[str, JSONValue]]:
  840. schemas: list[dict[str, JSONValue]] = []
  841. for ref in selected_tools:
  842. schema: dict[str, JSONValue] = {
  843. "tool_code": ref.tool_code,
  844. "tool_binding_id": ref.tool_binding_id,
  845. "required": ref.required,
  846. "config_json": ref.config_json,
  847. }
  848. if ref.tool_binding_id is not None and self.tool_client is not None:
  849. try:
  850. detail = self.tool_client.get_tool_binding_detail(
  851. binding_id=ref.tool_binding_id)
  852. schema.update(
  853. {
  854. "name": detail.tool_definition.name,
  855. "description": detail.tool_definition.description,
  856. "tool_type": detail.tool_definition.tool_type,
  857. "input_schema_json": detail.tool_version.input_schema_json or {},
  858. "output_schema_json": detail.tool_version.output_schema_json or {},
  859. "timeout_ms": detail.tool_version.timeout_ms,
  860. }
  861. )
  862. except ToolServiceClientError as exc:
  863. schema["schema_error"] = str(exc)
  864. schemas.append(schema)
  865. return schemas
  866. def _invoke_react_tool_with_retry(
  867. self,
  868. *,
  869. agent_run: AgentRun,
  870. agent_version: AgentVersion,
  871. tool_ref: AgentToolRefContract) -> list[dict[str, JSONValue]]:
  872. retry_count = self._read_int(
  873. agent_version.model_config_json,
  874. "react_tool_retry_count",
  875. default=self.react_tool_retry_count)
  876. attempts: list[dict[str, JSONValue]] = []
  877. for attempt_index in range(max(retry_count, 0) + 1):
  878. current = self._invoke_selected_tools(
  879. agent_run=agent_run,
  880. agent_version=agent_version,
  881. selected_tools=[tool_ref])
  882. for item in current:
  883. item["attempt_index"] = attempt_index
  884. attempts.extend(current)
  885. if current and current[-1].get("status") == "completed":
  886. break
  887. return attempts
  888. def _format_react_observation(
  889. self,
  890. tool_invocations: list[dict[str, JSONValue]]) -> str:
  891. return f"Observation: {tool_invocations}"
  892. def _parse_react_action(self, content: str) -> dict[str, JSONValue]:
  893. try:
  894. value = json.loads(content)
  895. except json.JSONDecodeError:
  896. start_index = content.find("{")
  897. end_index = content.rfind("}")
  898. if start_index < 0 or end_index <= start_index:
  899. return {"action": "finish", "answer": content}
  900. try:
  901. value = json.loads(content[start_index : end_index + 1])
  902. except json.JSONDecodeError:
  903. return {"action": "finish", "answer": content}
  904. if not isinstance(value, dict):
  905. return {"action": "finish", "answer": content}
  906. return {str(item_key): item_value for item_key, item_value in value.items()}
  907. def _parse_react_action_from_response(
  908. self,
  909. response: ChatCompletionResponseContract) -> dict[str, JSONValue]:
  910. if response.tool_calls_json:
  911. action = self._parse_openai_tool_call(response.tool_calls_json[0])
  912. if action is not None:
  913. return action
  914. return self._parse_react_action(response.content)
  915. def _parse_openai_tool_call(
  916. self,
  917. tool_call: dict[str, JSONValue]) -> dict[str, JSONValue] | None:
  918. function_value = tool_call.get("function")
  919. if not isinstance(function_value, dict):
  920. return None
  921. tool_code = function_value.get("name")
  922. if not isinstance(tool_code, str) or not tool_code:
  923. return None
  924. raw_arguments = function_value.get("arguments")
  925. input_json: dict[str, JSONValue] = {}
  926. if isinstance(raw_arguments, str) and raw_arguments:
  927. try:
  928. decoded = json.loads(raw_arguments)
  929. except json.JSONDecodeError:
  930. decoded = {"raw_arguments": raw_arguments}
  931. if isinstance(decoded, dict):
  932. input_json = {str(item_key): item_value for item_key, item_value in decoded.items()}
  933. elif isinstance(raw_arguments, dict):
  934. input_json = {str(item_key): item_value for item_key, item_value in raw_arguments.items()}
  935. tool_call_id = tool_call.get("id")
  936. return {
  937. "action": "tool",
  938. "tool_code": tool_code,
  939. "input_json": input_json,
  940. "tool_call_id": tool_call_id if isinstance(tool_call_id, str) else None,
  941. "tool_call_protocol": "openai",
  942. }
  943. def _build_chat_completion_request(
  944. self,
  945. *,
  946. agent_run: AgentRun,
  947. agent_version: AgentVersion,
  948. messages: list[ChatMessageContract],
  949. selected_tools: list[AgentToolRefContract] | None = None) -> ChatCompletionRequestContract:
  950. function_calling_enabled = self._read_bool(
  951. agent_version.model_config_json,
  952. "function_calling_enabled",
  953. default=False) or self._read_bool(
  954. agent_version.model_config_json,
  955. "tool_calling_enabled",
  956. default=False)
  957. return ChatCompletionRequestContract(
  958. model=self._read_optional_string(agent_version.model_config_json, "model"),
  959. temperature=self._read_optional_float(
  960. agent_version.model_config_json,
  961. "temperature"),
  962. max_tokens=self._read_optional_int(
  963. agent_version.model_config_json,
  964. "max_tokens"),
  965. messages=messages,
  966. tools_json=(
  967. self._build_openai_tool_schemas(
  968. agent_run=agent_run,
  969. selected_tools=selected_tools or [])
  970. if function_calling_enabled
  971. else []
  972. ),
  973. tool_choice="auto" if function_calling_enabled and selected_tools else None,
  974. metadata_json={
  975. "agent_id": agent_run.agent_id,
  976. "agent_version_id": agent_version.id,
  977. "agent_run_id": agent_run.id,
  978. })
  979. def _build_openai_tool_schemas(
  980. self,
  981. *,
  982. agent_run: AgentRun,
  983. selected_tools: list[AgentToolRefContract]) -> list[dict[str, JSONValue]]:
  984. tool_schemas: list[dict[str, JSONValue]] = []
  985. for schema in self._build_react_tool_schemas(
  986. agent_run=agent_run,
  987. selected_tools=selected_tools):
  988. tool_code = schema.get("tool_code")
  989. if not isinstance(tool_code, str) or not tool_code:
  990. continue
  991. description = schema.get("description")
  992. input_schema = schema.get("input_schema_json")
  993. tool_schemas.append(
  994. {
  995. "type": "function",
  996. "function": {
  997. "name": tool_code,
  998. "description": description if isinstance(description, str) else "",
  999. "parameters": input_schema if isinstance(input_schema, dict) else {},
  1000. },
  1001. }
  1002. )
  1003. return tool_schemas
  1004. def _build_skill_input_json(
  1005. self,
  1006. *,
  1007. agent_run: AgentRun,
  1008. ref: AgentSkillRefContract) -> dict[str, JSONValue]:
  1009. input_json: dict[str, JSONValue] = dict(agent_run.input_json or {})
  1010. if agent_run.input_text:
  1011. input_json.setdefault("input_text", agent_run.input_text)
  1012. configured_input = ref.config_json.get("input_json")
  1013. if isinstance(configured_input, dict):
  1014. input_json.update(
  1015. {str(item_key): item_value for item_key, item_value in configured_input.items()}
  1016. )
  1017. return input_json
  1018. def _resolve_skill_id_by_code(self, *, skill_code: str | None) -> str | None:
  1019. if skill_code is None or self.skill_client is None:
  1020. return None
  1021. try:
  1022. skills = self.skill_client.list_skills()
  1023. except SkillServiceClientError:
  1024. return None
  1025. for skill in skills:
  1026. if skill.code == skill_code:
  1027. return skill.id
  1028. return None
  1029. def _build_input_preview(self, agent_run: AgentRun) -> str:
  1030. return f"{agent_run.input_text or ''} {agent_run.input_json or {}}".lower()
  1031. def _matches_selection_keywords(
  1032. self,
  1033. config_json: dict[str, JSONValue],
  1034. input_preview: str) -> bool:
  1035. keywords = config_json.get("selection_keywords")
  1036. if not isinstance(keywords, list):
  1037. return False
  1038. return any(
  1039. isinstance(keyword, str) and keyword.lower() in input_preview
  1040. for keyword in keywords
  1041. )
  1042. def _build_dry_run_output(self, *, agent_run: AgentRun, agent_version: AgentVersion) -> str:
  1043. input_preview = agent_run.input_text or str(agent_run.input_json or {})
  1044. return (
  1045. f"[dry-run] Agent role={agent_version.role} "
  1046. f"version={agent_version.version_no} received: {input_preview}"
  1047. )
  1048. def _read_optional_string(self, payload: dict[str, JSONValue], key: str) -> str | None:
  1049. value = payload.get(key)
  1050. if isinstance(value, str) and value:
  1051. return value
  1052. return None
  1053. def _read_optional_float(self, payload: dict[str, JSONValue], key: str) -> float | None:
  1054. value = payload.get(key)
  1055. if isinstance(value, (int, float)) and not isinstance(value, bool):
  1056. return float(value)
  1057. return None
  1058. def _read_optional_int(self, payload: dict[str, JSONValue], key: str) -> int | None:
  1059. value = payload.get(key)
  1060. if isinstance(value, int) and not isinstance(value, bool):
  1061. return value
  1062. return None
  1063. def _read_relevant_memories(
  1064. self,
  1065. *,
  1066. agent_run: AgentRun,
  1067. agent_version: AgentVersion) -> tuple[list[MemorySearchResultContract], dict[str, JSONValue]]:
  1068. if self.memory_client is None:
  1069. return [], {"memory_read_enabled": False, "memory_read_reason": "client_missing"}
  1070. if not self._read_bool(agent_version.memory_policy_json, "enabled", default=True):
  1071. return [], {"memory_read_enabled": False, "memory_read_reason": "policy_disabled"}
  1072. query = agent_run.input_text or str(agent_run.input_json or "")
  1073. if not query:
  1074. return [], {"memory_read_enabled": True, "memory_read_count": 0}
  1075. scope = self._resolve_memory_scope(agent_run=agent_run, agent_version=agent_version)
  1076. if scope is None:
  1077. return [], {
  1078. "memory_read_enabled": True,
  1079. "memory_read_count": 0,
  1080. "memory_read_reason": "scope_unavailable",
  1081. }
  1082. scope_type, scope_id = scope
  1083. try:
  1084. results = self.memory_client.search_memories(
  1085. MemorySearchRequestContract(
  1086. query=query,
  1087. scope_type=scope_type,
  1088. scope_id=scope_id,
  1089. owner_agent_id=agent_run.agent_id,
  1090. session_id=agent_run.session_id,
  1091. limit=self._read_int(
  1092. agent_version.memory_policy_json,
  1093. "read_top_k",
  1094. default=8))
  1095. )
  1096. except MemoryClientError as exc:
  1097. return [], {
  1098. "memory_read_enabled": True,
  1099. "memory_read_count": 0,
  1100. "memory_read_error": str(exc),
  1101. }
  1102. return results, {
  1103. "memory_read_enabled": True,
  1104. "memory_read_count": len(results),
  1105. "memory_scope_type": scope_type,
  1106. "memory_scope_id": scope_id,
  1107. }
  1108. def _write_interaction_memory(
  1109. self,
  1110. *,
  1111. agent_run: AgentRun,
  1112. agent_version: AgentVersion,
  1113. output_text: str) -> dict[str, JSONValue]:
  1114. if self.memory_client is None:
  1115. return {"memory_write_enabled": False, "memory_write_reason": "client_missing"}
  1116. if not self._read_bool(agent_version.memory_policy_json, "write_enabled", default=True):
  1117. return {"memory_write_enabled": False, "memory_write_reason": "policy_disabled"}
  1118. scope = self._resolve_memory_scope(agent_run=agent_run, agent_version=agent_version)
  1119. if scope is None:
  1120. return {"memory_write_enabled": True, "memory_write_reason": "scope_unavailable"}
  1121. scope_type, scope_id = scope
  1122. try:
  1123. memory = self.memory_client.create_memory(
  1124. MemoryCreateContract(
  1125. scope_type=scope_type,
  1126. scope_id=scope_id,
  1127. memory_type="conversation",
  1128. content_text=self._format_interaction_memory(
  1129. agent_run=agent_run,
  1130. output_text=output_text),
  1131. content_json={
  1132. "agent_run_id": agent_run.id,
  1133. "agent_version_id": agent_version.id,
  1134. "input_text": agent_run.input_text,
  1135. "output_text": output_text,
  1136. },
  1137. metadata_json={
  1138. "source": "agent-service",
  1139. "role": agent_version.role,
  1140. "version_no": agent_version.version_no,
  1141. },
  1142. owner_agent_id=agent_run.agent_id,
  1143. session_id=agent_run.session_id,
  1144. source_ref=f"agent_run:{agent_run.id}",
  1145. importance_score=self._read_nested_int(
  1146. agent_version.memory_policy_json,
  1147. "config_json",
  1148. "write_importance_score",
  1149. default=50))
  1150. )
  1151. except MemoryClientError as exc:
  1152. return {
  1153. "memory_write_enabled": True,
  1154. "memory_write_error": str(exc),
  1155. }
  1156. return {
  1157. "memory_write_enabled": True,
  1158. "memory_written_id": memory.id,
  1159. "memory_scope_type": scope_type,
  1160. "memory_scope_id": scope_id,
  1161. }
  1162. def _resolve_memory_scope(
  1163. self,
  1164. *,
  1165. agent_run: AgentRun,
  1166. agent_version: AgentVersion) -> tuple[MemoryScopeType, str] | None:
  1167. scope_value = self._read_optional_string(
  1168. agent_version.memory_policy_json,
  1169. "memory_scope") or "session"
  1170. if scope_value == "global":
  1171. return "global", "global"
  1172. if scope_value == "agent":
  1173. return "agent", agent_run.agent_id
  1174. if scope_value == "session" and agent_run.session_id:
  1175. return "session", agent_run.session_id
  1176. if scope_value == "user":
  1177. user_id = self._read_input_json_string(agent_run=agent_run, key="user_id")
  1178. if user_id is not None:
  1179. return "user", user_id
  1180. if scope_value == "team":
  1181. team_id = self._read_input_json_string(agent_run=agent_run, key="team_id")
  1182. if team_id is not None:
  1183. return "team", team_id
  1184. return None
  1185. def _format_memory_context(self, memory_results: list[MemorySearchResultContract]) -> str:
  1186. lines = ["Relevant memories:"]
  1187. for index, result in enumerate(memory_results, start=1):
  1188. lines.append(f"{index}. {result.item.content_text}")
  1189. return "\n".join(lines)
  1190. def _format_interaction_memory(self, *, agent_run: AgentRun, output_text: str) -> str:
  1191. input_text = agent_run.input_text or str(agent_run.input_json or {})
  1192. return f"User input: {input_text}\nAgent output: {output_text}"
  1193. def _read_bool(self, payload: dict[str, JSONValue], key: str, *, default: bool) -> bool:
  1194. value = payload.get(key)
  1195. if isinstance(value, bool):
  1196. return value
  1197. return default
  1198. def _read_int(self, payload: dict[str, JSONValue], key: str, *, default: int) -> int:
  1199. value = payload.get(key)
  1200. if isinstance(value, int) and not isinstance(value, bool):
  1201. return value
  1202. return default
  1203. def _read_nested_int(
  1204. self,
  1205. payload: dict[str, JSONValue],
  1206. parent_key: str,
  1207. child_key: str,
  1208. *,
  1209. default: int) -> int:
  1210. parent_value = payload.get(parent_key)
  1211. if not isinstance(parent_value, dict):
  1212. return default
  1213. return self._read_int(
  1214. cast(dict[str, JSONValue], parent_value),
  1215. child_key,
  1216. default=default)
  1217. def _read_input_json_string(self, *, agent_run: AgentRun, key: str) -> str | None:
  1218. if agent_run.input_json is None:
  1219. return None
  1220. value = agent_run.input_json.get(key)
  1221. if isinstance(value, str) and value:
  1222. return value
  1223. return None
  1224. def build_agent_application_service(
  1225. *,
  1226. db: Session,
  1227. settings: AgentServiceSettings) -> AgentApplicationService:
  1228. redis_client = try_build_redis_client(settings.redis_url)
  1229. return AgentApplicationService(
  1230. agent_repository=AgentDefinitionRepository(db),
  1231. agent_version_repository=AgentVersionRepository(db),
  1232. agent_run_repository=AgentRunRepository(db),
  1233. agent_tool_invocation_repository=AgentToolInvocationRepository(db),
  1234. model_gateway_client=ModelGatewayClient(
  1235. base_url=settings.model_gateway_service_url,
  1236. timeout_seconds=settings.model_gateway_timeout_seconds),
  1237. memory_client=MemoryClient(
  1238. base_url=settings.memory_service_url,
  1239. timeout_seconds=settings.memory_service_timeout_seconds),
  1240. tool_client=ToolServiceClient(
  1241. base_url=settings.tool_service_url,
  1242. timeout_seconds=settings.tool_service_timeout_seconds),
  1243. skill_client=SkillServiceClient(
  1244. base_url=settings.skill_service_url,
  1245. timeout_seconds=settings.skill_service_timeout_seconds),
  1246. event_client=EventServiceClient(
  1247. base_url=settings.event_service_url,
  1248. timeout_seconds=settings.event_service_timeout_seconds),
  1249. task_queue_publisher=(
  1250. TaskQueuePublisher(client=redis_client) if redis_client is not None else None
  1251. ),
  1252. react_max_steps=settings.react_max_steps)