@@ -218,11 +218,11 @@ impl Chunk {
218218#[ cfg( feature = "server-logic" ) ]
219219pub mod db {
220220 use super :: * ;
221- use crate :: common:: errors:: { ChunkNotFound , DatabaseError , E , SubmissionNotFound } ;
221+ use crate :: common:: errors:: { DatabaseError , E , SubmissionNotFound } ;
222222 use crate :: db:: { Connection , True , WriterConnection } ;
223223 use axum_prometheus:: metrics:: { counter, gauge} ;
224224 use sqlx:: { QueryBuilder , Sqlite } ;
225- use sqlx:: { query, query_as} ;
225+ use sqlx:: { query, query_as, query_scalar } ;
226226
227227 impl < ' q > sqlx:: Encode < ' q , Sqlite > for super :: ChunkIndex {
228228 fn encode_by_ref (
@@ -281,25 +281,18 @@ pub mod db {
281281 chunk_id : ChunkId ,
282282 output_content : Option < Vec < u8 > > ,
283283 mut conn : impl WriterConnection ,
284- ) -> Result < ( ) , E < DatabaseError , E < SubmissionNotFound , ChunkNotFound > > > {
285- let _chunk_size: Result < ChunkSize , E < DatabaseError , E < SubmissionNotFound , ChunkNotFound > > > =
286- conn. transaction ( move |mut tx| {
287- Box :: pin ( async move {
288- let completed_work =
289- complete_chunk_raw ( chunk_id, output_content, & mut tx) . await ?;
290- crate :: common:: submission:: db:: maybe_complete_submission (
291- chunk_id. submission_id ,
292- & mut tx,
293- )
294- . await
295- . map_err ( |e| match e {
296- E :: L ( e) => E :: L ( e) ,
297- E :: R ( e) => E :: R ( E :: L ( e) ) ,
298- } ) ?;
299- Ok ( completed_work. unwrap_or_default ( ) )
300- } )
284+ ) -> Result < ( ) , E < DatabaseError , SubmissionNotFound > > {
285+ conn. transaction ( move |mut tx| {
286+ Box :: pin ( async move {
287+ complete_chunk_raw ( chunk_id, output_content, & mut tx) . await ?;
288+ crate :: common:: submission:: db:: maybe_complete_submission (
289+ chunk_id. submission_id ,
290+ & mut tx,
291+ )
292+ . await
301293 } )
302- . await ;
294+ } )
295+ . await ?;
303296
304297 counter ! ( crate :: prometheus:: CHUNKS_COMPLETED_COUNTER ) . increment ( 1 ) ;
305298 Ok ( ( ) )
@@ -311,9 +304,9 @@ pub mod db {
311304 chunk_id : ChunkId ,
312305 output_content : Option < Vec < u8 > > ,
313306 mut tx : impl WriterConnection < Transaction = True > ,
314- ) -> sqlx:: Result < Option < ChunkSize > > {
307+ ) -> sqlx:: Result < ( ) > {
315308 let now = chrono:: prelude:: Utc :: now ( ) ;
316- query ! (
309+ let chunk_moved = query ! (
317310 "
318311 INSERT INTO chunks_completed
319312 (submission_id, chunk_index, output_content, completed_at)
@@ -330,29 +323,43 @@ pub mod db {
330323 chunk_id. submission_id,
331324 chunk_id. chunk_index,
332325 )
333- . fetch_one ( tx. get_inner ( ) )
334- . await ?;
335- // Defense in depth: Above query should never be called twice on the same chunk.
336- // If it _does_ happen, it means that either a consumer is attempting a chunk they didn't reserve,
337- // or we gave out the same reservation twice.
326+ . fetch_optional ( tx. get_inner ( ) )
327+ . await ?
328+ . is_some ( ) ;
329+ // Defense in depth: Above query could be called twice on the same chunk. For instance,
330+ // when the server was restarted and the reservations are forgotten, and the same chunk
331+ // was reserved again.
332+ //
333+ // In addition, cancelling or pausing a submission while a chunk is reserved also results
334+ // in the chunk not being in the `chunks` table. Which is fine, because cancelled
335+ // submissions count as failed, and for paused submissions we will retry the chunk when it
336+ // becomes unpaused.
338337 //
339- // By returning early if the chunk was not found,
340- // we ensure that even in these situations
341- // we never mess up the submission's `chunks_done` counter.
338+ // By only updating `chunks_done` when we actually moved a chunk, we ensure that we never
339+ // mess up the submission's `chunks_done` counter.
342340 //
343341 // This does mean we potentially run the same chunk twice, but that is fine because we
344342 // assume chunks to be processed idempotently.
345343 //
346344 // (Not doing that resulted in a hard-to-track-down bug in the past.
347345 // https://github.com/channable/opsqueue/issues/76
348346 // )
349- sqlx:: query_scalar!(
350- "UPDATE submissions SET chunks_done = chunks_done + 1 WHERE submissions.id = $1 RETURNING submissions.chunk_size;" ,
351- chunk_id. submission_id,
352- )
353- . fetch_one ( tx. get_inner ( ) )
354- . await
355- . map ( |opt| opt. map ( ChunkSize ) )
347+ if chunk_moved {
348+ sqlx:: query_scalar!(
349+ "UPDATE submissions SET chunks_done = chunks_done + 1 WHERE submissions.id = $1 RETURNING submissions.chunk_size;" ,
350+ chunk_id. submission_id,
351+ )
352+ . fetch_one ( tx. get_inner ( ) )
353+ . await ?;
354+ } else {
355+ tracing:: warn!(
356+ "Could not complete chunk {:?} because it was either: \
357+ completed, failed, cancelled, or paused before. Ignoring.",
358+ chunk_id
359+ ) ;
360+ }
361+
362+ Ok ( ( ) )
356363 }
357364
358365 #[ tracing:: instrument( skip( conn) ) ]
@@ -369,7 +376,7 @@ pub mod db {
369376 submission_id,
370377 chunk_index,
371378 } = chunk_id;
372- let fields = query ! (
379+ let retries = query_scalar ! (
373380 "
374381 UPDATE chunks SET retries = retries + 1
375382 WHERE submission_id = $1 AND chunk_index = $2
@@ -378,24 +385,37 @@ pub mod db {
378385 submission_id,
379386 chunk_index
380387 )
381- . fetch_one ( tx. get_inner ( ) )
388+ . fetch_optional ( tx. get_inner ( ) )
382389 . await ?;
383- tracing:: trace!( "Retries: {}" , fields. retries) ;
384- if fields. retries >= max_retries. into ( ) {
385- crate :: common:: submission:: db:: fail_submission_notx (
386- submission_id,
387- chunk_index,
388- failure,
389- & mut tx,
390- )
391- . await ?;
392-
393- Ok :: < _ , sqlx:: Error > ( true )
394- } else {
395- counter ! ( crate :: prometheus:: CHUNKS_RETRIED_COUNTER ) . increment ( 1 ) ;
396- // When retrying, the chunk re-enters ('stays') in the backlog,
397- // so we *don't* decrement the backlog gauge here.
398- Ok :: < _ , sqlx:: Error > ( false )
390+ match retries {
391+ Some ( retries) => {
392+ tracing:: trace!( "Retries: {}" , retries) ;
393+ if retries >= max_retries. into ( ) {
394+ crate :: common:: submission:: db:: fail_submission_notx (
395+ submission_id,
396+ chunk_index,
397+ failure,
398+ & mut tx,
399+ )
400+ . await ?;
401+
402+ Ok :: < _ , sqlx:: Error > ( true )
403+ } else {
404+ counter ! ( crate :: prometheus:: CHUNKS_RETRIED_COUNTER ) . increment ( 1 ) ;
405+ // When retrying, the chunk re-enters ('stays') in the backlog,
406+ // so we *don't* decrement the backlog gauge here.
407+ Ok :: < _ , sqlx:: Error > ( false )
408+ }
409+ }
410+ None => {
411+ tracing:: warn!(
412+ "Could not fail chunk {:?} because it was either: \
413+ completed, failed, cancelled, or paused before. Ignoring.",
414+ chunk_id
415+ ) ;
416+
417+ Ok :: < _ , sqlx:: Error > ( false )
418+ }
399419 }
400420 } )
401421 } )
0 commit comments