routes.py 16 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395
  1. from datetime import datetime
  2. from typing import TypeVar
  3. from core_domain import ServiceHealth
  4. from fastapi import APIRouter, Depends, HTTPException, Query
  5. from sqlalchemy import text
  6. from sqlalchemy.orm import Session
  7. from app.application.services import TeamApplicationService, build_team_application_service
  8. from app.bootstrap.settings import TeamServiceSettings
  9. from app.db.session import get_db
  10. from app.domain.repositories import (
  11. TeamDefinitionRepository,
  12. TeamRunRepository,
  13. TeamVersionRepository,
  14. )
  15. from app.schemas.team import (
  16. ApiResponse,
  17. DeleteData,
  18. PageResult,
  19. TeamConfigCreateRequestDto,
  20. TeamConfigDeleteRequestDto,
  21. TeamConfigDetailRequestDto,
  22. TeamConfigDto,
  23. TeamConfigListRequestDto,
  24. TeamConfigUpdateRequestDto,
  25. TeamCreateRequest,
  26. TeamCreateRequestDto,
  27. TeamDeleteRequestDto,
  28. TeamDetailRequestDto,
  29. TeamDto,
  30. TeamListRequestDto,
  31. TeamResponse,
  32. TeamRunCreateRequest,
  33. TeamRunCreateRequestDto,
  34. TeamRunDeleteRequestDto,
  35. TeamRunDetailRequestDto,
  36. TeamRunDto,
  37. TeamRunExecuteRequest,
  38. TeamRunExecuteData,
  39. TeamRunExecuteRequestDto,
  40. TeamRunExecuteResponse,
  41. TeamRunListRequestDto,
  42. TeamRunResponse,
  43. TeamRunStatusUpdateRequest,
  44. TeamRunStatusUpdateRequestDto,
  45. TeamStatusUpdateRequest,
  46. TeamUpdateRequestDto,
  47. TeamVersionCreateRequest,
  48. TeamVersionResponse,
  49. TeamWorkerExecuteNextRequest,
  50. TeamWorkerExecuteNextResponse,
  51. )
  52. router = APIRouter()
  53. T = TypeVar("T")
  54. def ok(data: T) -> ApiResponse[T]:
  55. return ApiResponse[T](
  56. data=data,
  57. requestId="",
  58. serverTime=datetime.utcnow())
  59. def get_team_settings() -> TeamServiceSettings:
  60. return TeamServiceSettings()
  61. def get_team_application_service(
  62. db: Session = Depends(get_db),
  63. settings: TeamServiceSettings = Depends(get_team_settings)) -> TeamApplicationService:
  64. return build_team_application_service(
  65. team_repository=TeamDefinitionRepository(db),
  66. team_version_repository=TeamVersionRepository(db),
  67. team_run_repository=TeamRunRepository(db),
  68. settings=settings)
  69. @router.get("/health", response_model=ServiceHealth)
  70. def health_check(db: Session = Depends(get_db)) -> ServiceHealth:
  71. db.execute(text("SELECT 1"))
  72. return ServiceHealth(service="team-service", status="ok", database="ok")
  73. @router.post("", response_model=TeamResponse)
  74. def create_team(
  75. payload: TeamCreateRequest,
  76. service: TeamApplicationService = Depends(get_team_application_service)) -> TeamResponse:
  77. entity = service.create_team(payload)
  78. return TeamResponse.from_entity(entity)
  79. @router.get("", response_model=list[TeamResponse])
  80. def list_teams(
  81. service: TeamApplicationService = Depends(get_team_application_service)) -> list[TeamResponse]:
  82. return [TeamResponse.from_entity(item) for item in service.list_teams()]
  83. @router.post("/list", response_model=ApiResponse[PageResult[TeamDto]])
  84. def list_teams_contract(
  85. payload: TeamListRequestDto,
  86. service: TeamApplicationService = Depends(get_team_application_service)) -> ApiResponse[PageResult[TeamDto]]:
  87. keyword = (payload.keyword or "").lower().strip()
  88. items = [
  89. item
  90. for item in service.list_teams()
  91. if (payload.status is None or item.status == payload.status)
  92. and (
  93. not keyword
  94. or keyword in item.name.lower()
  95. or keyword in item.team_type.lower()
  96. or keyword in (item.description or "").lower()
  97. )
  98. ]
  99. page_items = items[payload.offset:payload.offset + payload.pageSize]
  100. return ok(PageResult[TeamDto].from_items(
  101. items=[TeamDto.from_entity(item) for item in page_items],
  102. total=len(items),
  103. page=payload.page,
  104. page_size=payload.pageSize))
  105. @router.post("/create", response_model=ApiResponse[TeamDto])
  106. def create_team_contract(
  107. payload: TeamCreateRequestDto,
  108. service: TeamApplicationService = Depends(get_team_application_service)) -> ApiResponse[TeamDto]:
  109. return ok(TeamDto.from_entity(service.create_team_from_contract(payload)))
  110. @router.post("/detail", response_model=ApiResponse[TeamDto])
  111. def get_team_contract(
  112. payload: TeamDetailRequestDto,
  113. service: TeamApplicationService = Depends(get_team_application_service)) -> ApiResponse[TeamDto]:
  114. entity = service.get_team(team_id=payload.teamId)
  115. if entity is None:
  116. raise HTTPException(status_code=404, detail=f"team not found: {payload.teamId}")
  117. return ok(TeamDto.from_entity(entity))
  118. @router.post("/update", response_model=ApiResponse[TeamDto])
  119. def update_team_contract(
  120. payload: TeamUpdateRequestDto,
  121. service: TeamApplicationService = Depends(get_team_application_service)) -> ApiResponse[TeamDto]:
  122. entity = service.update_team_from_contract(payload)
  123. if entity is None:
  124. raise HTTPException(status_code=404, detail=f"team not found: {payload.teamId}")
  125. return ok(TeamDto.from_entity(entity))
  126. @router.post("/delete", response_model=ApiResponse[DeleteData])
  127. def delete_team_contract(
  128. payload: TeamDeleteRequestDto,
  129. service: TeamApplicationService = Depends(get_team_application_service)) -> ApiResponse[DeleteData]:
  130. deleted = service.delete_team_from_contract(payload)
  131. return ok(DeleteData(deleted=deleted, teamId=payload.teamId))
  132. @router.patch("/{team_id}/status", response_model=TeamResponse)
  133. def update_team_status(
  134. team_id: str,
  135. payload: TeamStatusUpdateRequest,
  136. service: TeamApplicationService = Depends(get_team_application_service)) -> TeamResponse:
  137. entity = service.update_team_status(team_id=team_id, payload=payload)
  138. if entity is None:
  139. raise HTTPException(status_code=404, detail=f"team not found: {team_id}")
  140. return TeamResponse.from_entity(entity)
  141. @router.post("/versions", response_model=TeamVersionResponse)
  142. def create_team_version(
  143. payload: TeamVersionCreateRequest,
  144. service: TeamApplicationService = Depends(get_team_application_service)) -> TeamVersionResponse:
  145. try:
  146. entity = service.create_team_version(payload)
  147. except ValueError as exc:
  148. raise HTTPException(status_code=422, detail=str(exc)) from exc
  149. return TeamVersionResponse.from_entity(entity)
  150. @router.get("/versions", response_model=list[TeamVersionResponse])
  151. def list_team_versions(
  152. team_id: str = Query(...),
  153. service: TeamApplicationService = Depends(get_team_application_service)) -> list[TeamVersionResponse]:
  154. return [
  155. TeamVersionResponse.from_entity(item)
  156. for item in service.list_team_versions(team_id=team_id)
  157. ]
  158. @router.post("/configs/list", response_model=ApiResponse[PageResult[TeamConfigDto]])
  159. def list_team_configs_contract(
  160. payload: TeamConfigListRequestDto,
  161. service: TeamApplicationService = Depends(get_team_application_service)) -> ApiResponse[PageResult[TeamConfigDto]]:
  162. items = service.list_team_configs(team_id=payload.teamId)
  163. page_items = items[payload.offset:payload.offset + payload.pageSize]
  164. return ok(PageResult[TeamConfigDto].from_items(
  165. items=[TeamConfigDto.from_entity(item) for item in page_items],
  166. total=len(items),
  167. page=payload.page,
  168. page_size=payload.pageSize))
  169. @router.post("/configs/create", response_model=ApiResponse[TeamConfigDto])
  170. def create_team_config_contract(
  171. payload: TeamConfigCreateRequestDto,
  172. service: TeamApplicationService = Depends(get_team_application_service)) -> ApiResponse[TeamConfigDto]:
  173. try:
  174. entity = service.create_team_config_from_contract(payload)
  175. except ValueError as exc:
  176. raise HTTPException(status_code=422, detail=str(exc)) from exc
  177. return ok(TeamConfigDto.from_entity(entity))
  178. @router.post("/configs/detail", response_model=ApiResponse[TeamConfigDto])
  179. def get_team_config_contract(
  180. payload: TeamConfigDetailRequestDto,
  181. service: TeamApplicationService = Depends(get_team_application_service)) -> ApiResponse[TeamConfigDto]:
  182. entity = service.get_team_config(config_id=payload.configId)
  183. if entity is None:
  184. raise HTTPException(status_code=404, detail=f"team config not found: {payload.configId}")
  185. return ok(TeamConfigDto.from_entity(entity))
  186. @router.post("/configs/update", response_model=ApiResponse[TeamConfigDto])
  187. def update_team_config_contract(
  188. payload: TeamConfigUpdateRequestDto,
  189. service: TeamApplicationService = Depends(get_team_application_service)) -> ApiResponse[TeamConfigDto]:
  190. try:
  191. entity = service.update_team_config_from_contract(payload)
  192. except ValueError as exc:
  193. raise HTTPException(status_code=422, detail=str(exc)) from exc
  194. if entity is None:
  195. raise HTTPException(status_code=404, detail=f"team config not found: {payload.configId}")
  196. return ok(TeamConfigDto.from_entity(entity))
  197. @router.post("/configs/delete", response_model=ApiResponse[DeleteData])
  198. def delete_team_config_contract(
  199. payload: TeamConfigDeleteRequestDto,
  200. service: TeamApplicationService = Depends(get_team_application_service)) -> ApiResponse[DeleteData]:
  201. deleted = service.delete_team_config(config_id=payload.configId)
  202. return ok(DeleteData(deleted=deleted, configId=payload.configId))
  203. @router.post("/runs", response_model=TeamRunResponse)
  204. def create_team_run(
  205. payload: TeamRunCreateRequest,
  206. service: TeamApplicationService = Depends(get_team_application_service)) -> TeamRunResponse:
  207. try:
  208. entity = service.create_team_run(payload)
  209. except ValueError as exc:
  210. raise HTTPException(status_code=422, detail=str(exc)) from exc
  211. return TeamRunResponse.from_entity(entity)
  212. @router.get("/runs", response_model=list[TeamRunResponse])
  213. def list_team_runs(
  214. team_id: str | None = Query(default=None),
  215. session_id: str | None = Query(default=None),
  216. service: TeamApplicationService = Depends(get_team_application_service)) -> list[TeamRunResponse]:
  217. return [
  218. TeamRunResponse.from_entity(item)
  219. for item in service.list_team_runs(
  220. team_id=team_id,
  221. session_id=session_id)
  222. ]
  223. @router.post("/runs/list", response_model=ApiResponse[PageResult[TeamRunDto]])
  224. def list_team_runs_contract(
  225. payload: TeamRunListRequestDto,
  226. service: TeamApplicationService = Depends(get_team_application_service)) -> ApiResponse[PageResult[TeamRunDto]]:
  227. items = [
  228. item
  229. for item in service.list_team_runs(team_id=payload.teamId, session_id=payload.sessionId)
  230. if payload.status is None or item.status == payload.status
  231. ]
  232. page_items = items[payload.offset:payload.offset + payload.pageSize]
  233. return ok(PageResult[TeamRunDto].from_items(
  234. items=[TeamRunDto.from_entity(item) for item in page_items],
  235. total=len(items),
  236. page=payload.page,
  237. page_size=payload.pageSize))
  238. @router.post("/runs/create", response_model=ApiResponse[TeamRunDto])
  239. def create_team_run_contract(
  240. payload: TeamRunCreateRequestDto,
  241. service: TeamApplicationService = Depends(get_team_application_service)) -> ApiResponse[TeamRunDto]:
  242. try:
  243. entity = service.create_team_run_from_contract(payload)
  244. except ValueError as exc:
  245. raise HTTPException(status_code=422, detail=str(exc)) from exc
  246. return ok(TeamRunDto.from_entity(entity))
  247. @router.post("/runs/detail", response_model=ApiResponse[TeamRunDto])
  248. def get_team_run_contract(
  249. payload: TeamRunDetailRequestDto,
  250. service: TeamApplicationService = Depends(get_team_application_service)) -> ApiResponse[TeamRunDto]:
  251. entity = service.get_team_run(team_run_id=payload.teamRunId)
  252. if entity is None:
  253. raise HTTPException(status_code=404, detail=f"team_run not found: {payload.teamRunId}")
  254. return ok(TeamRunDto.from_entity(entity))
  255. @router.post("/runs/status", response_model=ApiResponse[TeamRunDto])
  256. def update_team_run_status_contract(
  257. payload: TeamRunStatusUpdateRequestDto,
  258. service: TeamApplicationService = Depends(get_team_application_service)) -> ApiResponse[TeamRunDto]:
  259. entity = service.update_team_run_status_from_contract(payload)
  260. if entity is None:
  261. raise HTTPException(status_code=404, detail=f"team_run not found: {payload.teamRunId}")
  262. return ok(TeamRunDto.from_entity(entity))
  263. @router.post("/runs/execute", response_model=ApiResponse[TeamRunExecuteData])
  264. def execute_team_run_contract(
  265. payload: TeamRunExecuteRequestDto,
  266. service: TeamApplicationService = Depends(get_team_application_service)) -> ApiResponse[TeamRunExecuteData]:
  267. entity = service.execute_team_run(
  268. team_run_id=payload.teamRunId,
  269. payload=TeamRunExecuteRequest(
  270. worker_key=payload.workerKey,
  271. dry_run=payload.dryRun))
  272. if entity is None:
  273. raise HTTPException(status_code=404, detail=f"team_run not found: {payload.teamRunId}")
  274. output_json = entity.output_json or {}
  275. member_run_count = output_json.get("member_run_count")
  276. dry_run = output_json.get("dry_run")
  277. return ok(TeamRunExecuteData(
  278. run=TeamRunDto.from_entity(entity),
  279. memberRunCount=member_run_count if isinstance(member_run_count, int) else 0,
  280. dryRun=dry_run if isinstance(dry_run, bool) else payload.dryRun))
  281. @router.post("/runs/delete", response_model=ApiResponse[DeleteData])
  282. def delete_team_run_contract(
  283. payload: TeamRunDeleteRequestDto,
  284. service: TeamApplicationService = Depends(get_team_application_service)) -> ApiResponse[DeleteData]:
  285. deleted = service.delete_team_run(team_run_id=payload.teamRunId)
  286. return ok(DeleteData(deleted=deleted, teamRunId=payload.teamRunId))
  287. @router.post("/runs/{team_run_id}/status", response_model=TeamRunResponse)
  288. def update_team_run_status(
  289. team_run_id: str,
  290. payload: TeamRunStatusUpdateRequest,
  291. service: TeamApplicationService = Depends(get_team_application_service)) -> TeamRunResponse:
  292. entity = service.update_team_run_status(team_run_id=team_run_id, payload=payload)
  293. if entity is None:
  294. raise HTTPException(status_code=404, detail=f"team_run not found: {team_run_id}")
  295. return TeamRunResponse.from_entity(entity)
  296. @router.post("/runs/{team_run_id}/execute", response_model=TeamRunExecuteResponse)
  297. def execute_team_run(
  298. team_run_id: str,
  299. payload: TeamRunExecuteRequest,
  300. service: TeamApplicationService = Depends(get_team_application_service)) -> TeamRunExecuteResponse:
  301. entity = service.execute_team_run(team_run_id=team_run_id, payload=payload)
  302. if entity is None:
  303. raise HTTPException(status_code=404, detail=f"team_run not found: {team_run_id}")
  304. output_json = entity.output_json or {}
  305. member_run_count = output_json.get("member_run_count")
  306. dry_run = output_json.get("dry_run")
  307. return TeamRunExecuteResponse(
  308. run=TeamRunResponse.from_entity(entity),
  309. member_run_count=member_run_count if isinstance(member_run_count, int) else 0,
  310. dry_run=dry_run if isinstance(dry_run, bool) else payload.dry_run)
  311. @router.post("/workers/execute-next", response_model=TeamWorkerExecuteNextResponse)
  312. def execute_next_worker_task(
  313. payload: TeamWorkerExecuteNextRequest,
  314. settings: TeamServiceSettings = Depends(get_team_settings),
  315. service: TeamApplicationService = Depends(get_team_application_service)) -> TeamWorkerExecuteNextResponse:
  316. result = service.execute_next_claimed_team_run(
  317. worker_key=payload.worker_key,
  318. lease_seconds=payload.lease_seconds or settings.worker_lease_seconds,
  319. stale_running_seconds=settings.worker_stale_running_seconds,
  320. dry_run=payload.dry_run if payload.dry_run is not None else settings.worker_dry_run)
  321. if result is None:
  322. raise HTTPException(status_code=404, detail="queued team_run not found")
  323. entity, released_lease_count = result
  324. output_json = entity.output_json or {}
  325. member_run_count = output_json.get("member_run_count")
  326. dry_run = output_json.get("dry_run")
  327. return TeamWorkerExecuteNextResponse(
  328. run=TeamRunResponse.from_entity(entity),
  329. member_run_count=member_run_count if isinstance(member_run_count, int) else 0,
  330. dry_run=dry_run if isinstance(dry_run, bool) else settings.worker_dry_run,
  331. released_lease_count=released_lease_count)