Our Personal Data Server from scratch!
0

Configure Feed

Select the types of activity you want to include in your feed.

1use async_trait::async_trait; 2use chrono::{DateTime, Utc}; 3use sqlx::PgPool; 4use tranquil_db_traits::{ 5 AccountStatus, BrokenGenesisCommit, CommitEventData, DbError, EventBlocksCids, FullRecordInfo, 6 ImportBlock, ImportRecord, ImportRepoError, RecordInfo, RecordWithTakedown, RepoAccountInfo, 7 RepoEventType, RepoInfo, RepoListItem, RepoRepository, RepoWithoutRev, SequenceNumber, 8 SequencedEvent, UserNeedingRecordBlobsBackfill, UserWithoutBlocks, 9}; 10use tranquil_types::{AtUri, CidLink, Did, Handle, Nsid, Rkey}; 11use uuid::Uuid; 12 13use super::user::map_sqlx_error; 14 15struct RecordRow { 16 rkey: String, 17 record_cid: String, 18} 19 20struct SequencedEventRow { 21 seq: i64, 22 did: String, 23 created_at: DateTime<Utc>, 24 event_type: RepoEventType, 25 commit_cid: Option<String>, 26 prev_cid: Option<String>, 27 prev_data_cid: Option<String>, 28 ops: Option<serde_json::Value>, 29 blobs: Option<Vec<String>>, 30 blocks_cids: Option<Vec<String>>, 31 handle: Option<String>, 32 active: Option<bool>, 33 status: Option<String>, 34 rev: Option<String>, 35} 36 37pub struct PostgresRepoRepository { 38 pool: PgPool, 39} 40 41impl PostgresRepoRepository { 42 pub fn new(pool: PgPool) -> Self { 43 Self { pool } 44 } 45} 46 47#[async_trait] 48impl RepoRepository for PostgresRepoRepository { 49 async fn create_repo( 50 &self, 51 user_id: Uuid, 52 repo_root_cid: &CidLink, 53 repo_rev: &str, 54 ) -> Result<(), DbError> { 55 sqlx::query!( 56 "INSERT INTO repos (user_id, repo_root_cid, repo_rev) VALUES ($1, $2, $3)", 57 user_id, 58 repo_root_cid.as_str(), 59 repo_rev 60 ) 61 .execute(&self.pool) 62 .await 63 .map_err(map_sqlx_error)?; 64 65 Ok(()) 66 } 67 68 async fn update_repo_root( 69 &self, 70 user_id: Uuid, 71 repo_root_cid: &CidLink, 72 repo_rev: &str, 73 ) -> Result<(), DbError> { 74 sqlx::query!( 75 "UPDATE repos SET repo_root_cid = $1, repo_rev = $2, updated_at = NOW() WHERE user_id = $3", 76 repo_root_cid.as_str(), 77 repo_rev, 78 user_id 79 ) 80 .execute(&self.pool) 81 .await 82 .map_err(map_sqlx_error)?; 83 84 Ok(()) 85 } 86 87 async fn update_repo_rev(&self, user_id: Uuid, repo_rev: &str) -> Result<(), DbError> { 88 sqlx::query!( 89 "UPDATE repos SET repo_rev = $1 WHERE user_id = $2", 90 repo_rev, 91 user_id 92 ) 93 .execute(&self.pool) 94 .await 95 .map_err(map_sqlx_error)?; 96 97 Ok(()) 98 } 99 100 async fn delete_repo(&self, user_id: Uuid) -> Result<(), DbError> { 101 sqlx::query!("DELETE FROM repos WHERE user_id = $1", user_id) 102 .execute(&self.pool) 103 .await 104 .map_err(map_sqlx_error)?; 105 106 Ok(()) 107 } 108 109 async fn get_repo_root_for_update(&self, user_id: Uuid) -> Result<Option<CidLink>, DbError> { 110 let result = sqlx::query_scalar!( 111 "SELECT repo_root_cid FROM repos WHERE user_id = $1 FOR UPDATE NOWAIT", 112 user_id 113 ) 114 .fetch_optional(&self.pool) 115 .await 116 .map_err(map_sqlx_error)?; 117 118 Ok(result.map(CidLink::from)) 119 } 120 121 async fn get_repo(&self, user_id: Uuid) -> Result<Option<RepoInfo>, DbError> { 122 let row = sqlx::query!( 123 "SELECT user_id, repo_root_cid, repo_rev FROM repos WHERE user_id = $1", 124 user_id 125 ) 126 .fetch_optional(&self.pool) 127 .await 128 .map_err(map_sqlx_error)?; 129 130 Ok(row.map(|r| RepoInfo { 131 user_id: r.user_id, 132 repo_root_cid: CidLink::from(r.repo_root_cid), 133 repo_rev: r.repo_rev, 134 })) 135 } 136 137 async fn get_repo_root_by_did(&self, did: &Did) -> Result<Option<CidLink>, DbError> { 138 let result = sqlx::query_scalar!( 139 "SELECT r.repo_root_cid FROM repos r JOIN users u ON r.user_id = u.id WHERE u.did = $1", 140 did.as_str() 141 ) 142 .fetch_optional(&self.pool) 143 .await 144 .map_err(map_sqlx_error)?; 145 146 Ok(result.map(CidLink::from)) 147 } 148 149 async fn count_repos(&self) -> Result<i64, DbError> { 150 let count = sqlx::query_scalar!(r#"SELECT COUNT(*) as "count!" FROM repos"#) 151 .fetch_one(&self.pool) 152 .await 153 .map_err(map_sqlx_error)?; 154 155 Ok(count) 156 } 157 158 async fn get_repos_without_rev(&self) -> Result<Vec<RepoWithoutRev>, DbError> { 159 let rows = sqlx::query!("SELECT user_id, repo_root_cid FROM repos WHERE repo_rev IS NULL") 160 .fetch_all(&self.pool) 161 .await 162 .map_err(map_sqlx_error)?; 163 164 Ok(rows 165 .into_iter() 166 .map(|r| RepoWithoutRev { 167 user_id: r.user_id, 168 repo_root_cid: CidLink::from(r.repo_root_cid), 169 }) 170 .collect()) 171 } 172 173 async fn upsert_records( 174 &self, 175 repo_id: Uuid, 176 collections: &[Nsid], 177 rkeys: &[Rkey], 178 record_cids: &[CidLink], 179 repo_rev: &str, 180 ) -> Result<(), DbError> { 181 let collections_str: Vec<&str> = collections.iter().map(|c| c.as_str()).collect(); 182 let rkeys_str: Vec<&str> = rkeys.iter().map(|r| r.as_str()).collect(); 183 let cids_str: Vec<&str> = record_cids.iter().map(|c| c.as_str()).collect(); 184 185 sqlx::query!( 186 r#" 187 INSERT INTO records (repo_id, collection, rkey, record_cid, repo_rev) 188 SELECT $1, collection, rkey, record_cid, $5 189 FROM UNNEST($2::text[], $3::text[], $4::text[]) AS t(collection, rkey, record_cid) 190 ON CONFLICT (repo_id, collection, rkey) DO UPDATE 191 SET record_cid = EXCLUDED.record_cid, repo_rev = EXCLUDED.repo_rev, created_at = NOW() 192 "#, 193 repo_id, 194 &collections_str as &[&str], 195 &rkeys_str as &[&str], 196 &cids_str as &[&str], 197 repo_rev 198 ) 199 .execute(&self.pool) 200 .await 201 .map_err(map_sqlx_error)?; 202 203 Ok(()) 204 } 205 206 async fn delete_records( 207 &self, 208 repo_id: Uuid, 209 collections: &[Nsid], 210 rkeys: &[Rkey], 211 ) -> Result<(), DbError> { 212 let collections_str: Vec<&str> = collections.iter().map(|c| c.as_str()).collect(); 213 let rkeys_str: Vec<&str> = rkeys.iter().map(|r| r.as_str()).collect(); 214 215 sqlx::query!( 216 r#" 217 DELETE FROM records 218 WHERE repo_id = $1 219 AND (collection, rkey) IN (SELECT * FROM UNNEST($2::text[], $3::text[])) 220 "#, 221 repo_id, 222 &collections_str as &[&str], 223 &rkeys_str as &[&str] 224 ) 225 .execute(&self.pool) 226 .await 227 .map_err(map_sqlx_error)?; 228 229 Ok(()) 230 } 231 232 async fn delete_all_records(&self, repo_id: Uuid) -> Result<(), DbError> { 233 sqlx::query!("DELETE FROM records WHERE repo_id = $1", repo_id) 234 .execute(&self.pool) 235 .await 236 .map_err(map_sqlx_error)?; 237 238 Ok(()) 239 } 240 241 async fn get_record_cid( 242 &self, 243 repo_id: Uuid, 244 collection: &Nsid, 245 rkey: &Rkey, 246 ) -> Result<Option<CidLink>, DbError> { 247 let result = sqlx::query_scalar!( 248 "SELECT record_cid FROM records WHERE repo_id = $1 AND collection = $2 AND rkey = $3", 249 repo_id, 250 collection.as_str(), 251 rkey.as_str() 252 ) 253 .fetch_optional(&self.pool) 254 .await 255 .map_err(map_sqlx_error)?; 256 257 Ok(result.map(CidLink::from)) 258 } 259 260 async fn list_records( 261 &self, 262 repo_id: Uuid, 263 collection: &Nsid, 264 cursor: Option<&Rkey>, 265 limit: i64, 266 reverse: bool, 267 rkey_start: Option<&Rkey>, 268 rkey_end: Option<&Rkey>, 269 ) -> Result<Vec<RecordInfo>, DbError> { 270 let to_record_info = |rows: Vec<RecordRow>| { 271 rows.into_iter() 272 .map(|r| RecordInfo { 273 rkey: Rkey::from(r.rkey), 274 record_cid: CidLink::from(r.record_cid), 275 }) 276 .collect() 277 }; 278 279 let collection_str = collection.as_str(); 280 281 if let Some(cursor_val) = cursor { 282 let cursor_str = cursor_val.as_str(); 283 return match reverse { 284 false => { 285 let rows = sqlx::query_as!( 286 RecordRow, 287 r#"SELECT rkey, record_cid FROM records 288 WHERE repo_id = $1 AND collection = $2 AND rkey < $3 289 ORDER BY rkey DESC LIMIT $4"#, 290 repo_id, 291 collection_str, 292 cursor_str, 293 limit 294 ) 295 .fetch_all(&self.pool) 296 .await 297 .map_err(map_sqlx_error)?; 298 Ok(to_record_info(rows)) 299 } 300 true => { 301 let rows = sqlx::query_as!( 302 RecordRow, 303 r#"SELECT rkey, record_cid FROM records 304 WHERE repo_id = $1 AND collection = $2 AND rkey > $3 305 ORDER BY rkey ASC LIMIT $4"#, 306 repo_id, 307 collection_str, 308 cursor_str, 309 limit 310 ) 311 .fetch_all(&self.pool) 312 .await 313 .map_err(map_sqlx_error)?; 314 Ok(to_record_info(rows)) 315 } 316 }; 317 } 318 319 if let (Some(start), Some(end)) = (rkey_start, rkey_end) { 320 let start_str = start.as_str(); 321 let end_str = end.as_str(); 322 return match reverse { 323 false => { 324 let rows = sqlx::query_as!( 325 RecordRow, 326 r#"SELECT rkey, record_cid FROM records 327 WHERE repo_id = $1 AND collection = $2 AND rkey >= $3 AND rkey <= $4 328 ORDER BY rkey DESC LIMIT $5"#, 329 repo_id, 330 collection_str, 331 start_str, 332 end_str, 333 limit 334 ) 335 .fetch_all(&self.pool) 336 .await 337 .map_err(map_sqlx_error)?; 338 Ok(to_record_info(rows)) 339 } 340 true => { 341 let rows = sqlx::query_as!( 342 RecordRow, 343 r#"SELECT rkey, record_cid FROM records 344 WHERE repo_id = $1 AND collection = $2 AND rkey >= $3 AND rkey <= $4 345 ORDER BY rkey ASC LIMIT $5"#, 346 repo_id, 347 collection_str, 348 start_str, 349 end_str, 350 limit 351 ) 352 .fetch_all(&self.pool) 353 .await 354 .map_err(map_sqlx_error)?; 355 Ok(to_record_info(rows)) 356 } 357 }; 358 } 359 360 if let Some(start) = rkey_start { 361 let start_str = start.as_str(); 362 return match reverse { 363 false => { 364 let rows = sqlx::query_as!( 365 RecordRow, 366 r#"SELECT rkey, record_cid FROM records 367 WHERE repo_id = $1 AND collection = $2 AND rkey >= $3 368 ORDER BY rkey DESC LIMIT $4"#, 369 repo_id, 370 collection_str, 371 start_str, 372 limit 373 ) 374 .fetch_all(&self.pool) 375 .await 376 .map_err(map_sqlx_error)?; 377 Ok(to_record_info(rows)) 378 } 379 true => { 380 let rows = sqlx::query_as!( 381 RecordRow, 382 r#"SELECT rkey, record_cid FROM records 383 WHERE repo_id = $1 AND collection = $2 AND rkey >= $3 384 ORDER BY rkey ASC LIMIT $4"#, 385 repo_id, 386 collection_str, 387 start_str, 388 limit 389 ) 390 .fetch_all(&self.pool) 391 .await 392 .map_err(map_sqlx_error)?; 393 Ok(to_record_info(rows)) 394 } 395 }; 396 } 397 398 if let Some(end) = rkey_end { 399 let end_str = end.as_str(); 400 return match reverse { 401 false => { 402 let rows = sqlx::query_as!( 403 RecordRow, 404 r#"SELECT rkey, record_cid FROM records 405 WHERE repo_id = $1 AND collection = $2 AND rkey <= $3 406 ORDER BY rkey DESC LIMIT $4"#, 407 repo_id, 408 collection_str, 409 end_str, 410 limit 411 ) 412 .fetch_all(&self.pool) 413 .await 414 .map_err(map_sqlx_error)?; 415 Ok(to_record_info(rows)) 416 } 417 true => { 418 let rows = sqlx::query_as!( 419 RecordRow, 420 r#"SELECT rkey, record_cid FROM records 421 WHERE repo_id = $1 AND collection = $2 AND rkey <= $3 422 ORDER BY rkey ASC LIMIT $4"#, 423 repo_id, 424 collection_str, 425 end_str, 426 limit 427 ) 428 .fetch_all(&self.pool) 429 .await 430 .map_err(map_sqlx_error)?; 431 Ok(to_record_info(rows)) 432 } 433 }; 434 } 435 436 match reverse { 437 false => { 438 let rows = sqlx::query_as!( 439 RecordRow, 440 r#"SELECT rkey, record_cid FROM records 441 WHERE repo_id = $1 AND collection = $2 442 ORDER BY rkey DESC LIMIT $3"#, 443 repo_id, 444 collection_str, 445 limit 446 ) 447 .fetch_all(&self.pool) 448 .await 449 .map_err(map_sqlx_error)?; 450 Ok(to_record_info(rows)) 451 } 452 true => { 453 let rows = sqlx::query_as!( 454 RecordRow, 455 r#"SELECT rkey, record_cid FROM records 456 WHERE repo_id = $1 AND collection = $2 457 ORDER BY rkey ASC LIMIT $3"#, 458 repo_id, 459 collection_str, 460 limit 461 ) 462 .fetch_all(&self.pool) 463 .await 464 .map_err(map_sqlx_error)?; 465 Ok(to_record_info(rows)) 466 } 467 } 468 } 469 470 async fn get_all_records(&self, repo_id: Uuid) -> Result<Vec<FullRecordInfo>, DbError> { 471 let rows = sqlx::query!( 472 "SELECT collection, rkey, record_cid FROM records WHERE repo_id = $1", 473 repo_id 474 ) 475 .fetch_all(&self.pool) 476 .await 477 .map_err(map_sqlx_error)?; 478 479 Ok(rows 480 .into_iter() 481 .map(|r| FullRecordInfo { 482 collection: Nsid::from(r.collection), 483 rkey: Rkey::from(r.rkey), 484 record_cid: CidLink::from(r.record_cid), 485 }) 486 .collect()) 487 } 488 489 async fn list_collections(&self, repo_id: Uuid) -> Result<Vec<Nsid>, DbError> { 490 let rows = sqlx::query_scalar!( 491 "SELECT DISTINCT collection FROM records WHERE repo_id = $1", 492 repo_id 493 ) 494 .fetch_all(&self.pool) 495 .await 496 .map_err(map_sqlx_error)?; 497 498 Ok(rows.into_iter().map(Nsid::from).collect()) 499 } 500 501 async fn count_records(&self, repo_id: Uuid) -> Result<i64, DbError> { 502 let count = sqlx::query_scalar!( 503 r#"SELECT COUNT(*) as "count!" FROM records WHERE repo_id = $1"#, 504 repo_id 505 ) 506 .fetch_one(&self.pool) 507 .await 508 .map_err(map_sqlx_error)?; 509 510 Ok(count) 511 } 512 513 async fn count_all_records(&self) -> Result<i64, DbError> { 514 let count = sqlx::query_scalar!(r#"SELECT COUNT(*) as "count!" FROM records"#) 515 .fetch_one(&self.pool) 516 .await 517 .map_err(map_sqlx_error)?; 518 519 Ok(count) 520 } 521 522 async fn get_record_by_cid( 523 &self, 524 cid: &CidLink, 525 ) -> Result<Option<RecordWithTakedown>, DbError> { 526 let row = sqlx::query!( 527 "SELECT id, takedown_ref FROM records WHERE record_cid = $1", 528 cid.as_str() 529 ) 530 .fetch_optional(&self.pool) 531 .await 532 .map_err(map_sqlx_error)?; 533 534 Ok(row.map(|r| RecordWithTakedown { 535 id: r.id, 536 takedown_ref: r.takedown_ref, 537 })) 538 } 539 540 async fn set_record_takedown( 541 &self, 542 cid: &CidLink, 543 takedown_ref: Option<&str>, 544 ) -> Result<(), DbError> { 545 sqlx::query!( 546 "UPDATE records SET takedown_ref = $1 WHERE record_cid = $2", 547 takedown_ref, 548 cid.as_str() 549 ) 550 .execute(&self.pool) 551 .await 552 .map_err(map_sqlx_error)?; 553 554 Ok(()) 555 } 556 557 async fn insert_user_blocks( 558 &self, 559 user_id: Uuid, 560 block_cids: &[Vec<u8>], 561 repo_rev: &str, 562 ) -> Result<(), DbError> { 563 sqlx::query( 564 r#" 565 INSERT INTO user_blocks (user_id, block_cid, repo_rev) 566 SELECT $1, block_cid, $3 FROM UNNEST($2::bytea[]) AS t(block_cid) 567 ON CONFLICT (user_id, block_cid) DO NOTHING 568 "#, 569 ) 570 .bind(user_id) 571 .bind(block_cids) 572 .bind(repo_rev) 573 .execute(&self.pool) 574 .await 575 .map_err(map_sqlx_error)?; 576 577 Ok(()) 578 } 579 580 async fn delete_user_blocks( 581 &self, 582 user_id: Uuid, 583 block_cids: &[Vec<u8>], 584 ) -> Result<(), DbError> { 585 sqlx::query!( 586 "DELETE FROM user_blocks WHERE user_id = $1 AND block_cid = ANY($2)", 587 user_id, 588 block_cids 589 ) 590 .execute(&self.pool) 591 .await 592 .map_err(map_sqlx_error)?; 593 594 Ok(()) 595 } 596 597 async fn count_user_blocks(&self, user_id: Uuid) -> Result<i64, DbError> { 598 let count = sqlx::query_scalar!( 599 r#"SELECT COUNT(*) as "count!" FROM user_blocks WHERE user_id = $1"#, 600 user_id 601 ) 602 .fetch_one(&self.pool) 603 .await 604 .map_err(map_sqlx_error)?; 605 606 Ok(count) 607 } 608 609 async fn get_user_block_cids_since_rev( 610 &self, 611 user_id: Uuid, 612 since_rev: &str, 613 ) -> Result<Vec<Vec<u8>>, DbError> { 614 let rows: Vec<(Vec<u8>,)> = sqlx::query_as( 615 r#" 616 SELECT block_cid FROM user_blocks 617 WHERE user_id = $1 AND repo_rev > $2 618 ORDER BY repo_rev ASC 619 "#, 620 ) 621 .bind(user_id) 622 .bind(since_rev) 623 .fetch_all(&self.pool) 624 .await 625 .map_err(map_sqlx_error)?; 626 627 Ok(rows.into_iter().map(|(cid,)| cid).collect()) 628 } 629 630 async fn insert_commit_event(&self, data: &CommitEventData) -> Result<SequenceNumber, DbError> { 631 let seq = sqlx::query_scalar!( 632 r#" 633 INSERT INTO repo_seq (did, event_type, commit_cid, prev_cid, ops, blobs, blocks_cids, prev_data_cid, rev) 634 VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9) 635 RETURNING seq 636 "#, 637 data.did.as_str(), 638 data.event_type.as_str(), 639 data.commit_cid.as_ref().map(|c| c.as_str()), 640 data.prev_cid.as_ref().map(|c| c.as_str()), 641 data.ops, 642 data.blobs.as_deref(), 643 data.blocks_cids.as_deref(), 644 data.prev_data_cid.as_ref().map(|c| c.as_str()), 645 data.rev 646 ) 647 .fetch_one(&self.pool) 648 .await 649 .map_err(map_sqlx_error)?; 650 651 Ok(seq.into()) 652 } 653 654 async fn insert_identity_event( 655 &self, 656 did: &Did, 657 handle: Option<&Handle>, 658 ) -> Result<SequenceNumber, DbError> { 659 let handle_str = handle.map(|h| h.as_str()); 660 let seq = sqlx::query_scalar!( 661 r#" 662 INSERT INTO repo_seq (did, event_type, handle) 663 VALUES ($1, 'identity', $2) 664 RETURNING seq 665 "#, 666 did.as_str(), 667 handle_str 668 ) 669 .fetch_one(&self.pool) 670 .await 671 .map_err(map_sqlx_error)?; 672 673 sqlx::query(&format!("NOTIFY repo_updates, '{}'", seq)) 674 .execute(&self.pool) 675 .await 676 .map_err(map_sqlx_error)?; 677 678 Ok(seq.into()) 679 } 680 681 async fn insert_account_event( 682 &self, 683 did: &Did, 684 status: AccountStatus, 685 ) -> Result<SequenceNumber, DbError> { 686 let active = status.is_active(); 687 let status_str = status.for_firehose(); 688 let seq = sqlx::query_scalar!( 689 r#" 690 INSERT INTO repo_seq (did, event_type, active, status) 691 VALUES ($1, 'account', $2, $3) 692 RETURNING seq 693 "#, 694 did.as_str(), 695 active, 696 status_str 697 ) 698 .fetch_one(&self.pool) 699 .await 700 .map_err(map_sqlx_error)?; 701 702 sqlx::query(&format!("NOTIFY repo_updates, '{}'", seq)) 703 .execute(&self.pool) 704 .await 705 .map_err(map_sqlx_error)?; 706 707 Ok(seq.into()) 708 } 709 710 async fn insert_sync_event( 711 &self, 712 did: &Did, 713 commit_cid: &CidLink, 714 rev: Option<&str>, 715 ) -> Result<SequenceNumber, DbError> { 716 let seq = sqlx::query_scalar!( 717 r#" 718 INSERT INTO repo_seq (did, event_type, commit_cid, rev) 719 VALUES ($1, 'sync', $2, $3) 720 RETURNING seq 721 "#, 722 did.as_str(), 723 commit_cid.as_str(), 724 rev 725 ) 726 .fetch_one(&self.pool) 727 .await 728 .map_err(map_sqlx_error)?; 729 730 sqlx::query(&format!("NOTIFY repo_updates, '{}'", seq)) 731 .execute(&self.pool) 732 .await 733 .map_err(map_sqlx_error)?; 734 735 Ok(seq.into()) 736 } 737 738 async fn insert_genesis_commit_event( 739 &self, 740 did: &Did, 741 commit_cid: &CidLink, 742 mst_root_cid: &CidLink, 743 rev: &str, 744 ) -> Result<SequenceNumber, DbError> { 745 let ops = serde_json::json!([]); 746 let blobs: Vec<String> = vec![]; 747 let blocks_cids: Vec<String> = vec![mst_root_cid.to_string(), commit_cid.to_string()]; 748 let prev_cid: Option<&str> = None; 749 750 let seq = sqlx::query_scalar!( 751 r#" 752 INSERT INTO repo_seq (did, event_type, commit_cid, prev_cid, ops, blobs, blocks_cids, rev) 753 VALUES ($1, 'commit', $2, $3::TEXT, $4, $5, $6, $7) 754 RETURNING seq 755 "#, 756 did.as_str(), 757 commit_cid.as_str(), 758 prev_cid, 759 ops, 760 &blobs, 761 &blocks_cids, 762 rev 763 ) 764 .fetch_one(&self.pool) 765 .await 766 .map_err(map_sqlx_error)?; 767 768 sqlx::query(&format!("NOTIFY repo_updates, '{}'", seq)) 769 .execute(&self.pool) 770 .await 771 .map_err(map_sqlx_error)?; 772 773 Ok(seq.into()) 774 } 775 776 async fn update_seq_blocks_cids( 777 &self, 778 seq: SequenceNumber, 779 blocks_cids: &[String], 780 ) -> Result<(), DbError> { 781 sqlx::query!( 782 "UPDATE repo_seq SET blocks_cids = $1 WHERE seq = $2", 783 blocks_cids, 784 seq.as_i64() 785 ) 786 .execute(&self.pool) 787 .await 788 .map_err(map_sqlx_error)?; 789 790 Ok(()) 791 } 792 793 async fn delete_sequences_except( 794 &self, 795 did: &Did, 796 keep_seq: SequenceNumber, 797 ) -> Result<(), DbError> { 798 sqlx::query!( 799 "DELETE FROM repo_seq WHERE did = $1 AND seq != $2", 800 did.as_str(), 801 keep_seq.as_i64() 802 ) 803 .execute(&self.pool) 804 .await 805 .map_err(map_sqlx_error)?; 806 807 Ok(()) 808 } 809 810 async fn get_max_seq(&self) -> Result<SequenceNumber, DbError> { 811 let seq = sqlx::query_scalar!(r#"SELECT COALESCE(MAX(seq), 0) as "max!" FROM repo_seq"#) 812 .fetch_one(&self.pool) 813 .await 814 .map_err(map_sqlx_error)?; 815 816 Ok(seq.into()) 817 } 818 819 async fn get_min_seq_since( 820 &self, 821 since: DateTime<Utc>, 822 ) -> Result<Option<SequenceNumber>, DbError> { 823 let seq = sqlx::query_scalar!( 824 "SELECT MIN(seq) FROM repo_seq WHERE created_at >= $1", 825 since 826 ) 827 .fetch_one(&self.pool) 828 .await 829 .map_err(map_sqlx_error)?; 830 831 Ok(seq.map(SequenceNumber::from)) 832 } 833 834 async fn get_account_with_repo(&self, did: &Did) -> Result<Option<RepoAccountInfo>, DbError> { 835 let row = sqlx::query!( 836 r#"SELECT u.id, u.did, u.deactivated_at, u.takedown_ref, r.repo_root_cid as "repo_root_cid?" 837 FROM users u 838 LEFT JOIN repos r ON r.user_id = u.id 839 WHERE u.did = $1"#, 840 did.as_str() 841 ) 842 .fetch_optional(&self.pool) 843 .await 844 .map_err(map_sqlx_error)?; 845 846 Ok(row.map(|r| RepoAccountInfo { 847 user_id: r.id, 848 did: Did::from(r.did), 849 deactivated_at: r.deactivated_at, 850 takedown_ref: r.takedown_ref, 851 repo_root_cid: r.repo_root_cid.map(CidLink::from), 852 })) 853 } 854 855 async fn get_events_since_seq( 856 &self, 857 since_seq: SequenceNumber, 858 limit: Option<i64>, 859 ) -> Result<Vec<SequencedEvent>, DbError> { 860 let map_row = |r: SequencedEventRow| { 861 let status = r 862 .status 863 .as_deref() 864 .and_then(AccountStatus::parse) 865 .or_else(|| r.active.filter(|a| *a).map(|_| AccountStatus::Active)); 866 SequencedEvent { 867 seq: r.seq.into(), 868 did: Did::from(r.did), 869 created_at: r.created_at, 870 event_type: r.event_type, 871 commit_cid: r.commit_cid.map(CidLink::from), 872 prev_cid: r.prev_cid.map(CidLink::from), 873 prev_data_cid: r.prev_data_cid.map(CidLink::from), 874 ops: r.ops, 875 blobs: r.blobs, 876 blocks_cids: r.blocks_cids, 877 handle: r.handle.map(Handle::from), 878 active: r.active, 879 status, 880 rev: r.rev, 881 } 882 }; 883 match limit { 884 Some(lim) => { 885 let rows = sqlx::query_as!( 886 SequencedEventRow, 887 r#"SELECT seq, did, created_at, event_type as "event_type: RepoEventType", commit_cid, prev_cid, prev_data_cid, 888 ops, blobs, blocks_cids, handle, active, status, rev 889 FROM repo_seq 890 WHERE seq > $1 891 ORDER BY seq ASC 892 LIMIT $2"#, 893 since_seq.as_i64(), 894 lim 895 ) 896 .fetch_all(&self.pool) 897 .await 898 .map_err(map_sqlx_error)?; 899 Ok(rows.into_iter().map(map_row).collect()) 900 } 901 None => { 902 let rows = sqlx::query_as!( 903 SequencedEventRow, 904 r#"SELECT seq, did, created_at, event_type as "event_type: RepoEventType", commit_cid, prev_cid, prev_data_cid, 905 ops, blobs, blocks_cids, handle, active, status, rev 906 FROM repo_seq 907 WHERE seq > $1 908 ORDER BY seq ASC"#, 909 since_seq.as_i64() 910 ) 911 .fetch_all(&self.pool) 912 .await 913 .map_err(map_sqlx_error)?; 914 Ok(rows.into_iter().map(map_row).collect()) 915 } 916 } 917 } 918 919 async fn get_events_in_seq_range( 920 &self, 921 start_seq: SequenceNumber, 922 end_seq: SequenceNumber, 923 ) -> Result<Vec<SequencedEvent>, DbError> { 924 let rows = sqlx::query!( 925 r#"SELECT seq, did, created_at, event_type as "event_type: RepoEventType", commit_cid, prev_cid, prev_data_cid, 926 ops, blobs, blocks_cids, handle, active, status, rev 927 FROM repo_seq 928 WHERE seq > $1 AND seq < $2 929 ORDER BY seq ASC"#, 930 start_seq.as_i64(), 931 end_seq.as_i64() 932 ) 933 .fetch_all(&self.pool) 934 .await 935 .map_err(map_sqlx_error)?; 936 Ok(rows 937 .into_iter() 938 .map(|r| { 939 let status = r 940 .status 941 .as_deref() 942 .and_then(AccountStatus::parse) 943 .or_else(|| r.active.filter(|a| *a).map(|_| AccountStatus::Active)); 944 SequencedEvent { 945 seq: r.seq.into(), 946 did: Did::from(r.did), 947 created_at: r.created_at, 948 event_type: r.event_type, 949 commit_cid: r.commit_cid.map(CidLink::from), 950 prev_cid: r.prev_cid.map(CidLink::from), 951 prev_data_cid: r.prev_data_cid.map(CidLink::from), 952 ops: r.ops, 953 blobs: r.blobs, 954 blocks_cids: r.blocks_cids, 955 handle: r.handle.map(Handle::from), 956 active: r.active, 957 status, 958 rev: r.rev, 959 } 960 }) 961 .collect()) 962 } 963 964 async fn get_event_by_seq( 965 &self, 966 seq: SequenceNumber, 967 ) -> Result<Option<SequencedEvent>, DbError> { 968 let row = sqlx::query!( 969 r#"SELECT seq, did, created_at, event_type as "event_type: RepoEventType", commit_cid, prev_cid, prev_data_cid, 970 ops, blobs, blocks_cids, handle, active, status, rev 971 FROM repo_seq 972 WHERE seq = $1"#, 973 seq.as_i64() 974 ) 975 .fetch_optional(&self.pool) 976 .await 977 .map_err(map_sqlx_error)?; 978 Ok(row.map(|r| { 979 let status = r 980 .status 981 .as_deref() 982 .and_then(AccountStatus::parse) 983 .or_else(|| r.active.filter(|a| *a).map(|_| AccountStatus::Active)); 984 SequencedEvent { 985 seq: r.seq.into(), 986 did: Did::from(r.did), 987 created_at: r.created_at, 988 event_type: r.event_type, 989 commit_cid: r.commit_cid.map(CidLink::from), 990 prev_cid: r.prev_cid.map(CidLink::from), 991 prev_data_cid: r.prev_data_cid.map(CidLink::from), 992 ops: r.ops, 993 blobs: r.blobs, 994 blocks_cids: r.blocks_cids, 995 handle: r.handle.map(Handle::from), 996 active: r.active, 997 status, 998 rev: r.rev, 999 } 1000 })) 1001 } 1002 1003 async fn get_events_since_cursor( 1004 &self, 1005 cursor: SequenceNumber, 1006 limit: i64, 1007 ) -> Result<Vec<SequencedEvent>, DbError> { 1008 let rows = sqlx::query!( 1009 r#"SELECT seq, did, created_at, event_type as "event_type: RepoEventType", commit_cid, prev_cid, prev_data_cid, 1010 ops, blobs, blocks_cids, handle, active, status, rev 1011 FROM repo_seq 1012 WHERE seq > $1 1013 ORDER BY seq ASC 1014 LIMIT $2"#, 1015 cursor.as_i64(), 1016 limit 1017 ) 1018 .fetch_all(&self.pool) 1019 .await 1020 .map_err(map_sqlx_error)?; 1021 Ok(rows 1022 .into_iter() 1023 .map(|r| { 1024 let status = r 1025 .status 1026 .as_deref() 1027 .and_then(AccountStatus::parse) 1028 .or_else(|| r.active.filter(|a| *a).map(|_| AccountStatus::Active)); 1029 SequencedEvent { 1030 seq: r.seq.into(), 1031 did: Did::from(r.did), 1032 created_at: r.created_at, 1033 event_type: r.event_type, 1034 commit_cid: r.commit_cid.map(CidLink::from), 1035 prev_cid: r.prev_cid.map(CidLink::from), 1036 prev_data_cid: r.prev_data_cid.map(CidLink::from), 1037 ops: r.ops, 1038 blobs: r.blobs, 1039 blocks_cids: r.blocks_cids, 1040 handle: r.handle.map(Handle::from), 1041 active: r.active, 1042 status, 1043 rev: r.rev, 1044 } 1045 }) 1046 .collect()) 1047 } 1048 1049 async fn get_events_since_rev( 1050 &self, 1051 did: &Did, 1052 since_rev: &str, 1053 ) -> Result<Vec<EventBlocksCids>, DbError> { 1054 let rows = sqlx::query!( 1055 r#"SELECT blocks_cids, commit_cid 1056 FROM repo_seq 1057 WHERE did = $1 AND rev > $2 1058 ORDER BY seq DESC"#, 1059 did.as_str(), 1060 since_rev 1061 ) 1062 .fetch_all(&self.pool) 1063 .await 1064 .map_err(map_sqlx_error)?; 1065 1066 Ok(rows 1067 .into_iter() 1068 .map(|r| EventBlocksCids { 1069 blocks_cids: r.blocks_cids, 1070 commit_cid: r.commit_cid.map(CidLink::from), 1071 }) 1072 .collect()) 1073 } 1074 1075 async fn list_repos_paginated( 1076 &self, 1077 cursor_did: Option<&Did>, 1078 limit: i64, 1079 ) -> Result<Vec<RepoListItem>, DbError> { 1080 let cursor_str = cursor_did.map(|d| d.as_str()).unwrap_or(""); 1081 let rows = sqlx::query!( 1082 r#"SELECT u.did, u.deactivated_at, u.takedown_ref, r.repo_root_cid, r.repo_rev 1083 FROM repos r 1084 JOIN users u ON r.user_id = u.id 1085 WHERE u.did > $1 1086 ORDER BY u.did ASC 1087 LIMIT $2"#, 1088 cursor_str, 1089 limit 1090 ) 1091 .fetch_all(&self.pool) 1092 .await 1093 .map_err(map_sqlx_error)?; 1094 1095 Ok(rows 1096 .into_iter() 1097 .map(|r| RepoListItem { 1098 did: Did::from(r.did), 1099 deactivated_at: r.deactivated_at, 1100 takedown_ref: r.takedown_ref, 1101 repo_root_cid: CidLink::from(r.repo_root_cid), 1102 repo_rev: r.repo_rev, 1103 }) 1104 .collect()) 1105 } 1106 1107 async fn get_repo_root_cid_by_user_id( 1108 &self, 1109 user_id: Uuid, 1110 ) -> Result<Option<CidLink>, DbError> { 1111 let cid = sqlx::query_scalar!( 1112 "SELECT repo_root_cid FROM repos WHERE user_id = $1", 1113 user_id 1114 ) 1115 .fetch_optional(&self.pool) 1116 .await 1117 .map_err(map_sqlx_error)?; 1118 Ok(cid.map(CidLink::from)) 1119 } 1120 1121 async fn notify_update(&self, seq: SequenceNumber) -> Result<(), DbError> { 1122 sqlx::query(&format!("NOTIFY repo_updates, '{}'", seq.as_i64())) 1123 .execute(&self.pool) 1124 .await 1125 .map_err(map_sqlx_error)?; 1126 Ok(()) 1127 } 1128 1129 async fn import_repo_data( 1130 &self, 1131 user_id: Uuid, 1132 blocks: &[ImportBlock], 1133 records: &[ImportRecord], 1134 ) -> Result<(), ImportRepoError> { 1135 let mut tx = self 1136 .pool 1137 .begin() 1138 .await 1139 .map_err(|e| ImportRepoError::Database(e.to_string()))?; 1140 1141 let repo = sqlx::query!( 1142 "SELECT repo_root_cid FROM repos WHERE user_id = $1 FOR UPDATE NOWAIT", 1143 user_id 1144 ) 1145 .fetch_optional(&mut *tx) 1146 .await 1147 .map_err(|e| { 1148 if let sqlx::Error::Database(ref db_err) = e 1149 && db_err.code().as_deref() == Some("55P03") 1150 { 1151 return ImportRepoError::ConcurrentModification; 1152 } 1153 ImportRepoError::Database(e.to_string()) 1154 })?; 1155 1156 if repo.is_none() { 1157 return Err(ImportRepoError::RepoNotFound); 1158 } 1159 1160 let block_chunks: Vec<Vec<&ImportBlock>> = blocks 1161 .iter() 1162 .collect::<Vec<_>>() 1163 .chunks(100) 1164 .map(|c| c.to_vec()) 1165 .collect(); 1166 1167 for chunk in block_chunks { 1168 for block in chunk { 1169 sqlx::query!( 1170 "INSERT INTO blocks (cid, data) VALUES ($1, $2) ON CONFLICT (cid) DO NOTHING", 1171 &block.cid_bytes, 1172 &block.data 1173 ) 1174 .execute(&mut *tx) 1175 .await 1176 .map_err(|e| ImportRepoError::Database(e.to_string()))?; 1177 } 1178 } 1179 1180 sqlx::query!("DELETE FROM records WHERE repo_id = $1", user_id) 1181 .execute(&mut *tx) 1182 .await 1183 .map_err(|e| ImportRepoError::Database(e.to_string()))?; 1184 1185 for record in records { 1186 sqlx::query!( 1187 r#" 1188 INSERT INTO records (repo_id, collection, rkey, record_cid) 1189 VALUES ($1, $2, $3, $4) 1190 ON CONFLICT (repo_id, collection, rkey) DO UPDATE SET record_cid = $4 1191 "#, 1192 user_id, 1193 record.collection.as_str(), 1194 record.rkey.as_str(), 1195 record.record_cid.as_str() 1196 ) 1197 .execute(&mut *tx) 1198 .await 1199 .map_err(|e| ImportRepoError::Database(e.to_string()))?; 1200 } 1201 1202 tx.commit() 1203 .await 1204 .map_err(|e| ImportRepoError::Database(e.to_string()))?; 1205 1206 Ok(()) 1207 } 1208 1209 async fn apply_commit( 1210 &self, 1211 input: tranquil_db_traits::ApplyCommitInput, 1212 ) -> Result<tranquil_db_traits::ApplyCommitResult, tranquil_db_traits::ApplyCommitError> { 1213 use tranquil_db_traits::ApplyCommitError; 1214 1215 let mut tx = self 1216 .pool 1217 .begin() 1218 .await 1219 .map_err(|e| ApplyCommitError::Database(e.to_string()))?; 1220 1221 let lock_result: Result<Option<_>, sqlx::Error> = sqlx::query!( 1222 "SELECT repo_root_cid FROM repos WHERE user_id = $1 FOR UPDATE NOWAIT", 1223 input.user_id 1224 ) 1225 .fetch_optional(&mut *tx) 1226 .await; 1227 1228 match lock_result { 1229 Err(e) => { 1230 if let Some(db_err) = e.as_database_error() 1231 && db_err.code().as_deref() == Some("55P03") 1232 { 1233 return Err(ApplyCommitError::ConcurrentModification); 1234 } 1235 return Err(ApplyCommitError::Database(format!( 1236 "Failed to acquire repo lock: {}", 1237 e 1238 ))); 1239 } 1240 Ok(Some(row)) => { 1241 if let Some(expected_root) = &input.expected_root_cid 1242 && row.repo_root_cid != expected_root.as_str() 1243 { 1244 return Err(ApplyCommitError::ConcurrentModification); 1245 } 1246 } 1247 Ok(None) => { 1248 return Err(ApplyCommitError::RepoNotFound); 1249 } 1250 } 1251 1252 let is_account_active: bool = 1253 sqlx::query_scalar("SELECT deactivated_at IS NULL FROM users WHERE id = $1") 1254 .bind(input.user_id) 1255 .fetch_optional(&mut *tx) 1256 .await 1257 .map_err(|e| ApplyCommitError::Database(e.to_string()))? 1258 .flatten() 1259 .unwrap_or(false); 1260 1261 sqlx::query("UPDATE repos SET repo_root_cid = $1, repo_rev = $2 WHERE user_id = $3") 1262 .bind(&input.new_root_cid) 1263 .bind(&input.new_rev) 1264 .bind(input.user_id) 1265 .execute(&mut *tx) 1266 .await 1267 .map_err(|e| ApplyCommitError::Database(e.to_string()))?; 1268 1269 if !input.new_block_cids.is_empty() { 1270 sqlx::query( 1271 r#" 1272 INSERT INTO user_blocks (user_id, block_cid, repo_rev) 1273 SELECT $1, block_cid, $3 FROM UNNEST($2::bytea[]) AS t(block_cid) 1274 ON CONFLICT (user_id, block_cid) DO NOTHING 1275 "#, 1276 ) 1277 .bind(input.user_id) 1278 .bind(&input.new_block_cids) 1279 .bind(&input.new_rev) 1280 .execute(&mut *tx) 1281 .await 1282 .map_err(|e| ApplyCommitError::Database(e.to_string()))?; 1283 } 1284 1285 if !input.obsolete_block_cids.is_empty() { 1286 sqlx::query( 1287 r#" 1288 DELETE FROM user_blocks 1289 WHERE user_id = $1 1290 AND block_cid = ANY($2) 1291 "#, 1292 ) 1293 .bind(input.user_id) 1294 .bind(&input.obsolete_block_cids) 1295 .execute(&mut *tx) 1296 .await 1297 .map_err(|e| ApplyCommitError::Database(e.to_string()))?; 1298 } 1299 1300 if !input.record_upserts.is_empty() { 1301 let collections: Vec<&str> = input 1302 .record_upserts 1303 .iter() 1304 .map(|r| r.collection.as_str()) 1305 .collect(); 1306 let rkeys: Vec<&str> = input 1307 .record_upserts 1308 .iter() 1309 .map(|r| r.rkey.as_str()) 1310 .collect(); 1311 let cids: Vec<&str> = input 1312 .record_upserts 1313 .iter() 1314 .map(|r| r.cid.as_str()) 1315 .collect(); 1316 1317 sqlx::query( 1318 r#" 1319 INSERT INTO records (repo_id, collection, rkey, record_cid, repo_rev) 1320 SELECT $1, t.collection, t.rkey, t.cid, $5 1321 FROM UNNEST($2::text[], $3::text[], $4::text[]) AS t(collection, rkey, cid) 1322 ON CONFLICT (repo_id, collection, rkey) DO UPDATE SET record_cid = EXCLUDED.record_cid, repo_rev = EXCLUDED.repo_rev 1323 "#, 1324 ) 1325 .bind(input.user_id) 1326 .bind(&collections) 1327 .bind(&rkeys) 1328 .bind(&cids) 1329 .bind(&input.new_rev) 1330 .execute(&mut *tx) 1331 .await 1332 .map_err(|e| ApplyCommitError::Database(e.to_string()))?; 1333 } 1334 1335 if !input.record_deletes.is_empty() { 1336 let collections: Vec<&str> = input 1337 .record_deletes 1338 .iter() 1339 .map(|r| r.collection.as_str()) 1340 .collect(); 1341 let rkeys: Vec<&str> = input 1342 .record_deletes 1343 .iter() 1344 .map(|r| r.rkey.as_str()) 1345 .collect(); 1346 1347 sqlx::query( 1348 r#" 1349 DELETE FROM records 1350 WHERE repo_id = $1 1351 AND (collection, rkey) IN (SELECT collection, rkey FROM UNNEST($2::text[], $3::text[]) AS t(collection, rkey)) 1352 "#, 1353 ) 1354 .bind(input.user_id) 1355 .bind(&collections) 1356 .bind(&rkeys) 1357 .execute(&mut *tx) 1358 .await 1359 .map_err(|e| ApplyCommitError::Database(e.to_string()))?; 1360 } 1361 1362 let event = &input.commit_event; 1363 let seq: i64 = sqlx::query_scalar( 1364 r#" 1365 INSERT INTO repo_seq (did, event_type, commit_cid, prev_cid, ops, blobs, blocks_cids, prev_data_cid, rev) 1366 VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9) 1367 RETURNING seq 1368 "#, 1369 ) 1370 .bind(&event.did) 1371 .bind(event.event_type.as_str()) 1372 .bind(&event.commit_cid) 1373 .bind(&event.prev_cid) 1374 .bind(&event.ops) 1375 .bind(&event.blobs) 1376 .bind(&event.blocks_cids) 1377 .bind(&event.prev_data_cid) 1378 .bind(&event.rev) 1379 .fetch_one(&mut *tx) 1380 .await 1381 .map_err(|e| ApplyCommitError::Database(e.to_string()))?; 1382 1383 sqlx::query(&format!("NOTIFY repo_updates, '{}'", seq)) 1384 .execute(&mut *tx) 1385 .await 1386 .map_err(|e| ApplyCommitError::Database(e.to_string()))?; 1387 1388 tx.commit() 1389 .await 1390 .map_err(|e| ApplyCommitError::Database(e.to_string()))?; 1391 1392 Ok(tranquil_db_traits::ApplyCommitResult { 1393 seq, 1394 is_account_active, 1395 }) 1396 } 1397 1398 async fn get_broken_genesis_commits( 1399 &self, 1400 ) -> Result<Vec<tranquil_db_traits::BrokenGenesisCommit>, DbError> { 1401 let rows = sqlx::query!( 1402 r#" 1403 SELECT seq, did, commit_cid 1404 FROM repo_seq 1405 WHERE event_type = 'commit' 1406 AND prev_cid IS NULL 1407 AND (blocks_cids IS NULL OR array_length(blocks_cids, 1) IS NULL OR array_length(blocks_cids, 1) = 0) 1408 "# 1409 ) 1410 .fetch_all(&self.pool) 1411 .await 1412 .map_err(map_sqlx_error)?; 1413 1414 Ok(rows 1415 .into_iter() 1416 .map(|r| BrokenGenesisCommit { 1417 seq: r.seq.into(), 1418 did: Did::from(r.did), 1419 commit_cid: r.commit_cid.map(CidLink::from), 1420 }) 1421 .collect()) 1422 } 1423 1424 async fn get_users_without_blocks(&self) -> Result<Vec<UserWithoutBlocks>, DbError> { 1425 let rows: Vec<(Uuid, String, Option<String>)> = sqlx::query_as( 1426 r#" 1427 SELECT u.id as user_id, r.repo_root_cid, r.repo_rev 1428 FROM users u 1429 JOIN repos r ON r.user_id = u.id 1430 WHERE NOT EXISTS (SELECT 1 FROM user_blocks ub WHERE ub.user_id = u.id) 1431 "#, 1432 ) 1433 .fetch_all(&self.pool) 1434 .await 1435 .map_err(map_sqlx_error)?; 1436 1437 Ok(rows 1438 .into_iter() 1439 .map(|(user_id, repo_root_cid, repo_rev)| UserWithoutBlocks { 1440 user_id, 1441 repo_root_cid: CidLink::from(repo_root_cid), 1442 repo_rev, 1443 }) 1444 .collect()) 1445 } 1446 1447 async fn get_users_needing_record_blobs_backfill( 1448 &self, 1449 limit: i64, 1450 ) -> Result<Vec<tranquil_db_traits::UserNeedingRecordBlobsBackfill>, DbError> { 1451 let rows = sqlx::query!( 1452 r#" 1453 SELECT DISTINCT u.id as user_id, u.did 1454 FROM users u 1455 JOIN records r ON r.repo_id = u.id 1456 WHERE NOT EXISTS (SELECT 1 FROM record_blobs rb WHERE rb.repo_id = u.id) 1457 LIMIT $1 1458 "#, 1459 limit 1460 ) 1461 .fetch_all(&self.pool) 1462 .await 1463 .map_err(map_sqlx_error)?; 1464 1465 Ok(rows 1466 .into_iter() 1467 .map(|r| UserNeedingRecordBlobsBackfill { 1468 user_id: r.user_id, 1469 did: Did::from(r.did), 1470 }) 1471 .collect()) 1472 } 1473 1474 async fn insert_record_blobs( 1475 &self, 1476 repo_id: Uuid, 1477 record_uris: &[AtUri], 1478 blob_cids: &[CidLink], 1479 ) -> Result<(), DbError> { 1480 let uris_str: Vec<&str> = record_uris.iter().map(|u| u.as_str()).collect(); 1481 let cids_str: Vec<&str> = blob_cids.iter().map(|c| c.as_str()).collect(); 1482 1483 sqlx::query!( 1484 r#" 1485 INSERT INTO record_blobs (repo_id, record_uri, blob_cid) 1486 SELECT $1, record_uri, blob_cid 1487 FROM UNNEST($2::text[], $3::text[]) AS t(record_uri, blob_cid) 1488 ON CONFLICT (repo_id, record_uri, blob_cid) DO NOTHING 1489 "#, 1490 repo_id, 1491 &uris_str as &[&str], 1492 &cids_str as &[&str] 1493 ) 1494 .execute(&self.pool) 1495 .await 1496 .map_err(map_sqlx_error)?; 1497 1498 Ok(()) 1499 } 1500}