Skip to main content

radicle_feed/storage/
postgres.rs

1use 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    /// Inserts a new timeline entry into the database, in case of conflict it tries to update the last operation ID.
78    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}