1use crate::db::{
7 ReadSetRevisionError, ReadSetRevisionProof, ReadSetStoreIdentity, ReadSetStoreRevision,
8 codec::{ByteDecodeError, ByteReader},
9};
10use candid::CandidType;
11use serde::Deserialize;
12use std::{error::Error as StdError, fmt};
13
14pub const MAX_RESUMABLE_JOB_STATE_BYTES: usize = 256 * 1024;
16pub const MAX_RESUMABLE_JOB_RECEIPT_BYTES: usize = 64 * 1024;
18pub const MAX_RESUMABLE_JOB_IDEMPOTENCY_KEY_BYTES: usize = 256;
20pub const MAX_RESUMABLE_JOB_CONTINUATION_BYTES: usize = 16 * 1024;
22
23#[derive(CandidType, Clone, Copy, Debug, Deserialize, Eq, Hash, Ord, PartialEq, PartialOrd)]
25pub struct ResumableJobId([u8; 32]);
26
27impl ResumableJobId {
28 pub fn try_from_bytes(bytes: [u8; 32]) -> Result<Self, ResumableJobError> {
30 if bytes == [0; 32] {
31 return Err(ResumableJobError::InvalidJobId);
32 }
33 Ok(Self(bytes))
34 }
35
36 #[must_use]
38 pub const fn to_bytes(self) -> [u8; 32] {
39 self.0
40 }
41}
42
43#[derive(CandidType, Clone, Debug, Deserialize, Eq, PartialEq)]
45pub struct ResumableJobIdempotencyKey(String);
46
47impl ResumableJobIdempotencyKey {
48 pub fn new(value: impl Into<String>) -> Result<Self, ResumableJobError> {
50 let value = value.into();
51 if value.is_empty() || value.len() > MAX_RESUMABLE_JOB_IDEMPOTENCY_KEY_BYTES {
52 return Err(ResumableJobError::InvalidIdempotencyKey);
53 }
54 Ok(Self(value))
55 }
56
57 #[must_use]
59 pub const fn as_str(&self) -> &str {
60 self.0.as_str()
61 }
62
63 pub(in crate::db) const fn validate(&self) -> Result<(), ResumableJobError> {
64 if self.0.is_empty() || self.0.len() > MAX_RESUMABLE_JOB_IDEMPOTENCY_KEY_BYTES {
65 return Err(ResumableJobError::InvalidIdempotencyKey);
66 }
67 Ok(())
68 }
69}
70
71#[derive(CandidType, Clone, Copy, Debug, Deserialize, Eq, PartialEq)]
73pub enum ResumableJobStatus {
74 Active,
76 Completed,
78 Invalidated,
80}
81
82#[derive(CandidType, Clone, Debug, Deserialize, Eq, PartialEq)]
84pub struct ResumableJobState {
85 pub job_id: ResumableJobId,
87 pub sequence: u64,
89 pub status: ResumableJobStatus,
91 pub proof: ReadSetRevisionProof,
93 pub continuation: Option<String>,
95 pub application_state: Vec<u8>,
97}
98
99#[derive(CandidType, Clone, Debug, Deserialize, Eq, PartialEq)]
101pub struct ResumableJobAdvanceRequest {
102 pub job_id: ResumableJobId,
104 pub expected_sequence: u64,
106 pub idempotency_key: ResumableJobIdempotencyKey,
108}
109
110impl ResumableJobAdvanceRequest {
111 #[must_use]
113 pub const fn new(
114 job_id: ResumableJobId,
115 expected_sequence: u64,
116 idempotency_key: ResumableJobIdempotencyKey,
117 ) -> Self {
118 Self {
119 job_id,
120 expected_sequence,
121 idempotency_key,
122 }
123 }
124}
125
126#[derive(CandidType, Clone, Debug, Deserialize, Eq, PartialEq)]
128pub struct ResumableJobAdvance {
129 pub continuation: Option<String>,
131 pub application_state: Vec<u8>,
133 pub application_receipt: Vec<u8>,
135}
136
137impl ResumableJobAdvance {
138 pub fn new(
140 continuation: Option<String>,
141 application_state: Vec<u8>,
142 application_receipt: Vec<u8>,
143 ) -> Result<Self, ResumableJobError> {
144 let advance = Self {
145 continuation,
146 application_state,
147 application_receipt,
148 };
149 advance.validate()?;
150 Ok(advance)
151 }
152
153 pub(in crate::db) fn validate(&self) -> Result<(), ResumableJobError> {
154 validate_continuation(self.continuation.as_deref())?;
155 if self.application_state.len() > MAX_RESUMABLE_JOB_STATE_BYTES
156 || self.application_receipt.len() > MAX_RESUMABLE_JOB_RECEIPT_BYTES
157 {
158 return Err(ResumableJobError::PayloadTooLarge);
159 }
160 Ok(())
161 }
162}
163
164#[derive(CandidType, Clone, Copy, Debug, Deserialize, Eq, PartialEq)]
166pub enum ResumableJobAdvanceStatus {
167 Advanced,
169 Invalidated,
171}
172
173#[derive(CandidType, Clone, Debug, Deserialize, Eq, PartialEq)]
175pub struct ResumableJobAdvanceReceipt {
176 pub request_sequence: u64,
178 pub committed_sequence: u64,
180 pub status: ResumableJobAdvanceStatus,
182 pub continuation: Option<String>,
184 pub application_receipt: Vec<u8>,
186 idempotency_key: ResumableJobIdempotencyKey,
187}
188
189impl ResumableJobAdvanceReceipt {
190 #[must_use]
192 pub const fn idempotency_key(&self) -> &ResumableJobIdempotencyKey {
193 &self.idempotency_key
194 }
195}
196
197#[derive(CandidType, Clone, Debug, Deserialize, Eq, PartialEq)]
199pub enum ResumableJobError {
200 InvalidJobId,
202 InvalidIdempotencyKey,
204 PayloadTooLarge,
206 AlreadyExists,
208 NotFound,
210 StaleSequence { expected: u64, actual: u64 },
212 Invalidated,
214 Completed,
216 NotTerminal,
218 SourceProof(ReadSetRevisionError),
220 CapacityExceeded,
222 CorruptProgressStore,
224 IncompatibleProgressFormat,
226 Internal,
228 ExecutionBudgetExceeded {
230 resource: u64,
231 limit: u64,
232 observed: u64,
233 scope: u64,
234 lane: u64,
235 normalized_shape_fingerprint_prefix: u64,
236 },
237}
238
239impl fmt::Display for ResumableJobError {
240 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
241 formatter.write_str("resumable job operation failed")
242 }
243}
244
245impl From<ByteDecodeError> for ResumableJobError {
246 fn from(_: ByteDecodeError) -> Self {
247 Self::CorruptProgressStore
248 }
249}
250
251impl From<ReadSetRevisionError> for ResumableJobError {
252 fn from(error: ReadSetRevisionError) -> Self {
253 Self::SourceProof(error)
254 }
255}
256
257impl StdError for ResumableJobError {}
258
259#[derive(Debug)]
261pub enum CompareProofAndAdvanceError<E> {
262 Protocol(ResumableJobError),
264 Operation(E),
266}
267
268impl<E: fmt::Display> fmt::Display for CompareProofAndAdvanceError<E> {
269 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
270 match self {
271 Self::Protocol(error) => error.fmt(formatter),
272 Self::Operation(error) => error.fmt(formatter),
273 }
274 }
275}
276
277impl<E: StdError + 'static> StdError for CompareProofAndAdvanceError<E> {}
278
279impl<E> From<ResumableJobError> for CompareProofAndAdvanceError<E> {
280 fn from(error: ResumableJobError) -> Self {
281 Self::Protocol(error)
282 }
283}
284
285#[derive(Clone, Debug, Eq, PartialEq)]
286pub(in crate::db) struct ResumableJobRecord {
287 state: ResumableJobState,
288 last_receipt: Option<ResumableJobAdvanceReceipt>,
289}
290
291impl ResumableJobRecord {
292 pub(in crate::db) fn new(
293 job_id: ResumableJobId,
294 proof: ReadSetRevisionProof,
295 application_state: Vec<u8>,
296 ) -> Result<Self, ResumableJobError> {
297 let state = ResumableJobState {
298 job_id,
299 sequence: 0,
300 status: ResumableJobStatus::Active,
301 proof,
302 continuation: None,
303 application_state,
304 };
305 let record = Self {
306 state,
307 last_receipt: None,
308 };
309 record.validate()?;
310 Ok(record)
311 }
312
313 pub(in crate::db) const fn state(&self) -> &ResumableJobState {
314 &self.state
315 }
316
317 pub(in crate::db) const fn last_receipt(&self) -> Option<&ResumableJobAdvanceReceipt> {
318 self.last_receipt.as_ref()
319 }
320
321 pub(in crate::db) fn apply_advance(
322 &self,
323 request: &ResumableJobAdvanceRequest,
324 advance: ResumableJobAdvance,
325 ) -> Result<(Self, ResumableJobAdvanceReceipt), ResumableJobError> {
326 advance.validate()?;
327 let committed_sequence = self
328 .state
329 .sequence
330 .checked_add(1)
331 .ok_or(ResumableJobError::CapacityExceeded)?;
332 let receipt = ResumableJobAdvanceReceipt {
333 request_sequence: request.expected_sequence,
334 committed_sequence,
335 status: ResumableJobAdvanceStatus::Advanced,
336 continuation: advance.continuation.clone(),
337 application_receipt: advance.application_receipt,
338 idempotency_key: request.idempotency_key.clone(),
339 };
340 let record = Self {
341 state: ResumableJobState {
342 job_id: self.state.job_id,
343 sequence: committed_sequence,
344 status: if advance.continuation.is_some() {
345 ResumableJobStatus::Active
346 } else {
347 ResumableJobStatus::Completed
348 },
349 proof: self.state.proof.clone(),
350 continuation: advance.continuation,
351 application_state: advance.application_state,
352 },
353 last_receipt: Some(receipt.clone()),
354 };
355 record.validate()?;
356 Ok((record, receipt))
357 }
358
359 pub(in crate::db) fn invalidate(
360 &self,
361 request: &ResumableJobAdvanceRequest,
362 ) -> Result<(Self, ResumableJobAdvanceReceipt), ResumableJobError> {
363 let committed_sequence = self
364 .state
365 .sequence
366 .checked_add(1)
367 .ok_or(ResumableJobError::CapacityExceeded)?;
368 let receipt = ResumableJobAdvanceReceipt {
369 request_sequence: request.expected_sequence,
370 committed_sequence,
371 status: ResumableJobAdvanceStatus::Invalidated,
372 continuation: None,
373 application_receipt: Vec::new(),
374 idempotency_key: request.idempotency_key.clone(),
375 };
376 let record = Self {
377 state: ResumableJobState {
378 sequence: committed_sequence,
379 status: ResumableJobStatus::Invalidated,
380 continuation: None,
381 ..self.state.clone()
382 },
383 last_receipt: Some(receipt.clone()),
384 };
385 record.validate()?;
386 Ok((record, receipt))
387 }
388
389 pub(in crate::db) fn validate(&self) -> Result<(), ResumableJobError> {
390 if self.state.job_id.to_bytes() == [0; 32] {
391 return Err(ResumableJobError::InvalidJobId);
392 }
393 self.state.proof.validate()?;
394 validate_continuation(self.state.continuation.as_deref())?;
395 if self.state.application_state.len() > MAX_RESUMABLE_JOB_STATE_BYTES {
396 return Err(ResumableJobError::PayloadTooLarge);
397 }
398 if let Some(receipt) = &self.last_receipt {
399 receipt.idempotency_key.validate()?;
400 validate_continuation(receipt.continuation.as_deref())?;
401 let state_matches_receipt = match (self.state.status, receipt.status) {
402 (
403 ResumableJobStatus::Active | ResumableJobStatus::Completed,
404 ResumableJobAdvanceStatus::Advanced,
405 ) => self.state.continuation == receipt.continuation,
406 (ResumableJobStatus::Invalidated, ResumableJobAdvanceStatus::Invalidated) => {
407 self.state.continuation.is_none() && receipt.continuation.is_none()
408 }
409 _ => false,
410 };
411 if receipt.application_receipt.len() > MAX_RESUMABLE_JOB_RECEIPT_BYTES
412 || receipt.committed_sequence != self.state.sequence
413 || receipt.request_sequence.checked_add(1) != Some(receipt.committed_sequence)
414 || !state_matches_receipt
415 {
416 return Err(ResumableJobError::CorruptProgressStore);
417 }
418 } else if self.state.sequence != 0
419 || self.state.status != ResumableJobStatus::Active
420 || self.state.continuation.is_some()
421 {
422 return Err(ResumableJobError::CorruptProgressStore);
423 }
424 Ok(())
425 }
426}
427
428fn validate_continuation(continuation: Option<&str>) -> Result<(), ResumableJobError> {
429 if continuation.is_some_and(|value| value.len() > MAX_RESUMABLE_JOB_CONTINUATION_BYTES) {
430 return Err(ResumableJobError::PayloadTooLarge);
431 }
432 Ok(())
433}
434
435pub(in crate::db) fn encode_resumable_job_payload(
436 record: &ResumableJobRecord,
437) -> Result<Vec<u8>, ResumableJobError> {
438 record.validate()?;
439 let mut bytes = Vec::new();
440 bytes.extend_from_slice(&record.state.job_id.to_bytes());
441 bytes.extend_from_slice(&record.state.sequence.to_be_bytes());
442 bytes.push(match record.state.status {
443 ResumableJobStatus::Active => 0,
444 ResumableJobStatus::Invalidated => 1,
445 ResumableJobStatus::Completed => 2,
446 });
447 write_proof(&mut bytes, &record.state.proof)?;
448 write_optional_string(&mut bytes, record.state.continuation.as_deref())?;
449 write_bytes(&mut bytes, &record.state.application_state)?;
450 match &record.last_receipt {
451 None => bytes.push(0),
452 Some(receipt) => {
453 bytes.push(1);
454 bytes.extend_from_slice(&receipt.request_sequence.to_be_bytes());
455 bytes.extend_from_slice(&receipt.committed_sequence.to_be_bytes());
456 bytes.push(match receipt.status {
457 ResumableJobAdvanceStatus::Advanced => 0,
458 ResumableJobAdvanceStatus::Invalidated => 1,
459 });
460 write_string(&mut bytes, receipt.idempotency_key.as_str())?;
461 write_optional_string(&mut bytes, receipt.continuation.as_deref())?;
462 write_bytes(&mut bytes, &receipt.application_receipt)?;
463 }
464 }
465 Ok(bytes)
466}
467
468pub(in crate::db) fn decode_resumable_job_payload(
469 bytes: &[u8],
470) -> Result<ResumableJobRecord, ResumableJobError> {
471 let mut reader = ByteReader::new(bytes);
472 let job_id = ResumableJobId::try_from_bytes(reader.read_array()?)?;
473 let sequence = reader.read_u64()?;
474 let status = match reader.read_u8()? {
475 0 => ResumableJobStatus::Active,
476 1 => ResumableJobStatus::Invalidated,
477 2 => ResumableJobStatus::Completed,
478 _ => return Err(ResumableJobError::CorruptProgressStore),
479 };
480 let proof = read_proof(&mut reader)?;
481 let continuation = read_optional_string(&mut reader, MAX_RESUMABLE_JOB_CONTINUATION_BYTES)?;
482 let application_state = reader
483 .read_bounded_len_prefixed_bytes(MAX_RESUMABLE_JOB_STATE_BYTES)?
484 .to_vec();
485 let last_receipt = match reader.read_u8()? {
486 0 => None,
487 1 => {
488 let request_sequence = reader.read_u64()?;
489 let committed_sequence = reader.read_u64()?;
490 let receipt_status = match reader.read_u8()? {
491 0 => ResumableJobAdvanceStatus::Advanced,
492 1 => ResumableJobAdvanceStatus::Invalidated,
493 _ => return Err(ResumableJobError::CorruptProgressStore),
494 };
495 let idempotency_key = ResumableJobIdempotencyKey::new(
496 reader.read_bounded_string(MAX_RESUMABLE_JOB_IDEMPOTENCY_KEY_BYTES)?,
497 )?;
498 let receipt_continuation =
499 read_optional_string(&mut reader, MAX_RESUMABLE_JOB_CONTINUATION_BYTES)?;
500 let application_receipt = reader
501 .read_bounded_len_prefixed_bytes(MAX_RESUMABLE_JOB_RECEIPT_BYTES)?
502 .to_vec();
503 Some(ResumableJobAdvanceReceipt {
504 request_sequence,
505 committed_sequence,
506 status: receipt_status,
507 continuation: receipt_continuation,
508 application_receipt,
509 idempotency_key,
510 })
511 }
512 _ => return Err(ResumableJobError::CorruptProgressStore),
513 };
514 reader.finish()?;
515 let record = ResumableJobRecord {
516 state: ResumableJobState {
517 job_id,
518 sequence,
519 status,
520 proof,
521 continuation,
522 application_state,
523 },
524 last_receipt,
525 };
526 record.validate()?;
527 Ok(record)
528}
529
530fn write_proof(bytes: &mut Vec<u8>, proof: &ReadSetRevisionProof) -> Result<(), ResumableJobError> {
531 proof.validate()?;
532 bytes.extend_from_slice(&proof.database_incarnation());
533 bytes.extend_from_slice(&proof.accepted_root_revision().to_be_bytes());
534 bytes.push(proof.accepted_root_fingerprint_method());
535 bytes.extend_from_slice(&proof.accepted_root_fingerprint());
536 let count =
537 u32::try_from(proof.stores().len()).map_err(|_| ResumableJobError::PayloadTooLarge)?;
538 bytes.extend_from_slice(&count.to_be_bytes());
539 for store in proof.stores() {
540 bytes.extend_from_slice(&store.store().to_bytes());
541 bytes.extend_from_slice(&store.data_revision().to_be_bytes());
542 bytes.extend_from_slice(&store.access_state_revision().to_be_bytes());
543 }
544 Ok(())
545}
546
547fn read_proof(reader: &mut ByteReader<'_>) -> Result<ReadSetRevisionProof, ResumableJobError> {
548 let database_incarnation = reader.read_array()?;
549 let accepted_root_revision = reader.read_u64()?;
550 let accepted_root_fingerprint_method = reader.read_u8()?;
551 let accepted_root_fingerprint = reader.read_array()?;
552 let count = reader.read_u32()? as usize;
553 if count == 0 || count > crate::db::MAX_READ_SET_PROOF_STORES {
554 return Err(ResumableJobError::CorruptProgressStore);
555 }
556 let mut stores = Vec::with_capacity(count);
557 for _ in 0..count {
558 stores.push(ReadSetStoreRevision::new(
559 ReadSetStoreIdentity::from_bytes(reader.read_array()?),
560 reader.read_u64()?,
561 reader.read_u64()?,
562 ));
563 }
564 ReadSetRevisionProof::from_parts(
565 database_incarnation,
566 accepted_root_revision,
567 accepted_root_fingerprint_method,
568 accepted_root_fingerprint,
569 stores,
570 )
571 .map_err(Into::into)
572}
573
574fn write_string(bytes: &mut Vec<u8>, value: &str) -> Result<(), ResumableJobError> {
575 write_bytes(bytes, value.as_bytes())
576}
577
578fn write_optional_string(
579 bytes: &mut Vec<u8>,
580 value: Option<&str>,
581) -> Result<(), ResumableJobError> {
582 match value {
583 None => bytes.push(0),
584 Some(value) => {
585 bytes.push(1);
586 write_string(bytes, value)?;
587 }
588 }
589 Ok(())
590}
591
592fn write_bytes(bytes: &mut Vec<u8>, value: &[u8]) -> Result<(), ResumableJobError> {
593 let len = u32::try_from(value.len()).map_err(|_| ResumableJobError::PayloadTooLarge)?;
594 bytes.extend_from_slice(&len.to_be_bytes());
595 bytes.extend_from_slice(value);
596 Ok(())
597}
598
599fn read_optional_string(
600 reader: &mut ByteReader<'_>,
601 max: usize,
602) -> Result<Option<String>, ResumableJobError> {
603 match reader.read_u8()? {
604 0 => Ok(None),
605 1 => Ok(Some(reader.read_bounded_string(max)?)),
606 _ => Err(ResumableJobError::CorruptProgressStore),
607 }
608}
609
610#[cfg(test)]
611mod tests {
612 use super::*;
613
614 fn proof() -> ReadSetRevisionProof {
615 ReadSetRevisionProof::from_parts(
616 [1; 16],
617 7,
618 1,
619 [2; 32],
620 vec![ReadSetStoreRevision::new(
621 ReadSetStoreIdentity::from_bytes([3; 32]),
622 11,
623 13,
624 )],
625 )
626 .expect("bounded canonical proof should admit")
627 }
628
629 fn job_id() -> ResumableJobId {
630 ResumableJobId::try_from_bytes([4; 32]).expect("nonzero job identity should admit")
631 }
632
633 #[test]
634 fn current_resumable_job_payload_round_trips_state_and_replay_receipt() {
635 let record = ResumableJobRecord::new(job_id(), proof(), vec![1, 2, 3])
636 .expect("initial resumable record should admit");
637 let request = ResumableJobAdvanceRequest::new(
638 job_id(),
639 0,
640 ResumableJobIdempotencyKey::new("page-0")
641 .expect("bounded idempotency key should admit"),
642 );
643 let advance = ResumableJobAdvance::new(
644 Some("opaque-continuation".to_string()),
645 vec![4, 5],
646 vec![6, 7],
647 )
648 .expect("bounded advance should admit");
649 let (advanced, _) = record
650 .apply_advance(&request, advance)
651 .expect("current request should advance");
652
653 let bytes = encode_resumable_job_payload(&advanced)
654 .expect("current resumable payload should encode");
655 assert!(!bytes.starts_with(b"DIDL"));
656 assert_eq!(
657 decode_resumable_job_payload(&bytes).expect("current resumable payload should decode"),
658 advanced,
659 );
660 }
661
662 #[test]
663 fn resumable_job_payload_rejects_truncation_and_trailing_bytes() {
664 let record = ResumableJobRecord::new(job_id(), proof(), Vec::new())
665 .expect("initial resumable record should admit");
666 let bytes =
667 encode_resumable_job_payload(&record).expect("current resumable payload should encode");
668
669 for end in 0..bytes.len() {
670 assert_eq!(
671 decode_resumable_job_payload(&bytes[..end]),
672 Err(ResumableJobError::CorruptProgressStore),
673 "truncation at byte {end} must retain corruption classification",
674 );
675 }
676 let mut trailing = bytes;
677 trailing.push(0);
678 assert_eq!(
679 decode_resumable_job_payload(&trailing),
680 Err(ResumableJobError::CorruptProgressStore),
681 );
682 }
683}