Coverage for /home/admin/Documents/AI/applications/lexigram-dev/lexigram/experimental/ai/lexigram-ai-workers/src/lexigram/ai/workers/types.py: 80%

108 statements  

« prev     ^ index     » next       coverage.py v7.15.4, created at 2026-08-25 07:19 +0800

1"""Type definitions for AI Workers — enums and dataclasses used across subpackages.""" 

2 

3from __future__ import annotations 

4 

5from dataclasses import dataclass, field 

6from datetime import UTC, datetime 

7from enum import StrEnum 

8from typing import TYPE_CHECKING, Any 

9 

10if TYPE_CHECKING: 

11 from collections.abc import Callable 

12 

13 from lexigram.contracts import JobProtocol 

14 

15 

16# ── DLQ types ──────────────────────────────────────────────────────── 

17 

18 

19class FailureCategory(StrEnum): 

20 """Categories of failures.""" 

21 

22 TRANSIENT = "transient" 

23 PERMANENT = "permanent" 

24 THROTTLED = "throttled" 

25 INVALID_INPUT = "invalid_input" 

26 UNKNOWN = "unknown" 

27 

28 

29class DLQAction(StrEnum): 

30 """Actions for DLQ items.""" 

31 

32 RETRY = "retry" 

33 ARCHIVE = "archive" 

34 DELETE = "delete" 

35 NOTIFY = "notify" 

36 

37 

38@dataclass 

39class DLQItem: 

40 """Item in the dead letter queue.""" 

41 

42 job_id: str 

43 original_job: JobProtocol 

44 failure_count: int 

45 first_failure: datetime 

46 last_failure: datetime 

47 last_error: str 

48 failure_category: FailureCategory = FailureCategory.UNKNOWN 

49 retry_count: int = 0 

50 max_retries: int = 5 

51 next_retry: datetime | None = None 

52 metadata: dict[str, Any] = field(default_factory=dict) 

53 

54 def can_retry(self) -> bool: 

55 """Check if item can be retried.""" 

56 if self.retry_count >= self.max_retries: 

57 return False 

58 if self.failure_category == FailureCategory.PERMANENT: 

59 return False 

60 return not (self.next_retry and datetime.now(UTC) < self.next_retry) 

61 

62 def calculate_backoff(self, base_delay: int = 60) -> int: 

63 """Calculate exponential backoff delay in seconds.""" 

64 return min(base_delay * (2**self.retry_count), 3600) 

65 

66 def to_dict(self) -> dict[str, Any]: 

67 """Convert to dictionary for serialization.""" 

68 return { 

69 "job_id": self.job_id, 

70 "failure_count": self.failure_count, 

71 "first_failure": self.first_failure.isoformat(), 

72 "last_failure": self.last_failure.isoformat(), 

73 "last_error": self.last_error, 

74 "failure_category": self.failure_category.value, 

75 "retry_count": self.retry_count, 

76 "max_retries": self.max_retries, 

77 "next_retry": self.next_retry.isoformat() if self.next_retry else None, 

78 "can_retry": self.can_retry(), 

79 "metadata": self.metadata, 

80 } 

81 

82 

83@dataclass 

84class DLQStats: 

85 """Statistics for DLQ.""" 

86 

87 total_items: int = 0 

88 by_category: dict[str, int] = field(default_factory=dict) 

89 retried_count: int = 0 

90 archived_count: int = 0 

91 deleted_count: int = 0 

92 permanent_failures: int = 0 

93 

94 def to_dict(self) -> dict[str, Any]: 

95 """Convert to dictionary.""" 

96 return { 

97 "total_items": self.total_items, 

98 "by_category": self.by_category, 

99 "retried_count": self.retried_count, 

100 "archived_count": self.archived_count, 

101 "deleted_count": self.deleted_count, 

102 "permanent_failures": self.permanent_failures, 

103 } 

104 

105 

106# ── Maintenance types ──────────────────────────────────────────────── 

107 

108 

109class MaintenanceTaskType(StrEnum): 

110 """Types of maintenance tasks.""" 

111 

112 INDEX_OPTIMIZATION = "index_optimization" 

113 CACHE_CLEANUP = "cache_cleanup" 

114 DOCUMENT_CLEANUP = "document_cleanup" 

115 METRICS_ROLLUP = "metrics_rollup" 

116 HEALTH_CHECK = "health_check" 

117 

118 

119class MaintenanceStatus(StrEnum): 

