repositories.py 8.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288
  1. from datetime import datetime, timedelta
  2. from sqlalchemy import func, select
  3. from sqlalchemy.orm import Session
  4. from core_domain import TeamRunStatus, TeamStatus, TeamVersionStatus
  5. from core_shared import JSONValue
  6. from app.db.models import TeamDefinition, TeamRun, TeamVersion
  7. class TeamDefinitionRepository:
  8. def __init__(self, db: Session) -> None:
  9. self.db = db
  10. def create(
  11. self,
  12. *,
  13. code: str,
  14. name: str,
  15. description: str | None,
  16. team_type: str,
  17. owner_user_id: str | None,
  18. metadata_json: dict[str, JSONValue] | None) -> TeamDefinition:
  19. entity = TeamDefinition(
  20. code=code,
  21. name=name,
  22. description=description,
  23. team_type=team_type,
  24. owner_user_id=owner_user_id,
  25. metadata_json=metadata_json)
  26. self.db.add(entity)
  27. self.db.commit()
  28. self.db.refresh(entity)
  29. return entity
  30. def list_all(self) -> list[TeamDefinition]:
  31. stmt = (
  32. select(TeamDefinition)
  33. .order_by(TeamDefinition.created_time.desc())
  34. )
  35. return list(self.db.scalars(stmt))
  36. def get_by_id(self, *, team_id: str) -> TeamDefinition | None:
  37. stmt = (
  38. select(TeamDefinition)
  39. .where(TeamDefinition.id == team_id)
  40. )
  41. return self.db.scalar(stmt)
  42. def save(self, entity: TeamDefinition) -> TeamDefinition:
  43. self.db.add(entity)
  44. self.db.commit()
  45. self.db.refresh(entity)
  46. return entity
  47. def delete(self, entity: TeamDefinition) -> None:
  48. self.db.delete(entity)
  49. self.db.commit()
  50. def update_status(
  51. self,
  52. *,
  53. team_id: str,
  54. status: TeamStatus) -> TeamDefinition | None:
  55. entity = self.get_by_id(team_id=team_id)
  56. if entity is None:
  57. return None
  58. entity.status = status
  59. self.db.commit()
  60. self.db.refresh(entity)
  61. return entity
  62. class TeamVersionRepository:
  63. def __init__(self, db: Session) -> None:
  64. self.db = db
  65. def create(
  66. self,
  67. *,
  68. team_id: str,
  69. status: TeamVersionStatus,
  70. coordination_mode: str,
  71. objective: str | None,
  72. member_refs_json: list[dict[str, JSONValue]],
  73. policy_json: dict[str, JSONValue]) -> TeamVersion:
  74. version_no = self._next_version_no(team_id)
  75. entity = TeamVersion(
  76. team_id=team_id,
  77. version_no=version_no,
  78. status=status,
  79. coordination_mode=coordination_mode,
  80. objective=objective,
  81. member_refs_json=member_refs_json,
  82. policy_json=policy_json,
  83. published_time=datetime.utcnow() if status == "published" else None)
  84. self.db.add(entity)
  85. self.db.commit()
  86. self.db.refresh(entity)
  87. return entity
  88. def list_by_team(self, *, team_id: str) -> list[TeamVersion]:
  89. stmt = (
  90. select(TeamVersion)
  91. .where(TeamVersion.team_id == team_id)
  92. .order_by(TeamVersion.version_no.desc())
  93. )
  94. return list(self.db.scalars(stmt))
  95. def list_all(self) -> list[TeamVersion]:
  96. stmt = select(TeamVersion).order_by(TeamVersion.created_time.desc())
  97. return list(self.db.scalars(stmt))
  98. def get_by_id(self, *, team_version_id: str) -> TeamVersion | None:
  99. stmt = (
  100. select(TeamVersion)
  101. .where(TeamVersion.id == team_version_id)
  102. )
  103. return self.db.scalar(stmt)
  104. def save(self, entity: TeamVersion) -> TeamVersion:
  105. self.db.add(entity)
  106. self.db.commit()
  107. self.db.refresh(entity)
  108. return entity
  109. def delete(self, entity: TeamVersion) -> None:
  110. self.db.delete(entity)
  111. self.db.commit()
  112. def get_latest_published(self, *, team_id: str) -> TeamVersion | None:
  113. stmt = (
  114. select(TeamVersion)
  115. .where(TeamVersion.team_id == team_id)
  116. .where(TeamVersion.status == "published")
  117. .order_by(TeamVersion.version_no.desc())
  118. .limit(1)
  119. )
  120. return self.db.scalar(stmt)
  121. def _next_version_no(self, team_id: str) -> int:
  122. stmt = select(func.max(TeamVersion.version_no)).where(TeamVersion.team_id == team_id)
  123. current_max = self.db.scalar(stmt)
  124. return (current_max or 0) + 1
  125. class TeamRunRepository:
  126. def __init__(self, db: Session) -> None:
  127. self.db = db
  128. def create(
  129. self,
  130. *,
  131. team_id: str,
  132. team_version_id: str,
  133. session_id: str | None,
  134. input_text: str | None,
  135. input_json: dict[str, JSONValue] | None) -> TeamRun:
  136. now = datetime.utcnow()
  137. entity = TeamRun(
  138. team_id=team_id,
  139. team_version_id=team_version_id,
  140. session_id=session_id,
  141. input_text=input_text,
  142. input_json=input_json,
  143. status="queued",
  144. queued_time=now)
  145. self.db.add(entity)
  146. self.db.commit()
  147. self.db.refresh(entity)
  148. return entity
  149. def list_by_scope(
  150. self,
  151. *,
  152. team_id: str | None = None,
  153. session_id: str | None = None) -> list[TeamRun]:
  154. stmt = select(TeamRun)
  155. if team_id is not None:
  156. stmt = stmt.where(TeamRun.team_id == team_id)
  157. if session_id is not None:
  158. stmt = stmt.where(TeamRun.session_id == session_id)
  159. stmt = stmt.order_by(TeamRun.created_time.desc())
  160. return list(self.db.scalars(stmt))
  161. def get_by_id(self, *, team_run_id: str) -> TeamRun | None:
  162. stmt = (
  163. select(TeamRun)
  164. .where(TeamRun.id == team_run_id)
  165. )
  166. return self.db.scalar(stmt)
  167. def delete(self, entity: TeamRun) -> None:
  168. self.db.delete(entity)
  169. self.db.commit()
  170. def claim_next_queued(
  171. self,
  172. *,
  173. worker_key: str,
  174. lease_expire_time: datetime) -> TeamRun | None:
  175. stmt = (
  176. select(TeamRun)
  177. .where(TeamRun.status == "queued")
  178. .order_by(TeamRun.created_time.asc())
  179. .with_for_update(skip_locked=True)
  180. .limit(1)
  181. )
  182. entity = self.db.scalar(stmt)
  183. if entity is None:
  184. return None
  185. now = datetime.utcnow()
  186. entity.status = "running"
  187. entity.worker_key = worker_key
  188. entity.started_time = entity.started_time or now
  189. entity.lease_expire_time = lease_expire_time
  190. self.db.commit()
  191. self.db.refresh(entity)
  192. return entity
  193. def release_expired_leases(
  194. self,
  195. *,
  196. now_time: datetime,
  197. stale_running_seconds: int,
  198. max_items: int = 100) -> int:
  199. stale_started_before = now_time - timedelta(seconds=stale_running_seconds)
  200. stmt = (
  201. select(TeamRun)
  202. .where(TeamRun.status == "running")
  203. .where(
  204. (TeamRun.lease_expire_time.is_not(None) & (TeamRun.lease_expire_time <= now_time))
  205. | (
  206. TeamRun.lease_expire_time.is_(None)
  207. & TeamRun.started_time.is_not(None)
  208. & (TeamRun.started_time <= stale_started_before)
  209. )
  210. )
  211. .order_by(TeamRun.lease_expire_time.asc())
  212. .limit(max_items)
  213. )
  214. entities = list(self.db.scalars(stmt))
  215. for entity in entities:
  216. entity.status = "queued"
  217. entity.worker_key = None
  218. entity.lease_expire_time = None
  219. entity.queued_time = now_time
  220. entity.started_time = None
  221. entity.finished_time = None
  222. if entities:
  223. self.db.commit()
  224. return len(entities)
  225. def update_status(
  226. self,
  227. *,
  228. team_run_id: str,
  229. status: TeamRunStatus,
  230. worker_key: str | None = None,
  231. output_text: str | None = None,
  232. output_json: dict[str, JSONValue] | None = None,
  233. error_code: str | None = None,
  234. error_message: str | None = None) -> TeamRun | None:
  235. entity = self.db.get(TeamRun, team_run_id)
  236. if entity is None:
  237. return None
  238. now = datetime.utcnow()
  239. entity.status = status
  240. entity.worker_key = worker_key
  241. entity.output_text = output_text
  242. entity.output_json = output_json
  243. entity.error_code = error_code
  244. entity.error_message = error_message
  245. if status == "running" and entity.started_time is None:
  246. entity.started_time = now
  247. if status != "running":
  248. entity.lease_expire_time = None
  249. if status in {"completed", "failed", "cancelled"}:
  250. entity.finished_time = now
  251. self.db.commit()
  252. self.db.refresh(entity)
  253. return entity