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}