app

Local-first trade for farms and co-ops
git clone https://radroots.dev/git/app.git
Log | Files | Refs | README | LICENSE

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 }