profiles.rs (10555B)
1 use harvestcircle_application::{ 2 BoxFuture, CachedProfile, ProfileRefreshStatus, ProfileRepository, 3 }; 4 use harvestcircle_domain::{ 5 EventId, Kind0ProfileCandidate, ProfileMetadata, PublicKey, SafeError, UnixTimestamp, 6 }; 7 use sqlx::Row; 8 9 use crate::Database; 10 use crate::db::{corrupt_storage, map_transaction_error, storage_unavailable}; 11 12 const PROFILE_PROJECTION: &str = "SELECT substr(event_id, 1, 33) AS event_id, \ 13 length(event_id) AS event_id_bytes, event_created_at_unix_s, \ 14 CASE WHEN name IS NULL THEN NULL ELSE substr(CAST(name AS BLOB), 1, 129) END AS name, \ 15 CASE WHEN name IS NULL THEN NULL ELSE length(CAST(name AS BLOB)) END AS name_bytes, \ 16 CASE WHEN display_name IS NULL THEN NULL ELSE substr(CAST(display_name AS BLOB), 1, 129) END AS display_name, \ 17 CASE WHEN display_name IS NULL THEN NULL ELSE length(CAST(display_name AS BLOB)) END AS display_name_bytes, \ 18 CASE WHEN nip05 IS NULL THEN NULL ELSE substr(CAST(nip05 AS BLOB), 1, 321) END AS nip05, \ 19 CASE WHEN nip05 IS NULL THEN NULL ELSE length(CAST(nip05 AS BLOB)) END AS nip05_bytes, \ 20 CASE WHEN about IS NULL THEN NULL ELSE substr(CAST(about AS BLOB), 1, 4097) END AS about, \ 21 CASE WHEN about IS NULL THEN NULL ELSE length(CAST(about AS BLOB)) END AS about_bytes, \ 22 CASE WHEN picture IS NULL THEN NULL ELSE substr(CAST(picture AS BLOB), 1, 2049) END AS picture, \ 23 CASE WHEN picture IS NULL THEN NULL ELSE length(CAST(picture AS BLOB)) END AS picture_bytes, \ 24 refreshed_at_unix_s, substr(CAST(refresh_status AS BLOB), 1, 13) AS refresh_status, \ 25 length(CAST(refresh_status AS BLOB)) AS refresh_status_bytes FROM profile_cache"; 26 27 impl ProfileRepository for Database { 28 fn load_profile( 29 &self, 30 public_key: PublicKey, 31 ) -> BoxFuture<'_, Result<Option<CachedProfile>, SafeError>> { 32 Box::pin(async move { 33 self.host() 34 .transaction(|transaction| { 35 Box::pin(async move { 36 let sql = 37 format!("{PROFILE_PROJECTION} WHERE subject_public_key = ? LIMIT 2"); 38 let rows = sqlx::query(sqlx::AssertSqlSafe(sql.as_str())) 39 .bind(public_key.as_bytes().as_slice()) 40 .fetch_all(&mut *transaction) 41 .await 42 .map_err(|_| corrupt_storage())?; 43 match rows.as_slice() { 44 [] => Ok(None), 45 [row] => decode_profile(row, public_key).map(Some), 46 _ => Err(corrupt_storage()), 47 } 48 }) 49 }) 50 .await 51 .map_err(map_transaction_error) 52 }) 53 } 54 55 fn save_profile<'a>( 56 &'a self, 57 profile: &'a CachedProfile, 58 ) -> BoxFuture<'a, Result<(), SafeError>> { 59 Box::pin(async move { 60 let candidate = profile.candidate(); 61 let metadata = candidate.metadata(); 62 let author = *candidate.author().as_bytes(); 63 let event_id = candidate.event_id().as_bytes(); 64 let event_created_at = candidate.created_at().as_seconds(); 65 let name = metadata.name().map(str::to_owned); 66 let display_name = metadata.display_name().map(str::to_owned); 67 let nip05 = metadata.nip05().map(str::to_owned); 68 let about = metadata.about().map(str::to_owned); 69 let picture = metadata.picture().map(str::to_owned); 70 let refreshed_at = profile.refreshed_at().as_seconds(); 71 let refresh_status = encode_refresh_status(profile.refresh_status()); 72 self.host() 73 .transaction(|transaction| { 74 Box::pin(async move { 75 sqlx::query( 76 "INSERT INTO profile_cache (subject_public_key, event_id, \ 77 event_created_at_unix_s, name, display_name, nip05, about, picture, \ 78 refreshed_at_unix_s, refresh_status) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?) \ 79 ON CONFLICT(subject_public_key) DO UPDATE SET event_id = excluded.event_id, \ 80 event_created_at_unix_s = excluded.event_created_at_unix_s, name = excluded.name, \ 81 display_name = excluded.display_name, nip05 = excluded.nip05, about = excluded.about, \ 82 picture = excluded.picture, refreshed_at_unix_s = excluded.refreshed_at_unix_s, \ 83 refresh_status = excluded.refresh_status WHERE \ 84 excluded.event_created_at_unix_s > profile_cache.event_created_at_unix_s \ 85 OR (excluded.event_created_at_unix_s = profile_cache.event_created_at_unix_s \ 86 AND excluded.event_id < profile_cache.event_id)", 87 ) 88 .bind(author.as_slice()) 89 .bind(event_id.as_slice()) 90 .bind(event_created_at) 91 .bind(&name) 92 .bind(&display_name) 93 .bind(&nip05) 94 .bind(&about) 95 .bind(&picture) 96 .bind(refreshed_at) 97 .bind(refresh_status) 98 .execute(&mut *transaction) 99 .await 100 .map(|_| ()) 101 .map_err(|_| storage_unavailable()) 102 }) 103 }) 104 .await 105 .map_err(map_transaction_error) 106 }) 107 } 108 109 fn record_refresh_status<'a>( 110 &'a self, 111 public_key: PublicKey, 112 refreshed_at: UnixTimestamp, 113 status: ProfileRefreshStatus, 114 ) -> BoxFuture<'a, Result<(), SafeError>> { 115 Box::pin(async move { 116 self.host() 117 .transaction(|transaction| { 118 Box::pin(async move { 119 sqlx::query( 120 "UPDATE profile_cache SET refreshed_at_unix_s = ?, refresh_status = ? \ 121 WHERE subject_public_key = ?", 122 ) 123 .bind(refreshed_at.as_seconds()) 124 .bind(encode_refresh_status(status)) 125 .bind(public_key.as_bytes().as_slice()) 126 .execute(&mut *transaction) 127 .await 128 .map(|_| ()) 129 .map_err(|_| storage_unavailable()) 130 }) 131 }) 132 .await 133 .map_err(map_transaction_error) 134 }) 135 } 136 137 fn remove_profile(&self, public_key: PublicKey) -> BoxFuture<'_, Result<(), SafeError>> { 138 Box::pin(async move { 139 self.host() 140 .transaction(|transaction| { 141 Box::pin(async move { 142 sqlx::query("DELETE FROM profile_cache WHERE subject_public_key = ?") 143 .bind(public_key.as_bytes().as_slice()) 144 .execute(&mut *transaction) 145 .await 146 .map(|_| ()) 147 .map_err(|_| storage_unavailable()) 148 }) 149 }) 150 .await 151 .map_err(map_transaction_error) 152 }) 153 } 154 } 155 156 fn decode_profile( 157 row: &sqlx::sqlite::SqliteRow, 158 author: PublicKey, 159 ) -> Result<CachedProfile, SafeError> { 160 let event_id = 161 bounded_blob(row, "event_id", "event_id_bytes", 32)?.ok_or_else(corrupt_storage)?; 162 let event_id: [u8; 32] = event_id.try_into().map_err(|_| corrupt_storage())?; 163 let created_at = UnixTimestamp::from_seconds( 164 row.try_get("event_created_at_unix_s") 165 .map_err(|_| corrupt_storage())?, 166 ) 167 .ok_or_else(corrupt_storage)?; 168 let metadata = ProfileMetadata::new( 169 bounded_text(row, "name", "name_bytes", 128)?, 170 bounded_text(row, "display_name", "display_name_bytes", 128)?, 171 bounded_text(row, "nip05", "nip05_bytes", 320)?, 172 bounded_text(row, "about", "about_bytes", 4_096)?, 173 bounded_text(row, "picture", "picture_bytes", 2_048)?, 174 ) 175 .map_err(|_| corrupt_storage())?; 176 let refreshed_at = UnixTimestamp::from_seconds( 177 row.try_get("refreshed_at_unix_s") 178 .map_err(|_| corrupt_storage())?, 179 ) 180 .ok_or_else(corrupt_storage)?; 181 let status = bounded_text(row, "refresh_status", "refresh_status_bytes", 12)? 182 .ok_or_else(corrupt_storage)?; 183 Ok(CachedProfile::new( 184 Kind0ProfileCandidate::new(EventId::from_bytes(event_id), author, created_at, metadata), 185 refreshed_at, 186 decode_refresh_status(&status)?, 187 )) 188 } 189 190 fn bounded_text( 191 row: &sqlx::sqlite::SqliteRow, 192 value_column: &str, 193 length_column: &str, 194 maximum: usize, 195 ) -> Result<Option<String>, SafeError> { 196 bounded_blob(row, value_column, length_column, maximum)? 197 .map(String::from_utf8) 198 .transpose() 199 .map_err(|_| corrupt_storage()) 200 } 201 202 fn bounded_blob( 203 row: &sqlx::sqlite::SqliteRow, 204 value_column: &str, 205 length_column: &str, 206 maximum: usize, 207 ) -> Result<Option<Vec<u8>>, SafeError> { 208 let value = row 209 .try_get::<Option<Vec<u8>>, _>(value_column) 210 .map_err(|_| corrupt_storage())?; 211 let length = row 212 .try_get::<Option<i64>, _>(length_column) 213 .map_err(|_| corrupt_storage())?; 214 match (value, length) { 215 (None, None) => Ok(None), 216 (Some(value), Some(length)) 217 if usize::try_from(length) 218 .ok() 219 .is_some_and(|length| length <= maximum && length == value.len()) => 220 { 221 Ok(Some(value)) 222 } 223 _ => Err(corrupt_storage()), 224 } 225 } 226 227 const fn encode_refresh_status(status: ProfileRefreshStatus) -> &'static str { 228 match status { 229 ProfileRefreshStatus::Success => "success", 230 ProfileRefreshStatus::Offline => "offline", 231 ProfileRefreshStatus::InvalidData => "invalid_data", 232 } 233 } 234 235 fn decode_refresh_status(value: &str) -> Result<ProfileRefreshStatus, SafeError> { 236 match value { 237 "success" => Ok(ProfileRefreshStatus::Success), 238 "offline" => Ok(ProfileRefreshStatus::Offline), 239 "invalid_data" => Ok(ProfileRefreshStatus::InvalidData), 240 _ => Err(corrupt_storage()), 241 } 242 }