Coverage for src/local_deep_research/web/queue/lifecycle_cleanup.py: 93%

41 statements  

« prev     ^ index     » next       coverage.py v7.15.1, created at 2026-07-19 23:35 +0000

1from collections.abc import Collection 

2from dataclasses import dataclass 

3from datetime import UTC, datetime 

4 

5from sqlalchemy.orm import Session 

6 

7from ...database.models import QueuedResearch, QueueStatus, TaskMetadata 

8 

9 

10@dataclass(frozen=True, slots=True) 

11class QueuedResearchCleanupResult: 

12 cleaned_ids: frozenset[str] 

13 protected_ids: frozenset[str] 

14 

15 

16def reconcile_research_queue_status( 

17 db_session: Session, *, force: bool = False 

18) -> bool: 

19 """Recompute research capacity counters without committing the transaction.""" 

20 queued_tasks = ( 

21 db_session.query(TaskMetadata) 

22 .filter_by(task_type="research", status="queued") 

23 .count() 

24 ) 

25 active_tasks = ( 

26 db_session.query(TaskMetadata) 

27 .filter_by(task_type="research", status="processing") 

28 .count() 

29 ) 

30 queue_status = db_session.query(QueueStatus).first() 

31 if queue_status is None: 

32 db_session.add( 

33 QueueStatus( 

34 active_tasks=active_tasks, 

35 queued_tasks=queued_tasks, 

36 last_checked=datetime.now(UTC), 

37 ) 

38 ) 

39 return True 

40 

41 changed = ( 

42 queue_status.active_tasks != active_tasks 

43 or queue_status.queued_tasks != queued_tasks 

44 ) 

45 if changed or force: 

46 queue_status.active_tasks = active_tasks 

47 queue_status.queued_tasks = queued_tasks 

48 queue_status.last_checked = datetime.now(UTC) 

49 return changed or force 

50 

51 

52def cleanup_queued_research_state( 

53 db_session: Session, 

54 research_ids: Collection[str], 

55 include_claimed: bool = False, 

56) -> QueuedResearchCleanupResult: 

57 """Remove queued research state while leaving transaction control to the caller.""" 

58 research_id_set = set(research_ids) 

59 if not research_id_set: 59 ↛ 60line 59 didn't jump to line 60 because the condition on line 59 was never true

60 return QueuedResearchCleanupResult(frozenset(), frozenset()) 

61 

62 claimed_ids: set[str] = set() 

63 if not include_claimed: 

64 claimed_ids = { 

65 research_id 

66 for (research_id,) in ( 

67 db_session.query(QueuedResearch.research_id) 

68 .filter( 

69 QueuedResearch.research_id.in_(research_id_set), 

70 QueuedResearch.is_processing.is_(True), 

71 ) 

72 .all() 

73 ) 

74 } 

75 cleanup_ids = research_id_set - claimed_ids 

76 if not cleanup_ids: 76 ↛ 77line 76 didn't jump to line 77 because the condition on line 76 was never true

77 return QueuedResearchCleanupResult(frozenset(), frozenset(claimed_ids)) 

78 

79 queued_rows = ( 

80 db_session.query(QueuedResearch) 

81 .filter(QueuedResearch.research_id.in_(cleanup_ids)) 

82 .all() 

83 ) 

84 usernames = {queued_row.username for queued_row in queued_rows} 

85 for queued_row in queued_rows: 

86 db_session.delete(queued_row) 

87 

88 ( 

89 db_session.query(TaskMetadata) 

90 .filter( 

91 TaskMetadata.task_id.in_(cleanup_ids), 

92 TaskMetadata.task_type == "research", 

93 ) 

94 .delete(synchronize_session=False) 

95 ) 

96 

97 for username in usernames: 

98 surviving_rows = ( 

99 db_session.query(QueuedResearch) 

100 .filter_by(username=username) 

101 .order_by(QueuedResearch.position, QueuedResearch.id) 

102 .all() 

103 ) 

104 for position, queued_row in enumerate(surviving_rows, start=1): 

105 queued_row.position = position 

106 

107 reconcile_research_queue_status(db_session, force=True) 

108 return QueuedResearchCleanupResult( 

109 frozenset(cleanup_ids), frozenset(claimed_ids) 

110 )