Skip to main content

remem/db/extraction/
lifecycle.rs

1use anyhow::Result;
2use rusqlite::{params, Connection, OptionalExtension};
3
4use crate::db::extraction_replay::{mark_replay_range_failed, mark_replay_range_replayed_if_done};
5
6use super::exhaust::exhaust_extraction_task;
7use super::loaders::{ensure_task_updated, load_claimed_extraction_task};
8use super::{ExtractionTask, EXTRACTION_TASK_MAX_ATTEMPTS};
9
10pub fn claim_next_extraction_task(
11    conn: &mut Connection,
12    lease_owner: &str,
13    lease_secs: i64,
14) -> Result<Option<ExtractionTask>> {
15    let now = chrono::Utc::now().timestamp();
16    let tx = conn.transaction()?;
17    let candidate: Option<i64> = tx
18        .query_row(
19            "SELECT id FROM extraction_tasks
20             WHERE status = 'pending'
21               AND (next_retry_epoch IS NULL OR next_retry_epoch <= ?1)
22             ORDER BY priority ASC, created_at_epoch ASC, id ASC
23             LIMIT 1",
24            params![now],
25            |row| row.get(0),
26        )
27        .optional()?;
28
29    let Some(task_id) = candidate else {
30        tx.commit()?;
31        return Ok(None);
32    };
33
34    let task = claim_extraction_task_by_id_in_transaction(&tx, task_id, lease_owner, lease_secs)?;
35    tx.commit()?;
36    Ok(task)
37}
38
39#[cfg(test)]
40pub(crate) fn claim_extraction_task_by_id(
41    conn: &mut Connection,
42    task_id: i64,
43    lease_owner: &str,
44    lease_secs: i64,
45) -> Result<Option<ExtractionTask>> {
46    let tx = conn.transaction()?;
47    let task = claim_extraction_task_by_id_in_transaction(&tx, task_id, lease_owner, lease_secs)?;
48    tx.commit()?;
49    Ok(task)
50}
51
52pub(crate) fn claim_extraction_task_by_id_in_transaction(
53    conn: &Connection,
54    task_id: i64,
55    lease_owner: &str,
56    lease_secs: i64,
57) -> Result<Option<ExtractionTask>> {
58    let now = chrono::Utc::now().timestamp();
59    let lease_expires = now + lease_secs.max(1);
60    let updated = conn.execute(
61        "UPDATE extraction_tasks
62         SET status = 'processing',
63             lease_owner = ?1,
64             lease_expires_epoch = ?2,
65             updated_at_epoch = ?3
66         WHERE id = ?4
67           AND status = 'pending'
68           AND (next_retry_epoch IS NULL OR next_retry_epoch <= ?3)",
69        params![lease_owner, lease_expires, now, task_id],
70    )?;
71    if updated == 0 {
72        return Ok(None);
73    }
74
75    Ok(Some(load_claimed_extraction_task(conn, task_id)?))
76}
77
78pub fn release_expired_extraction_task_leases(conn: &Connection) -> Result<usize> {
79    let now = chrono::Utc::now().timestamp();
80    let tx = conn.unchecked_transaction()?;
81    let expired = {
82        let mut stmt = tx.prepare(
83            "SELECT id, lease_owner
84             FROM extraction_tasks
85             WHERE status = 'processing'
86               AND lease_expires_epoch IS NOT NULL
87               AND lease_expires_epoch < ?1
88             ORDER BY id ASC",
89        )?;
90        let rows = stmt.query_map(params![now], |row| {
91            Ok((row.get::<_, i64>(0)?, row.get::<_, Option<String>>(1)?))
92        })?;
93        crate::db::query::collect_rows(rows)?
94    };
95
96    for (task_id, lease_owner) in &expired {
97        if let Some(exact_owner) = lease_owner
98            .as_deref()
99            .filter(|owner| crate::db::is_exact_replay_worker_owner(owner))
100        {
101            archive_claimed_exact_replay_task_in_transaction(
102                &tx,
103                *task_id,
104                exact_owner,
105                "exact replay worker lease expired; rerun the locked exact recovery command",
106                now,
107            )?;
108        } else {
109            let updated = tx.execute(
110                "UPDATE extraction_tasks
111                 SET status = 'pending',
112                     lease_owner = NULL,
113                     lease_expires_epoch = NULL,
114                     updated_at_epoch = ?1
115                 WHERE id = ?2
116                   AND status = 'processing'
117                   AND ((?3 IS NULL AND lease_owner IS NULL) OR lease_owner = ?3)",
118                params![now, task_id, lease_owner],
119            )?;
120            ensure_task_updated(updated, *task_id)?;
121        }
122    }
123    tx.commit()?;
124    Ok(expired.len())
125}
126
127pub(crate) fn archive_claimed_exact_replay_task(
128    conn: &Connection,
129    task_id: i64,
130    lease_owner: &str,
131    error: &str,
132) -> Result<()> {
133    let tx = conn.unchecked_transaction()?;
134    archive_claimed_exact_replay_task_in_transaction(
135        &tx,
136        task_id,
137        lease_owner,
138        error,
139        chrono::Utc::now().timestamp(),
140    )?;
141    tx.commit()?;
142    Ok(())
143}
144
145fn archive_claimed_exact_replay_task_in_transaction(
146    conn: &Connection,
147    task_id: i64,
148    lease_owner: &str,
149    error: &str,
150    now: i64,
151) -> Result<()> {
152    anyhow::ensure!(
153        crate::db::is_exact_replay_worker_owner(lease_owner),
154        "exact replay archive requires an exact replay worker owner"
155    );
156    let replay_range_id: i64 = conn.query_row(
157        "SELECT replay_range_id
158         FROM extraction_tasks
159         WHERE id = ?1 AND status = 'processing' AND lease_owner = ?2",
160        params![task_id, lease_owner],
161        |row| row.get(0),
162    )?;
163    let updated = conn.execute(
164        "UPDATE extraction_tasks
165         SET status = 'failed',
166             attempts = attempts + 1,
167             next_retry_epoch = NULL,
168             lease_owner = NULL,
169             lease_expires_epoch = NULL,
170             last_error = ?1,
171             failure_class = ?2,
172             failed_at_epoch = COALESCE(failed_at_epoch, ?3),
173             archived_at_epoch = ?3,
174             updated_at_epoch = ?3
175         WHERE id = ?4 AND status = 'processing' AND lease_owner = ?5",
176        params![
177            crate::db::truncate_str(error, 2000),
178            crate::db::classify_failure(error).as_str(),
179            now,
180            task_id,
181            lease_owner
182        ],
183    )?;
184    ensure_task_updated(updated, task_id)?;
185    crate::db::extraction_replay::archive_exact_replay_range_after_task_failure(
186        conn,
187        replay_range_id,
188        task_id,
189        error,
190        now,
191    )
192}
193
194pub fn mark_extraction_task_done(
195    conn: &Connection,
196    task_id: i64,
197    lease_owner: &str,
198    completed_high_watermark_event_id: Option<i64>,
199) -> Result<()> {
200    let now = chrono::Utc::now().timestamp();
201    let updated = conn.execute(
202        "UPDATE extraction_tasks
203         SET status = CASE
204                 WHEN ?4 IS NOT NULL
205                  AND high_watermark_event_id IS NOT NULL
206                  AND high_watermark_event_id > ?4 THEN 'pending'
207                 ELSE 'done'
208             END,
209             cursor_event_id = ?4,
210             lease_owner = NULL,
211             lease_expires_epoch = NULL,
212             next_retry_epoch = NULL,
213             last_error = NULL,
214             failure_class = NULL,
215             failed_at_epoch = NULL,
216             archived_at_epoch = NULL,
217             updated_at_epoch = ?1
218         WHERE id = ?2 AND lease_owner = ?3 AND status = 'processing'",
219        params![now, task_id, lease_owner, completed_high_watermark_event_id],
220    )?;
221    ensure_task_updated(updated, task_id)?;
222    mark_replay_range_replayed_if_done(conn, task_id, now)
223}
224
225pub fn mark_extraction_task_failed(
226    conn: &Connection,
227    task_id: i64,
228    lease_owner: &str,
229    err: &str,
230) -> Result<()> {
231    let now = chrono::Utc::now().timestamp();
232    let updated = conn.execute(
233        "UPDATE extraction_tasks
234         SET status = 'failed',
235             attempts = attempts + 1,
236             lease_owner = NULL,
237             lease_expires_epoch = NULL,
238             next_retry_epoch = NULL,
239             last_error = ?1,
240             failure_class = ?2,
241             failed_at_epoch = COALESCE(failed_at_epoch, ?3),
242             archived_at_epoch = NULL,
243             updated_at_epoch = ?3
244         WHERE id = ?4 AND lease_owner = ?5 AND status = 'processing'",
245        params![
246            crate::db::truncate_str(err, 2000),
247            crate::db::classify_failure(err).as_str(),
248            now,
249            task_id,
250            lease_owner
251        ],
252    )?;
253    ensure_task_updated(updated, task_id)?;
254    mark_replay_range_failed(conn, task_id, now, err)
255}
256
257pub fn defer_extraction_task(
258    conn: &Connection,
259    task_id: i64,
260    lease_owner: &str,
261    reason: &str,
262    backoff_secs: i64,
263) -> Result<()> {
264    let task = load_claimed_extraction_task(conn, task_id)?;
265    defer_claimed_extraction_task(conn, &task, lease_owner, reason, backoff_secs)
266}
267
268pub fn defer_claimed_extraction_task(
269    conn: &Connection,
270    task: &ExtractionTask,
271    lease_owner: &str,
272    reason: &str,
273    backoff_secs: i64,
274) -> Result<()> {
275    let now = chrono::Utc::now().timestamp();
276    let next_attempt = task.attempts + 1;
277    if next_attempt >= EXTRACTION_TASK_MAX_ATTEMPTS {
278        return exhaust_extraction_task(conn, task, lease_owner, next_attempt, reason, now);
279    }
280
281    let updated = conn.execute(
282        "UPDATE extraction_tasks
283         SET status = 'pending',
284             attempts = ?1,
285             lease_owner = NULL,
286             lease_expires_epoch = NULL,
287             next_retry_epoch = ?2,
288             last_error = ?3,
289             failure_class = NULL,
290             failed_at_epoch = NULL,
291             archived_at_epoch = NULL,
292             updated_at_epoch = ?4
293         WHERE id = ?5 AND lease_owner = ?6 AND status = 'processing'",
294        params![
295            next_attempt,
296            now + backoff_secs.max(1),
297            crate::db::truncate_str(reason, 2000),
298            now,
299            task.id,
300            lease_owner
301        ],
302    )?;
303    ensure_task_updated(updated, task.id)
304}
305
306pub fn wait_extraction_task(
307    conn: &Connection,
308    task_id: i64,
309    lease_owner: &str,
310    reason: &str,
311    backoff_secs: i64,
312) -> Result<()> {
313    let now = chrono::Utc::now().timestamp();
314    let updated = conn.execute(
315        "UPDATE extraction_tasks
316         SET status = 'pending',
317             lease_owner = NULL,
318             lease_expires_epoch = NULL,
319             next_retry_epoch = ?1,
320             last_error = ?2,
321             failure_class = NULL,
322             failed_at_epoch = NULL,
323             archived_at_epoch = NULL,
324             updated_at_epoch = ?3
325         WHERE id = ?4 AND lease_owner = ?5 AND status = 'processing'",
326        params![
327            now + backoff_secs.max(1),
328            crate::db::truncate_str(reason, 2000),
329            now,
330            task_id,
331            lease_owner
332        ],
333    )?;
334    ensure_task_updated(updated, task_id)
335}
336
337pub fn mark_extraction_task_failed_or_retry(
338    conn: &Connection,
339    task_id: i64,
340    lease_owner: &str,
341    err: &str,
342    backoff_secs: i64,
343) -> Result<()> {
344    let task = load_claimed_extraction_task(conn, task_id)?;
345    mark_claimed_extraction_task_failed_or_retry(conn, &task, lease_owner, err, backoff_secs)
346}
347
348pub fn mark_claimed_extraction_task_failed_or_retry(
349    conn: &Connection,
350    task: &ExtractionTask,
351    lease_owner: &str,
352    err: &str,
353    backoff_secs: i64,
354) -> Result<()> {
355    let now = chrono::Utc::now().timestamp();
356    let next_attempt = task.attempts + 1;
357    if crate::db::classify_failure(err) == crate::db::FailureClass::Permanent
358        || next_attempt >= EXTRACTION_TASK_MAX_ATTEMPTS
359    {
360        return exhaust_extraction_task(conn, task, lease_owner, next_attempt, err, now);
361    }
362
363    let updated = conn.execute(
364        "UPDATE extraction_tasks
365         SET status = 'pending',
366             attempts = ?1,
367             next_retry_epoch = ?2,
368             lease_owner = NULL,
369             lease_expires_epoch = NULL,
370             last_error = ?3,
371             failure_class = NULL,
372             failed_at_epoch = NULL,
373             archived_at_epoch = NULL,
374             updated_at_epoch = ?4
375         WHERE id = ?5 AND lease_owner = ?6 AND status = 'processing'",
376        params![
377            next_attempt,
378            now + backoff_secs.max(1),
379            crate::db::truncate_str(err, 2000),
380            now,
381            task.id,
382            lease_owner
383        ],
384    )?;
385    ensure_task_updated(updated, task.id)
386}