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
« 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."""
3from __future__ import annotations
5from dataclasses import dataclass, field
6from datetime import UTC, datetime
7from enum import StrEnum
8from typing import TYPE_CHECKING, Any
10if TYPE_CHECKING:
11 from collections.abc import Callable
13 from lexigram.contracts import JobProtocol
16# ── DLQ types ────────────────────────────────────────────────────────
19class FailureCategory(StrEnum):
20 """Categories of failures."""
22 TRANSIENT = "transient"
23 PERMANENT = "permanent"
24 THROTTLED = "throttled"
25 INVALID_INPUT = "invalid_input"
26 UNKNOWN = "unknown"
29class DLQAction(StrEnum):
30 """Actions for DLQ items."""
32 RETRY = "retry"
33 ARCHIVE = "archive"
34 DELETE = "delete"
35 NOTIFY = "notify"
38@dataclass
39class DLQItem:
40 """Item in the dead letter queue."""
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)
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)
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)
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 }
83@dataclass
84class DLQStats:
85 """Statistics for DLQ."""
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
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 }
106# ── Maintenance types ────────────────────────────────────────────────
109class MaintenanceTaskType(StrEnum):
110 """Types of maintenance tasks."""
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"
119class MaintenanceStatus(StrEnum):
120 """Maintenance task status."""
122 PENDING = "pending"
123 RUNNING = "running"
124 COMPLETED = "completed"
125 FAILED = "failed"
126 SKIPPED = "skipped"
129@dataclass
130class MaintenanceTask:
131 """Configuration for a maintenance task."""
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
141 last_run: datetime | None = None
142 last_status: MaintenanceStatus = MaintenanceStatus.PENDING
143 last_error: str | None = None
144 run_count: int = 0
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
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 }
172@dataclass
173class MaintenanceResult:
174 """Result of a maintenance task execution."""
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)
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 )
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 )
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 }
247__all__ = [
248 "DLQAction",
249 "DLQItem",
250 "DLQStats",
251 "FailureCategory",
252 "MaintenanceResult",
253 "MaintenanceStatus",
254 "MaintenanceTask",
255 "MaintenanceTaskType",
256]