forked from
tranquil.farm/tranquil-pds
Our Personal Data Server from scratch!
48 kB
1500 lines
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}