radicle_feed/storage/
postgres.rs1use diesel::dsl::count;
2use diesel::{
3 insert_into, Connection, ExpressionMethods, OptionalExtension, PgConnection, QueryDsl,
4 RunQueryDsl, SelectableHelper,
5};
6use snafu::ResultExt;
7
8use radicle::cob::{ObjectId, TypeName};
9use radicle::git::Oid;
10use radicle::prelude::RepoId;
11
12use crate::models::activity_feed::ActivityFeedOperation;
13use crate::models::entry::{OperationEntry, TimelineEntry};
14use crate::schema::activity_feed_timeline;
15use crate::schema::{activity_feed_operations, seeded_radicle_repository};
16use crate::storage::{FeedStorage, StorageStats};
17
18pub struct PostgresStorage {
19 connection: PgConnection,
20}
21
22impl PostgresStorage {
23 pub fn new(db_url: String) -> Result<Self, snafu::Whatever> {
24 let conn =
25 PgConnection::establish(&db_url).whatever_context("Unable to establish connection")?;
26
27 Ok(Self { connection: conn })
28 }
29}
30
31impl FeedStorage for PostgresStorage {
32 type Error = snafu::Whatever;
33
34 fn get_last_processed_operation(
35 &mut self,
36 rid: &RepoId,
37 cob_id: &ObjectId,
38 typename: &TypeName,
39 ) -> Result<Option<String>, Self::Error> {
40 let rows = activity_feed_timeline::table
41 .filter(activity_feed_timeline::repo.eq(rid.to_string()))
42 .filter(activity_feed_timeline::cob_id.eq(cob_id.to_string()))
43 .filter(activity_feed_timeline::typename.eq(typename.to_string()))
44 .select(activity_feed_timeline::last_operation_id)
45 .load::<Option<String>>(&mut self.connection)
46 .whatever_context("Unable to get_last_processed_operation")?;
47
48 if rows.is_empty() {
49 Ok(None)
50 } else {
51 Ok(rows[0].clone())
52 }
53 }
54
55 fn get_operation_by_id(
56 &mut self,
57 id: &radicle::git::Oid,
58 ) -> Result<ActivityFeedOperation, Self::Error> {
59 activity_feed_operations::table
60 .filter(activity_feed_operations::operation_id.eq(id.to_string()))
61 .select(ActivityFeedOperation::as_select())
62 .first::<ActivityFeedOperation>(&mut self.connection)
63 .whatever_context("Unable to get_operation_by_id")
64 }
65
66 fn operation_exists(&mut self, operation_id: &Oid) -> Result<bool, Self::Error> {
67 let rows = activity_feed_operations::table
68 .filter(activity_feed_operations::operation_id.eq(operation_id.to_string()))
69 .select(activity_feed_operations::operation_id)
70 .first::<String>(&mut self.connection)
71 .optional()
72 .whatever_context("Unable to load last feed operation")?;
73
74 Ok(rows.is_some())
75 }
76
77 fn insert_timeline_entry(&mut self, entry: &TimelineEntry) -> Result<(), Self::Error> {
79 insert_into(activity_feed_timeline::table)
80 .values(entry)
81 .on_conflict((
82 activity_feed_timeline::repo,
83 activity_feed_timeline::node,
84 activity_feed_timeline::cob_id,
85 ))
86 .do_update()
87 .set(activity_feed_timeline::last_operation_id.eq(entry.last_operation_id.clone()))
88 .execute(&mut self.connection)
89 .whatever_context("Unable to insert")?;
90
91 Ok(())
92 }
93
94 fn resolve_rid(&mut self, rid: &RepoId) -> Result<Option<String>, Self::Error> {
95 seeded_radicle_repository::table
96 .filter(seeded_radicle_repository::repository_id.eq(rid.to_string()))
97 .select(seeded_radicle_repository::alias)
98 .first::<Option<String>>(&mut self.connection)
99 .whatever_context("Unable to resolve rid")
100 }
101
102 fn insert_batch(&mut self, entries: &[OperationEntry]) -> Result<(), Self::Error> {
103 if entries.is_empty() {
104 return Ok(());
105 }
106
107 insert_into(activity_feed_operations::table)
108 .values(entries)
109 .execute(&mut self.connection)
110 .whatever_context("context")?;
111
112 Ok(())
113 }
114
115 fn get_stats(&mut self) -> Result<StorageStats, Self::Error> {
116 let mut stats = StorageStats::default();
117
118 let statement_total_operations = activity_feed_operations::table
119 .count()
120 .first::<i64>(&mut self.connection)
121 .whatever_context("context")?;
122 stats.total_operations = statement_total_operations as u64;
123
124 let statement_operations_by_type = activity_feed_operations::table
125 .count()
126 .group_by(activity_feed_operations::typename)
127 .select((
128 activity_feed_operations::typename,
129 count(activity_feed_operations::id),
130 ))
131 .load::<(String, i64)>(&mut self.connection)
132 .whatever_context("context")?;
133
134 for (typename, count) in statement_operations_by_type {
135 stats.operations_by_type.insert(typename, count as u64);
136 }
137
138 let statement_tracked_objects = activity_feed_operations::table
139 .count()
140 .get_results::<i64>(&mut self.connection)
141 .whatever_context("context")?;
142
143 for count in statement_tracked_objects {
144 stats.tracked_objects = count as u64;
145 }
146
147 Ok(stats)
148 }
149}