120 """Maintenance task status.""" 

121 

122 PENDING = "pending" 

123 RUNNING = "running" 

124 COMPLETED = "completed" 

125 FAILED = "failed" 

126 SKIPPED = "skipped" 

127 

128 

129@dataclass 

130class MaintenanceTask: 

131 """Configuration for a maintenance task.""" 

132 

133 name: str 

134 task_type: MaintenanceTaskType 

135 handler: Callable[[], Any] 

136 schedule_cron: str | None = None 

137 interval_seconds: int | None = None 

138 enabled: bool = True 

139 timeout: float = 300.0 

140 

141 last_run: datetime | None = None 

142 last_status: MaintenanceStatus = MaintenanceStatus.PENDING 

143 last_error: str | None = None 

144 run_count: int = 0 

145 

146 def should_run(self) -> bool: 

147 """Check if task should run based on schedule.""" 

148 if not self.enabled: 

149 return False 

150 if self.last_run is None: 

151 return True 

152 if self.interval_seconds: 

153 elapsed = (datetime.now(UTC) - self.last_run).total_seconds() 

154 return elapsed >= self.interval_seconds 

155 return False 

156 

157 def to_dict(self) -> dict[str, Any]: 

158 """Convert to dictionary for serialization.""" 

159 return { 

160 "name": self.name, 

161 "task_type": self.task_type.value, 

162 "schedule_cron": self.schedule_cron, 

163 "interval_seconds": self.interval_seconds, 

164 "enabled": self.enabled, 

165 "last_run": self.last_run.isoformat() if self.last_run else None, 

166 "last_status": self.last_status.value, 

167 "last_error": self.last_error, 

168 "run_count": self.run_count, 

169 } 

170 

171 

172@dataclass 

173class MaintenanceResult: 

174 """Result of a maintenance task execution.""" 

175 

176 task_name: str 

177 task_type: MaintenanceTaskType 

178 status: MaintenanceStatus 

179 started_at: datetime 

180 completed_at: datetime 

181 duration_seconds: float 

182 items_processed: int = 0 

183 items_deleted: int = 0 

184 error: str | None = None 

185 metadata: dict[str, Any] = field(default_factory=dict) 

186 

187 @classmethod 

188 def success( 

189 cls, 

190 task_name: str, 

191 task_type: MaintenanceTaskType, 

192 started_at: datetime, 

193 items_processed: int = 0, 

194 items_deleted: int = 0, 

195 metadata: dict[str, Any] | None = None, 

196 ) -> MaintenanceResult: 

197 """Create successful result.""" 

198 now = datetime.now(UTC) 

199 return cls( 

200 task_name=task_name, 

201 task_type=task_type, 

202 status=MaintenanceStatus.COMPLETED, 

203 started_at=started_at, 

204 completed_at=now, 

205 duration_seconds=(now - started_at).total_seconds(), 

206 items_processed=items_processed, 

207 items_deleted=items_deleted, 

208 metadata=metadata or {}, 

209 ) 

210 

211 @classmethod 

212 def failure( 

213 cls, 

214 task_name: str, 

215 task_type: MaintenanceTaskType, 

216 started_at: datetime, 

217 error: str, 

218 ) -> MaintenanceResult: 

219 """Create failed result.""" 

220 now = datetime.now(UTC) 

221 return cls( 

222 task_name=task_name, 

223 task_type=task_type, 

224 status=MaintenanceStatus.FAILED, 

225 started_at=started_at, 

226 completed_at=now, 

227 duration_seconds=(now - started_at).total_seconds(), 

228 error=error, 

229 ) 

230 

231 def to_dict(self) -> dict[str, Any]: 

232 """Convert to dictionary for serialization.""" 

233 return { 

234 "task_name": self.task_name, 

235 "task_type": self.task_type.value, 

236 "status": self.status.value, 

237 "started_at": self.started_at.isoformat(), 

238 "completed_at": self.completed_at.isoformat(), 

239 "duration_seconds": self.duration_seconds, 

240 "items_processed": self.items_processed, 

241 "items_deleted": self.items_deleted, 

242 "error": self.error, 

243 "metadata": self.metadata, 

244 } 

245 

246 

247__all__ = [ 

248 "DLQAction", 

249 "DLQItem", 

250 "DLQStats", 

251 "FailureCategory", 

252 "MaintenanceResult", 

253 "MaintenanceStatus", 

254 "MaintenanceTask", 

255 "MaintenanceTaskType", 

256]