1use anyhow::Result;
2use rusqlite::Connection;
3use serde::Serialize;
4
5mod maintenance;
6mod query;
7mod sql;
8
9#[cfg(test)]
10mod tests;
11
12use maintenance::{
13 archive_surface, purge_archived_extraction_tasks, purge_archived_replay_ranges,
14 purge_simple_surface, requeue_due_extraction_tasks, requeue_due_jobs,
15 retry_due_extraction_replay_ranges, ArchiveSurface,
16};
17use query::{query_surface_stats, SurfaceQuery};
18use sql::{
19 column_exists, count_archived_rows, count_purgeable_extraction_tasks, cutoff_epoch,
20 failure_columns_available, table_exists,
21};
22
23pub const FAILURE_RETENTION_DAYS: i64 = 14;
24pub const ARCHIVED_FAILURE_PURGE_DAYS: i64 = 90;
25pub const MAX_FAILURE_AUTO_RETRIES: i64 = 3;
26
27const SECONDS_PER_DAY: i64 = 86_400;
28const FAILURE_RETRY_BASE_SECS: i64 = 300;
29
30#[derive(Debug, Clone, Copy, PartialEq, Eq)]
31pub enum FailureClass {
32 Transient,
33 Permanent,
34}
35
36impl FailureClass {
37 pub fn as_str(self) -> &'static str {
38 match self {
39 FailureClass::Transient => "transient",
40 FailureClass::Permanent => "permanent",
41 }
42 }
43}
44
45#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)]
46pub struct FailureLifecycleStats {
47 pub pending_observation: FailureSurfaceStats,
48 pub extraction_task: FailureSurfaceStats,
49 pub extraction_replay_range: FailureSurfaceStats,
50 pub job: FailureSurfaceStats,
51}
52
53#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)]
54pub struct FailureSurfaceStats {
55 pub actionable_total: i64,
56 pub actionable_7d: i64,
57 pub transient: i64,
58 pub permanent: i64,
59 pub exhausted: i64,
60 pub archived: i64,
61 pub historical_archived: i64,
62 pub historical_purged: i64,
63 pub oldest_actionable_epoch: Option<i64>,
64}
65
66#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)]
67pub struct FailureLifecycleMaintenance {
68 pub retried_extraction_replay_ranges: usize,
69 pub retried_extraction_tasks: usize,
70 pub retried_jobs: usize,
71 pub coalesced_jobs: usize,
72 pub archived_pending_observations: usize,
73 pub archived_extraction_tasks: usize,
74 pub archived_extraction_replay_ranges: usize,
75 pub archived_jobs: usize,
76}
77
78#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)]
79pub struct ArchivedFailurePurgePlan {
80 pub pending_observations: usize,
81 pub extraction_replay_ranges: usize,
82 pub extraction_tasks: usize,
83 pub jobs: usize,
84}
85
86pub fn classify_failure(error: &str) -> FailureClass {
87 let lower = error.to_ascii_lowercase();
88 if lower.contains("database schema is locked")
89 || lower.contains("database is locked")
90 || lower.contains("sqlite_busy")
91 || lower.contains("sqlite busy")
92 {
93 return FailureClass::Transient;
94 }
95
96 if [
97 "schema",
98 "vocab",
99 "malformed",
100 "invalid payload",
101 "invalid json",
102 "unsupported version",
103 "missing evidence",
104 "not implemented",
105 "retired",
106 ]
107 .iter()
108 .any(|needle| lower.contains(needle))
109 {
110 return FailureClass::Permanent;
111 }
112
113 FailureClass::Transient
114}
115
116pub fn query_failure_lifecycle_stats(
117 conn: &Connection,
118 now_epoch: i64,
119) -> Result<FailureLifecycleStats> {
120 let job_failed_predicate = job_failed_predicate(conn)?;
121 Ok(FailureLifecycleStats {
122 pending_observation: query_surface_stats(
123 conn,
124 SurfaceQuery {
125 surface: "pending_observation",
126 table: "pending_observations",
127 failed_predicate: "status = 'failed'".into(),
128 attempt_column: "attempt_count",
129 created_column: "created_at_epoch",
130 updated_column: "updated_at_epoch",
131 },
132 now_epoch,
133 )?,
134 extraction_task: query_surface_stats(
135 conn,
136 SurfaceQuery {
137 surface: "extraction_task",
138 table: "extraction_tasks",
139 failed_predicate: "status = 'failed'".into(),
140 attempt_column: "attempts",
141 created_column: "created_at_epoch",
142 updated_column: "updated_at_epoch",
143 },
144 now_epoch,
145 )?,
146 extraction_replay_range: query_surface_stats(
147 conn,
148 SurfaceQuery {
149 surface: "extraction_replay_range",
150 table: "extraction_replay_ranges",
151 failed_predicate: "status IN ('pending', 'failed', 'quarantined')".into(),
152 attempt_column: "attempts",
153 created_column: "created_at_epoch",
154 updated_column: "updated_at_epoch",
155 },
156 now_epoch,
157 )?,
158 job: query_surface_stats(
159 conn,
160 SurfaceQuery {
161 surface: "job",
162 table: "jobs",
163 failed_predicate: job_failed_predicate,
164 attempt_column: "attempt_count",
165 created_column: "created_at_epoch",
166 updated_column: "updated_at_epoch",
167 },
168 now_epoch,
169 )?,
170 })
171}
172
173fn job_failed_predicate(conn: &Connection) -> Result<std::borrow::Cow<'static, str>> {
174 if table_exists(conn, "jobs")?
175 && column_exists(conn, "jobs", "job_type")?
176 && column_exists(conn, "jobs", "failure_class")?
177 && column_exists(conn, "jobs", "last_error")?
178 {
179 Ok("state = 'failed'
180 AND NOT (
181 job_type = 'summary'
182 AND failure_class = 'permanent'
183 AND last_error = 'legacy summary job rejected during GH684 summary retirement upgrade; SessionRollup owns session summary output'
184 )"
185 .into())
186 } else {
187 Ok("state = 'failed'".into())
188 }
189}
190
191pub fn maintain_failure_lifecycle(conn: &Connection) -> Result<FailureLifecycleMaintenance> {
192 if !failure_columns_available(conn)? {
193 return Ok(FailureLifecycleMaintenance::default());
194 }
195 let now = chrono::Utc::now().timestamp();
196 let job_recovery = requeue_due_jobs(conn, now)?;
197 let mut result = FailureLifecycleMaintenance {
198 retried_extraction_replay_ranges: retry_due_extraction_replay_ranges(conn, now)?,
199 retried_extraction_tasks: requeue_due_extraction_tasks(conn, now)?,
200 retried_jobs: job_recovery.requeued,
201 coalesced_jobs: job_recovery.coalesced,
202 ..FailureLifecycleMaintenance::default()
203 };
204 let archived = archive_eligible_failures(conn, now, FAILURE_RETENTION_DAYS)?;
205 result.archived_pending_observations = archived.pending_observations;
206 result.archived_extraction_tasks = archived.extraction_tasks;
207 result.archived_extraction_replay_ranges = archived.extraction_replay_ranges;
208 result.archived_jobs = archived.jobs;
209
210 if result.retried_extraction_replay_ranges > 0
211 || result.retried_extraction_tasks > 0
212 || result.retried_jobs > 0
213 || result.coalesced_jobs > 0
214 || result.archived_pending_observations > 0
215 || result.archived_extraction_tasks > 0
216 || result.archived_extraction_replay_ranges > 0
217 || result.archived_jobs > 0
218 {
219 crate::log::info(
220 "failure_lifecycle",
221 &format!(
222 "maintenance retried replay_ranges={} extraction_tasks={} jobs={} coalesced_jobs={} archived pending_observations={} extraction_tasks={} replay_ranges={} jobs={}",
223 result.retried_extraction_replay_ranges,
224 result.retried_extraction_tasks,
225 result.retried_jobs,
226 result.coalesced_jobs,
227 result.archived_pending_observations,
228 result.archived_extraction_tasks,
229 result.archived_extraction_replay_ranges,
230 result.archived_jobs
231 ),
232 );
233 }
234
235 Ok(result)
236}
237
238pub fn archive_eligible_failures(
239 conn: &Connection,
240 now_epoch: i64,
241 retention_days: i64,
242) -> Result<ArchivedFailurePurgePlan> {
243 if !failure_columns_available(conn)? {
244 return Ok(ArchivedFailurePurgePlan::default());
245 }
246 let cutoff = cutoff_epoch(now_epoch, retention_days);
247 let tx = conn.unchecked_transaction()?;
248 let pending = archive_surface(
249 &tx,
250 ArchiveSurface {
251 surface: "pending_observation",
252 table: "pending_observations",
253 failed_predicate: "status = 'failed'",
254 eligible_extra: "COALESCE(failure_class, 'transient') = 'permanent'",
255 },
256 cutoff,
257 now_epoch,
258 )?;
259 let extraction_tasks = archive_surface(
260 &tx,
261 ArchiveSurface {
262 surface: "extraction_task",
263 table: "extraction_tasks",
264 failed_predicate: "status = 'failed'",
265 eligible_extra: "(failure_class = 'permanent' OR attempts >= 3)",
266 },
267 cutoff,
268 now_epoch,
269 )?;
270 let replay_ranges = archive_surface(
271 &tx,
272 ArchiveSurface {
273 surface: "extraction_replay_range",
274 table: "extraction_replay_ranges",
275 failed_predicate: "status IN ('pending', 'failed', 'quarantined')",
276 eligible_extra: "(failure_class = 'permanent' OR attempts >= 3)",
277 },
278 cutoff,
279 now_epoch,
280 )?;
281 let jobs = archive_surface(
282 &tx,
283 ArchiveSurface {
284 surface: "job",
285 table: "jobs",
286 failed_predicate: "state = 'failed'",
287 eligible_extra: "(failure_class = 'permanent' OR attempt_count >= 3)",
288 },
289 cutoff,
290 now_epoch,
291 )?;
292 tx.commit()?;
293 Ok(ArchivedFailurePurgePlan {
294 pending_observations: pending,
295 extraction_replay_ranges: replay_ranges,
296 extraction_tasks,
297 jobs,
298 })
299}
300
301pub fn count_archived_failures_to_purge_at(
302 conn: &Connection,
303 now_epoch: i64,
304 horizon_days: i64,
305) -> Result<ArchivedFailurePurgePlan> {
306 if !failure_columns_available(conn)? {
307 return Ok(ArchivedFailurePurgePlan::default());
308 }
309 let cutoff = cutoff_epoch(now_epoch, horizon_days);
310 Ok(ArchivedFailurePurgePlan {
311 pending_observations: count_archived_rows(
312 conn,
313 "pending_observations",
314 "status = 'failed'",
315 cutoff,
316 )?,
317 extraction_replay_ranges: count_archived_rows(
318 conn,
319 "extraction_replay_ranges",
320 "status IN ('pending', 'failed', 'quarantined')",
321 cutoff,
322 )?,
323 extraction_tasks: count_purgeable_extraction_tasks(conn, cutoff)?,
324 jobs: count_archived_rows(conn, "jobs", "state = 'failed'", cutoff)?,
325 })
326}
327
328pub fn purge_archived_failures_at(
329 conn: &Connection,
330 now_epoch: i64,
331 horizon_days: i64,
332) -> Result<ArchivedFailurePurgePlan> {
333 if !failure_columns_available(conn)? {
334 return Ok(ArchivedFailurePurgePlan::default());
335 }
336 if !conn.is_autocommit() {
337 return purge_archived_failures_in_transaction(conn, now_epoch, horizon_days);
338 }
339 let tx = rusqlite::Transaction::new_unchecked(conn, rusqlite::TransactionBehavior::Immediate)?;
340 let purged = purge_archived_failures_in_transaction(&tx, now_epoch, horizon_days)?;
341 tx.commit()?;
342 Ok(purged)
343}
344
345fn purge_archived_failures_in_transaction(
346 conn: &Connection,
347 now_epoch: i64,
348 horizon_days: i64,
349) -> Result<ArchivedFailurePurgePlan> {
350 let cutoff = cutoff_epoch(now_epoch, horizon_days);
351 let pending_observations = purge_simple_surface(
352 conn,
353 "pending_observation",
354 "pending_observations",
355 "status = 'failed'",
356 cutoff,
357 now_epoch,
358 )?;
359 let extraction_replay_ranges = purge_archived_replay_ranges(conn, cutoff, now_epoch)?;
360 let extraction_tasks = purge_archived_extraction_tasks(conn, cutoff, now_epoch)?;
361 let jobs = purge_simple_surface(conn, "job", "jobs", "state = 'failed'", cutoff, now_epoch)?;
362 Ok(ArchivedFailurePurgePlan {
363 pending_observations,
364 extraction_replay_ranges,
365 extraction_tasks,
366 jobs,
367 })
368}