Skip to main content

remem/db/
failure_lifecycle.rs

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}