Something went wrong. Try again.
Learn how to use Rust to build ATProto powered applications
Something went wrong. Try again.
14 kB · 408 lines
Rust
at main
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409use actix_web::web::Data;use async_sqlite::{ Pool, rusqlite, rusqlite::{Error, Row},};use atrium_api::types::string::Did;use chrono::{DateTime, Datelike, Utc};use rusqlite::types::Type;use serde::{Deserialize, Serialize};use std::{fmt::Debug, sync::Arc};
/// Creates the tables in the db.pub async fn create_tables_in_database(pool: &Pool) -> Result<(), async_sqlite::Error> { pool.conn(move |conn| { conn.execute("PRAGMA foreign_keys = ON", []).unwrap();
// status conn.execute( "CREATE TABLE IF NOT EXISTS status ( uri TEXT PRIMARY KEY, authorDid TEXT NOT NULL, status TEXT NOT NULL, createdAt INTEGER NOT NULL, indexedAt INTEGER NOT NULL )", [], ) .unwrap();
// auth_session conn.execute( "CREATE TABLE IF NOT EXISTS auth_session ( key TEXT PRIMARY KEY, session TEXT NOT NULL )", [], ) .unwrap();
// auth_state conn.execute( "CREATE TABLE IF NOT EXISTS auth_state ( key TEXT PRIMARY KEY, state TEXT NOT NULL )", [], ) .unwrap(); Ok(()) }) .await?; Ok(())}
///Status table datatype#[derive(Debug, Clone, Deserialize, Serialize)]pub struct StatusFromDb { pub uri: String, pub author_did: String, pub status: String, pub created_at: DateTime<Utc>, pub indexed_at: DateTime<Utc>, pub handle: Option<String>,}
//Status methodsimpl StatusFromDb { /// Creates a new [StatusFromDb] pub fn new(uri: String, author_did: String, status: String) -> Self { let now = chrono::Utc::now(); Self { uri, author_did, status, created_at: now, indexed_at: now, handle: None, } }
/// Helper to map from [Row] to [StatusDb] fn map_from_row(row: &Row) -> Result<Self, rusqlite::Error> { Ok(Self { uri: row.get(0)?, author_did: row.get(1)?, status: row.get(2)?, //DateTimes are stored as INTEGERS then parsed into a DateTime<UTC> created_at: { let timestamp: i64 = row.get(3)?; DateTime::from_timestamp(timestamp, 0).ok_or_else(|| { Error::InvalidColumnType(3, "Invalid timestamp".parse().unwrap(), Type::Text) })? }, //DateTimes are stored as INTEGERS then parsed into a DateTime<UTC> indexed_at: { let timestamp: i64 = row.get(4)?; DateTime::from_timestamp(timestamp, 0).ok_or_else(|| { Error::InvalidColumnType(4, "Invalid timestamp".parse().unwrap(), Type::Text) })? }, handle: None, }) }
/// Helper for the UI to see if indexed_at date is today or not pub fn is_today(&self) -> bool { let now = Utc::now();
self.indexed_at.day() == now.day() && self.indexed_at.month() == now.month() && self.indexed_at.year() == now.year() }
/// Saves the [StatusDb] pub async fn save(&self, pool: Data<Arc<Pool>>) -> Result<(), async_sqlite::Error> { let cloned_self = self.clone(); pool.conn(move |conn| { Ok(conn.execute( "INSERT INTO status (uri, authorDid, status, createdAt, indexedAt) VALUES (?1, ?2, ?3, ?4, ?5)", [ &cloned_self.uri, &cloned_self.author_did, &cloned_self.status, &cloned_self.created_at.timestamp().to_string(), &cloned_self.indexed_at.timestamp().to_string(), ], )?) }) .await?; Ok(()) }
/// Saves or updates a status by its did(uri) pub async fn save_or_update(&self, pool: &Pool) -> Result<(), async_sqlite::Error> { let cloned_self = self.clone(); pool.conn(move |conn| { //We check to see if the session already exists, if so we need to update not insert let mut stmt = conn.prepare("SELECT COUNT(*) FROM status WHERE uri = ?1")?; let count: i64 = stmt.query_row([&cloned_self.uri], |row| row.get(0))?; match count > 0 { true => { let mut update_stmt = conn.prepare("UPDATE status SET status = ?2, indexedAt = ? WHERE uri = ?1")?; update_stmt.execute([&cloned_self.uri, &cloned_self.status, &cloned_self.indexed_at.timestamp().to_string()])?; Ok(()) } false => { conn.execute( "INSERT INTO status (uri, authorDid, status, createdAt, indexedAt) VALUES (?1, ?2, ?3, ?4, ?5)", [ &cloned_self.uri, &cloned_self.author_did, &cloned_self.status, &cloned_self.created_at.timestamp().to_string(), &cloned_self.indexed_at.timestamp().to_string(), ], )?; Ok(()) } } }) .await?; Ok(()) } pub async fn delete_by_uri(pool: &Pool, uri: String) -> Result<(), async_sqlite::Error> { pool.conn(move |conn| { let mut stmt = conn.prepare("DELETE FROM status WHERE uri = ?1")?; stmt.execute([&uri]) }) .await?; Ok(()) }
/// Loads the last 10 statuses we have saved pub async fn load_latest_statuses( pool: &Data<Arc<Pool>>, ) -> Result<Vec<Self>, async_sqlite::Error> { Ok(pool .conn(move |conn| { let mut stmt = conn.prepare("SELECT * FROM status ORDER BY indexedAt DESC LIMIT 10")?; let status_iter = stmt .query_map([], |row| Ok(Self::map_from_row(row).unwrap())) .unwrap();
let mut statuses = Vec::new(); for status in status_iter { statuses.push(status?); } Ok(statuses) }) .await?) }
/// Loads the logged-in users current status pub async fn my_status( pool: &Data<Arc<Pool>>, did: &str, ) -> Result<Option<Self>, async_sqlite::Error> { let did = did.to_string(); pool.conn(move |conn| { let mut stmt = conn.prepare( "SELECT * FROM status WHERE authorDid = ?1 ORDER BY createdAt DESC LIMIT 1", )?; stmt.query_row([did.as_str()], |row| Self::map_from_row(row)) .map(Some) .or_else(|err| { if err == rusqlite::Error::QueryReturnedNoRows { Ok(None) } else { Err(err) } }) }) .await }
/// ui helper to show a handle or did if the handle cannot be found pub fn author_display_name(&self) -> String { match self.handle.as_ref() { Some(handle) => handle.to_string(), None => self.author_did.to_string(), } }}
/// AuthSession table data type#[derive(Debug, Clone, Deserialize, Serialize)]pub struct AuthSession { pub key: String, pub session: String,}
impl AuthSession { /// Creates a new [AuthSession] pub fn new<V>(key: String, session: V) -> Self where V: Serialize, { let session = serde_json::to_string(&session).unwrap(); Self { key: key.to_string(), session, } }
/// Helper to map from [Row] to [AuthSession] fn map_from_row(row: &Row) -> Result<Self, Error> { let key: String = row.get(0)?; let session: String = row.get(1)?; Ok(Self { key, session }) }
/// Gets a session by the users did(key) pub async fn get_by_did(pool: &Pool, did: String) -> Result<Option<Self>, async_sqlite::Error> { let did = Did::new(did).unwrap(); pool.conn(move |conn| { let mut stmt = conn.prepare("SELECT * FROM auth_session WHERE key = ?1")?; stmt.query_row([did.as_str()], |row| Self::map_from_row(row)) .map(Some) .or_else(|err| { if err == Error::QueryReturnedNoRows { Ok(None) } else { Err(err) } }) }) .await }
/// Saves or updates the session by its did(key) pub async fn save_or_update(&self, pool: &Pool) -> Result<(), async_sqlite::Error> { let cloned_self = self.clone(); pool.conn(move |conn| { //We check to see if the session already exists, if so we need to update not insert let mut stmt = conn.prepare("SELECT COUNT(*) FROM auth_session WHERE key = ?1")?; let count: i64 = stmt.query_row([&cloned_self.key], |row| row.get(0))?; match count > 0 { true => { let mut update_stmt = conn.prepare("UPDATE auth_session SET session = ?2 WHERE key = ?1")?; update_stmt.execute([&cloned_self.key, &cloned_self.session])?; Ok(()) } false => { conn.execute( "INSERT INTO auth_session (key, session) VALUES (?1, ?2)", [&cloned_self.key, &cloned_self.session], )?; Ok(()) } } }) .await?; Ok(()) }
/// Deletes the session by did pub async fn delete_by_did(pool: &Pool, did: String) -> Result<(), async_sqlite::Error> { pool.conn(move |conn| { let mut stmt = conn.prepare("DELETE FROM auth_session WHERE key = ?1")?; stmt.execute([&did]) }) .await?; Ok(()) }
/// Deletes all the sessions pub async fn delete_all(pool: &Pool) -> Result<(), async_sqlite::Error> { pool.conn(move |conn| { let mut stmt = conn.prepare("DELETE FROM auth_session")?; stmt.execute([]) }) .await?; Ok(()) }}
/// AuthState table datatype#[derive(Debug, Clone, Deserialize, Serialize)]pub struct AuthState { pub key: String, pub state: String,}
impl AuthState { /// Creates a new [AuthState] pub fn new<V>(key: String, state: V) -> Self where V: Serialize, { let state = serde_json::to_string(&state).unwrap(); Self { key: key.to_string(), state, } }
/// Helper to map from [Row] to [AuthState] fn map_from_row(row: &Row) -> Result<Self, Error> { let key: String = row.get(0)?; let state: String = row.get(1)?; Ok(Self { key, state }) }
/// Gets a state by the users key pub async fn get_by_key(pool: &Pool, key: String) -> Result<Option<Self>, async_sqlite::Error> { pool.conn(move |conn| { let mut stmt = conn.prepare("SELECT * FROM auth_state WHERE key = ?1")?; stmt.query_row([key.as_str()], |row| Self::map_from_row(row)) .map(Some) .or_else(|err| { if err == Error::QueryReturnedNoRows { Ok(None) } else { Err(err) } }) }) .await }
/// Saves or updates the state by its key pub async fn save_or_update(&self, pool: &Pool) -> Result<(), async_sqlite::Error> { let cloned_self = self.clone(); pool.conn(move |conn| { //We check to see if the state already exists, if so we need to update let mut stmt = conn.prepare("SELECT COUNT(*) FROM auth_state WHERE key = ?1")?; let count: i64 = stmt.query_row([&cloned_self.key], |row| row.get(0))?; match count > 0 { true => { let mut update_stmt = conn.prepare("UPDATE auth_state SET state = ?2 WHERE key = ?1")?; update_stmt.execute([&cloned_self.key, &cloned_self.state])?; Ok(()) } false => { conn.execute( "INSERT INTO auth_state (key, state) VALUES (?1, ?2)", [&cloned_self.key, &cloned_self.state], )?; Ok(()) } } }) .await?; Ok(()) }
pub async fn delete_by_key(pool: &Pool, key: String) -> Result<(), async_sqlite::Error> { pool.conn(move |conn| { let mut stmt = conn.prepare("DELETE FROM auth_state WHERE key = ?1")?; stmt.execute([&key]) }) .await?; Ok(()) }
pub async fn delete_all(pool: &Pool) -> Result<(), async_sqlite::Error> { pool.conn(move |conn| { let mut stmt = conn.prepare("DELETE FROM auth_state")?; stmt.execute([]) }) .await?; Ok(()) }}