journal.rs (24007B)
1 use harvestcircle_application::{ 2 BoxFuture, DurableIdentityOperation, DurableOperationKind, DurableOperationPhase, 3 DurableOperationReceipt, DurableOperationRepository, DurableOperationStart, DurableRequestId, 4 DurableTerminalOutcome, OperationDiagnostic, OperationPriorState, 5 }; 6 use harvestcircle_domain::{ 7 PublicKey, SafeError, SafeErrorCode, SafeMessage, SignerAvailability, UnixTimestamp, 8 }; 9 use radroots_service_sqlite::ServiceSqliteTransaction; 10 use sqlx::Row; 11 12 use crate::contract::{ 13 HARVESTCIRCLE_DURABLE_OPERATION_CAPACITY, HARVESTCIRCLE_DURABLE_OPERATION_CLEANUP_BATCH, 14 HARVESTCIRCLE_TERMINAL_RECEIPT_RETENTION_SECONDS, 15 }; 16 use crate::db::{corrupt_storage, map_transaction_error, storage_unavailable}; 17 use crate::{Database, HARVESTCIRCLE_UNFINISHED_DURABLE_OPERATION_CAPACITY}; 18 19 const DURABLE_OPERATION_PROJECTION: &str = "SELECT \ 20 substr(CAST(request_id AS BLOB), 1, 37) AS request_id, length(CAST(request_id AS BLOB)) AS request_id_bytes, \ 21 substr(CAST(operation_kind AS BLOB), 1, 7) AS operation_kind, length(CAST(operation_kind AS BLOB)) AS operation_kind_bytes, \ 22 substr(account_public_key, 1, 33) AS account_public_key, length(account_public_key) AS account_public_key_bytes, \ 23 expected_revision, substr(CAST(phase AS BLOB), 1, 21) AS phase, length(CAST(phase AS BLOB)) AS phase_bytes, \ 24 CASE WHEN prior_selected_public_key IS NULL THEN NULL ELSE substr(prior_selected_public_key, 1, 33) END AS prior_selected_public_key, \ 25 CASE WHEN prior_selected_public_key IS NULL THEN NULL ELSE length(prior_selected_public_key) END AS prior_selected_public_key_bytes, \ 26 updated_at_unix_s, \ 27 CASE WHEN diagnostic_code IS NULL THEN NULL ELSE substr(CAST(diagnostic_code AS BLOB), 1, 21) END AS diagnostic_code, \ 28 CASE WHEN diagnostic_code IS NULL THEN NULL ELSE length(CAST(diagnostic_code AS BLOB)) END AS diagnostic_code_bytes, \ 29 CASE WHEN terminal_outcome IS NULL THEN NULL ELSE substr(CAST(terminal_outcome AS BLOB), 1, 10) END AS terminal_outcome, \ 30 CASE WHEN terminal_outcome IS NULL THEN NULL ELSE length(CAST(terminal_outcome AS BLOB)) END AS terminal_outcome_bytes, \ 31 CASE WHEN prior_binding_availability IS NULL THEN NULL ELSE substr(CAST(prior_binding_availability AS BLOB), 1, 19) END AS prior_binding_availability, \ 32 CASE WHEN prior_binding_availability IS NULL THEN NULL ELSE length(CAST(prior_binding_availability AS BLOB)) END AS prior_binding_availability_bytes, \ 33 resulting_revision, completed_at_unix_s FROM durable_operations"; 34 35 impl DurableOperationRepository for Database { 36 #[allow(clippy::too_many_arguments)] 37 fn begin_durable_operation<'a>( 38 &'a self, 39 request_id: &'a DurableRequestId, 40 kind: DurableOperationKind, 41 identity: PublicKey, 42 expected_revision: Option<u64>, 43 prior: OperationPriorState, 44 updated_at: UnixTimestamp, 45 ) -> BoxFuture<'a, Result<DurableOperationStart, SafeError>> { 46 Box::pin(async move { 47 let expected_revision = encode_revision(expected_revision)?; 48 let request_id = request_id.as_str().to_owned(); 49 self.host() 50 .transaction(|transaction| { 51 Box::pin(async move { 52 if let Some(operation) = query_operation(transaction, &request_id).await? { 53 if operation.kind() != kind 54 || operation.identity() != identity 55 || operation.expected_revision() 56 != expected_revision.map(|value| value as u64) 57 || operation.prior() != prior 58 { 59 return Err(operation_conflict()); 60 } 61 return Ok(DurableOperationStart::Existing(operation)); 62 } 63 cleanup_expired_terminal_receipts(transaction, updated_at).await?; 64 let unfinished: i64 = sqlx::query_scalar( 65 "SELECT count(*) FROM (SELECT 1 FROM durable_operations \ 66 WHERE terminal_outcome IS NULL LIMIT 1025)", 67 ) 68 .fetch_one(&mut *transaction) 69 .await 70 .map_err(|_| storage_unavailable())?; 71 if usize::try_from(unfinished).ok().is_none_or(|count| { 72 count >= HARVESTCIRCLE_UNFINISHED_DURABLE_OPERATION_CAPACITY 73 }) { 74 return Err(operation_capacity_exhausted()); 75 } 76 let total: i64 = sqlx::query_scalar( 77 "SELECT count(*) FROM (SELECT 1 FROM durable_operations LIMIT 4097)", 78 ) 79 .fetch_one(&mut *transaction) 80 .await 81 .map_err(|_| storage_unavailable())?; 82 if usize::try_from(total) 83 .ok() 84 .is_none_or(|count| count >= HARVESTCIRCLE_DURABLE_OPERATION_CAPACITY) 85 { 86 return Err(operation_capacity_exhausted()); 87 } 88 let result = sqlx::query( 89 "INSERT INTO durable_operations (request_id, operation_kind, \ 90 account_public_key, binding_public_key, expected_revision, phase, \ 91 prior_selected_public_key, updated_at_unix_s, prior_binding_availability) \ 92 VALUES (?, ?, ?, ?, ?, 'intent_recorded', ?, ?, ?)", 93 ) 94 .bind(&request_id) 95 .bind(encode_kind(kind)) 96 .bind(identity.as_bytes().as_slice()) 97 .bind(identity.as_bytes().as_slice()) 98 .bind(expected_revision) 99 .bind(prior.selected_identity().map(|value| value.as_bytes().to_vec())) 100 .bind(updated_at.as_seconds()) 101 .bind(prior.binding_availability().map(encode_availability)) 102 .execute(&mut *transaction) 103 .await 104 .map_err(|_| storage_unavailable())?; 105 if result.rows_affected() != 1 { 106 return Err(corrupt_storage()); 107 } 108 let operation = query_operation(transaction, &request_id) 109 .await? 110 .ok_or_else(corrupt_storage)?; 111 Ok(DurableOperationStart::Started(operation)) 112 }) 113 }) 114 .await 115 .map_err(map_transaction_error) 116 }) 117 } 118 119 fn load_durable_operation<'a>( 120 &'a self, 121 request_id: &'a DurableRequestId, 122 ) -> BoxFuture<'a, Result<Option<DurableIdentityOperation>, SafeError>> { 123 Box::pin(async move { 124 let request_id = request_id.as_str().to_owned(); 125 self.host() 126 .transaction(|transaction| { 127 Box::pin(async move { query_operation(transaction, &request_id).await }) 128 }) 129 .await 130 .map_err(map_transaction_error) 131 }) 132 } 133 134 fn advance_durable_operation<'a>( 135 &'a self, 136 request_id: &'a DurableRequestId, 137 expected_phase: DurableOperationPhase, 138 next_phase: DurableOperationPhase, 139 updated_at: UnixTimestamp, 140 diagnostic: Option<OperationDiagnostic>, 141 ) -> BoxFuture<'a, Result<DurableIdentityOperation, SafeError>> { 142 Box::pin(async move { 143 let request_id = request_id.as_str().to_owned(); 144 self.host() 145 .transaction(|transaction| { 146 Box::pin(async move { 147 let result = sqlx::query( 148 "UPDATE durable_operations SET phase = ?, updated_at_unix_s = ?, \ 149 diagnostic_code = ? WHERE request_id = ? AND phase = ? \ 150 AND terminal_outcome IS NULL", 151 ) 152 .bind(encode_phase(next_phase)) 153 .bind(updated_at.as_seconds()) 154 .bind(diagnostic.map(encode_diagnostic)) 155 .bind(&request_id) 156 .bind(encode_phase(expected_phase)) 157 .execute(&mut *transaction) 158 .await 159 .map_err(|_| storage_unavailable())?; 160 if result.rows_affected() != 1 { 161 return Err(operation_conflict()); 162 } 163 query_operation(transaction, &request_id) 164 .await? 165 .ok_or_else(corrupt_storage) 166 }) 167 }) 168 .await 169 .map_err(map_transaction_error) 170 }) 171 } 172 173 fn finalize_durable_operation<'a>( 174 &'a self, 175 request_id: &'a DurableRequestId, 176 expected_phase: DurableOperationPhase, 177 outcome: DurableTerminalOutcome, 178 resulting_revision: Option<u64>, 179 updated_at: UnixTimestamp, 180 ) -> BoxFuture<'a, Result<DurableOperationReceipt, SafeError>> { 181 Box::pin(async move { 182 let resulting_revision = encode_revision(resulting_revision)?; 183 let request_id = request_id.as_str().to_owned(); 184 self.host() 185 .transaction(|transaction| { 186 Box::pin(async move { 187 if let Some(existing) = query_operation(transaction, &request_id).await? 188 && let Some(receipt) = existing.terminal() 189 { 190 return if receipt.outcome() == outcome 191 && receipt.resulting_revision() 192 == resulting_revision.map(|value| value as u64) 193 { 194 Ok(receipt.clone()) 195 } else { 196 Err(operation_conflict()) 197 }; 198 } 199 let result = sqlx::query( 200 "UPDATE durable_operations SET phase = 'finalized', \ 201 terminal_outcome = ?, resulting_revision = ?, updated_at_unix_s = ?, \ 202 completed_at_unix_s = ? \ 203 WHERE request_id = ? AND phase = ? AND terminal_outcome IS NULL", 204 ) 205 .bind(encode_outcome(outcome)) 206 .bind(resulting_revision) 207 .bind(updated_at.as_seconds()) 208 .bind(updated_at.as_seconds()) 209 .bind(&request_id) 210 .bind(encode_phase(expected_phase)) 211 .execute(&mut *transaction) 212 .await 213 .map_err(|_| storage_unavailable())?; 214 if result.rows_affected() != 1 { 215 return Err(operation_conflict()); 216 } 217 query_operation(transaction, &request_id) 218 .await? 219 .and_then(|operation| operation.terminal().cloned()) 220 .ok_or_else(corrupt_storage) 221 }) 222 }) 223 .await 224 .map_err(map_transaction_error) 225 }) 226 } 227 228 fn list_unfinished_durable_operations( 229 &self, 230 ) -> BoxFuture<'_, Result<Vec<DurableIdentityOperation>, SafeError>> { 231 Box::pin(async move { 232 self.host() 233 .transaction(|transaction| { 234 Box::pin(async move { 235 let sql = format!( 236 "{DURABLE_OPERATION_PROJECTION} WHERE terminal_outcome IS NULL \ 237 ORDER BY request_id LIMIT {}", 238 HARVESTCIRCLE_UNFINISHED_DURABLE_OPERATION_CAPACITY + 1 239 ); 240 let rows = sqlx::query(sqlx::AssertSqlSafe(sql.as_str())) 241 .fetch_all(&mut *transaction) 242 .await 243 .map_err(|_| corrupt_storage())?; 244 if rows.len() > HARVESTCIRCLE_UNFINISHED_DURABLE_OPERATION_CAPACITY { 245 return Err(corrupt_storage()); 246 } 247 rows.iter().map(decode_operation).collect() 248 }) 249 }) 250 .await 251 .map_err(map_transaction_error) 252 }) 253 } 254 } 255 256 async fn query_operation( 257 transaction: &mut ServiceSqliteTransaction<'_>, 258 request_id: &str, 259 ) -> Result<Option<DurableIdentityOperation>, SafeError> { 260 let sql = format!("{DURABLE_OPERATION_PROJECTION} WHERE request_id = ? LIMIT 2"); 261 let rows = sqlx::query(sqlx::AssertSqlSafe(sql.as_str())) 262 .bind(request_id) 263 .fetch_all(&mut *transaction) 264 .await 265 .map_err(|_| corrupt_storage())?; 266 match rows.as_slice() { 267 [] => Ok(None), 268 [row] => decode_operation(row).map(Some), 269 _ => Err(corrupt_storage()), 270 } 271 } 272 273 fn decode_operation(row: &sqlx::sqlite::SqliteRow) -> Result<DurableIdentityOperation, SafeError> { 274 let request_id = 275 DurableRequestId::parse(required_text(row, "request_id", "request_id_bytes", 36)?)?; 276 let kind = decode_kind(&required_text( 277 row, 278 "operation_kind", 279 "operation_kind_bytes", 280 6, 281 )?)?; 282 let identity = required_key(row, "account_public_key", "account_public_key_bytes")?; 283 let expected_revision = decode_revision( 284 row.try_get("expected_revision") 285 .map_err(|_| corrupt_storage())?, 286 )?; 287 let phase = decode_phase(&required_text(row, "phase", "phase_bytes", 20)?)?; 288 let prior_selected = optional_key( 289 row, 290 "prior_selected_public_key", 291 "prior_selected_public_key_bytes", 292 )?; 293 let updated_at = UnixTimestamp::from_seconds( 294 row.try_get("updated_at_unix_s") 295 .map_err(|_| corrupt_storage())?, 296 ) 297 .ok_or_else(corrupt_storage)?; 298 let diagnostic = optional_text(row, "diagnostic_code", "diagnostic_code_bytes", 20)? 299 .map(|value| decode_diagnostic(&value)) 300 .transpose()?; 301 let outcome = optional_text(row, "terminal_outcome", "terminal_outcome_bytes", 9)? 302 .map(|value| decode_outcome(&value)) 303 .transpose()?; 304 let prior_availability = optional_text( 305 row, 306 "prior_binding_availability", 307 "prior_binding_availability_bytes", 308 18, 309 )? 310 .map(|value| decode_availability(&value)) 311 .transpose()?; 312 let resulting_revision = decode_revision( 313 row.try_get("resulting_revision") 314 .map_err(|_| corrupt_storage())?, 315 )?; 316 let completed_at = row 317 .try_get::<Option<i64>, _>("completed_at_unix_s") 318 .map_err(|_| corrupt_storage())? 319 .map(|value| UnixTimestamp::from_seconds(value).ok_or_else(corrupt_storage)) 320 .transpose()?; 321 let terminal = match (outcome, completed_at) { 322 (Some(outcome), Some(completed_at)) => Some(DurableOperationReceipt::new( 323 request_id.clone(), 324 identity, 325 outcome, 326 resulting_revision, 327 completed_at, 328 )), 329 (None, None) if resulting_revision.is_none() => None, 330 _ => return Err(corrupt_storage()), 331 }; 332 Ok(DurableIdentityOperation::new( 333 request_id, 334 kind, 335 identity, 336 expected_revision, 337 phase, 338 OperationPriorState::new(prior_selected, prior_availability), 339 updated_at, 340 diagnostic, 341 terminal, 342 )) 343 } 344 345 async fn cleanup_expired_terminal_receipts( 346 transaction: &mut ServiceSqliteTransaction<'_>, 347 now: UnixTimestamp, 348 ) -> Result<(), SafeError> { 349 let cutoff = now 350 .as_seconds() 351 .saturating_sub(HARVESTCIRCLE_TERMINAL_RECEIPT_RETENTION_SECONDS); 352 let result = sqlx::query( 353 "DELETE FROM durable_operations WHERE request_id IN (\ 354 SELECT request_id FROM durable_operations \ 355 WHERE completed_at_unix_s IS NOT NULL AND completed_at_unix_s < ? \ 356 ORDER BY completed_at_unix_s, request_id LIMIT 256\ 357 )", 358 ) 359 .bind(cutoff) 360 .execute(&mut *transaction) 361 .await 362 .map_err(|_| storage_unavailable())?; 363 if result.rows_affected() 364 > u64::try_from(HARVESTCIRCLE_DURABLE_OPERATION_CLEANUP_BATCH) 365 .expect("cleanup bound fits in u64") 366 { 367 return Err(corrupt_storage()); 368 } 369 Ok(()) 370 } 371 372 fn required_key( 373 row: &sqlx::sqlite::SqliteRow, 374 value: &str, 375 length: &str, 376 ) -> Result<PublicKey, SafeError> { 377 optional_key(row, value, length)?.ok_or_else(corrupt_storage) 378 } 379 380 fn optional_key( 381 row: &sqlx::sqlite::SqliteRow, 382 value: &str, 383 length: &str, 384 ) -> Result<Option<PublicKey>, SafeError> { 385 optional_blob(row, value, length, 32)? 386 .map(|bytes| { 387 let bytes: [u8; 32] = bytes.try_into().map_err(|_| corrupt_storage())?; 388 PublicKey::from_bytes(bytes).map_err(|_| corrupt_storage()) 389 }) 390 .transpose() 391 } 392 393 fn required_text( 394 row: &sqlx::sqlite::SqliteRow, 395 value: &str, 396 length: &str, 397 maximum: usize, 398 ) -> Result<String, SafeError> { 399 optional_text(row, value, length, maximum)?.ok_or_else(corrupt_storage) 400 } 401 402 fn optional_text( 403 row: &sqlx::sqlite::SqliteRow, 404 value: &str, 405 length: &str, 406 maximum: usize, 407 ) -> Result<Option<String>, SafeError> { 408 optional_blob(row, value, length, maximum)? 409 .map(String::from_utf8) 410 .transpose() 411 .map_err(|_| corrupt_storage()) 412 } 413 414 fn optional_blob( 415 row: &sqlx::sqlite::SqliteRow, 416 value_column: &str, 417 length_column: &str, 418 maximum: usize, 419 ) -> Result<Option<Vec<u8>>, SafeError> { 420 let value = row 421 .try_get::<Option<Vec<u8>>, _>(value_column) 422 .map_err(|_| corrupt_storage())?; 423 let length = row 424 .try_get::<Option<i64>, _>(length_column) 425 .map_err(|_| corrupt_storage())?; 426 match (value, length) { 427 (None, None) => Ok(None), 428 (Some(value), Some(length)) 429 if usize::try_from(length) 430 .ok() 431 .is_some_and(|length| length <= maximum && length == value.len()) => 432 { 433 Ok(Some(value)) 434 } 435 _ => Err(corrupt_storage()), 436 } 437 } 438 439 fn encode_revision(value: Option<u64>) -> Result<Option<i64>, SafeError> { 440 value 441 .map(i64::try_from) 442 .transpose() 443 .map_err(|_| operation_conflict()) 444 } 445 446 fn decode_revision(value: Option<i64>) -> Result<Option<u64>, SafeError> { 447 value 448 .map(u64::try_from) 449 .transpose() 450 .map_err(|_| corrupt_storage()) 451 } 452 453 const fn encode_kind(value: DurableOperationKind) -> &'static str { 454 match value { 455 DurableOperationKind::Create => "create", 456 DurableOperationKind::Import => "import", 457 DurableOperationKind::Repair => "repair", 458 DurableOperationKind::Remove => "remove", 459 } 460 } 461 462 fn decode_kind(value: &str) -> Result<DurableOperationKind, SafeError> { 463 match value { 464 "create" => Ok(DurableOperationKind::Create), 465 "import" => Ok(DurableOperationKind::Import), 466 "repair" => Ok(DurableOperationKind::Repair), 467 "remove" => Ok(DurableOperationKind::Remove), 468 _ => Err(corrupt_storage()), 469 } 470 } 471 472 const fn encode_phase(value: DurableOperationPhase) -> &'static str { 473 match value { 474 DurableOperationPhase::IntentRecorded => "intent_recorded", 475 DurableOperationPhase::CredentialWritten => "credential_written", 476 DurableOperationPhase::MetadataCommitted => "metadata_committed", 477 DurableOperationPhase::SelectionCommitted => "selection_committed", 478 DurableOperationPhase::CompensationPending => "compensation_pending", 479 DurableOperationPhase::CredentialDeleted => "credential_deleted", 480 DurableOperationPhase::MetadataDeleted => "metadata_deleted", 481 DurableOperationPhase::Finalized => "finalized", 482 } 483 } 484 485 fn decode_phase(value: &str) -> Result<DurableOperationPhase, SafeError> { 486 match value { 487 "intent_recorded" => Ok(DurableOperationPhase::IntentRecorded), 488 "credential_written" => Ok(DurableOperationPhase::CredentialWritten), 489 "metadata_committed" => Ok(DurableOperationPhase::MetadataCommitted), 490 "selection_committed" => Ok(DurableOperationPhase::SelectionCommitted), 491 "compensation_pending" => Ok(DurableOperationPhase::CompensationPending), 492 "credential_deleted" => Ok(DurableOperationPhase::CredentialDeleted), 493 "metadata_deleted" => Ok(DurableOperationPhase::MetadataDeleted), 494 "finalized" => Ok(DurableOperationPhase::Finalized), 495 _ => Err(corrupt_storage()), 496 } 497 } 498 499 const fn encode_outcome(value: DurableTerminalOutcome) -> &'static str { 500 match value { 501 DurableTerminalOutcome::Completed => "completed", 502 DurableTerminalOutcome::Cancelled => "cancelled", 503 DurableTerminalOutcome::Failed => "failed", 504 } 505 } 506 507 fn decode_outcome(value: &str) -> Result<DurableTerminalOutcome, SafeError> { 508 match value { 509 "completed" => Ok(DurableTerminalOutcome::Completed), 510 "cancelled" => Ok(DurableTerminalOutcome::Cancelled), 511 "failed" => Ok(DurableTerminalOutcome::Failed), 512 _ => Err(corrupt_storage()), 513 } 514 } 515 516 const fn encode_diagnostic(value: OperationDiagnostic) -> &'static str { 517 match value { 518 OperationDiagnostic::StorageUnavailable => "storage_unavailable", 519 OperationDiagnostic::KeyringUnavailable => "keyring_unavailable", 520 OperationDiagnostic::CredentialMissing => "credential_missing", 521 OperationDiagnostic::CompensationFailed => "compensation_failed", 522 OperationDiagnostic::Conflict => "conflict", 523 OperationDiagnostic::Expired => "expired", 524 } 525 } 526 527 fn decode_diagnostic(value: &str) -> Result<OperationDiagnostic, SafeError> { 528 match value { 529 "storage_unavailable" => Ok(OperationDiagnostic::StorageUnavailable), 530 "keyring_unavailable" => Ok(OperationDiagnostic::KeyringUnavailable), 531 "credential_missing" => Ok(OperationDiagnostic::CredentialMissing), 532 "compensation_failed" => Ok(OperationDiagnostic::CompensationFailed), 533 "conflict" => Ok(OperationDiagnostic::Conflict), 534 "expired" => Ok(OperationDiagnostic::Expired), 535 _ => Err(corrupt_storage()), 536 } 537 } 538 539 const fn encode_availability(value: SignerAvailability) -> &'static str { 540 match value { 541 SignerAvailability::Available => "available", 542 SignerAvailability::CredentialMissing => "credential_missing", 543 SignerAvailability::StoreUnavailable => "store_unavailable", 544 } 545 } 546 547 fn decode_availability(value: &str) -> Result<SignerAvailability, SafeError> { 548 match value { 549 "available" => Ok(SignerAvailability::Available), 550 "credential_missing" => Ok(SignerAvailability::CredentialMissing), 551 "store_unavailable" => Ok(SignerAvailability::StoreUnavailable), 552 _ => Err(corrupt_storage()), 553 } 554 } 555 556 const fn operation_conflict() -> SafeError { 557 SafeError::new( 558 SafeErrorCode::InvalidApplicationState, 559 SafeMessage::new("The durable operation conflicts with existing state."), 560 ) 561 } 562 563 const fn operation_capacity_exhausted() -> SafeError { 564 SafeError::new( 565 SafeErrorCode::InvalidApplicationState, 566 SafeMessage::new("The unfinished operation capacity is exhausted."), 567 ) 568 }