| @@ -8,5 +8,12 @@ | |||
| <jdbc-url>jdbc:postgresql://localhost:5432/postgres</jdbc-url> | |||
| <working-dir>$ProjectFileDir$</working-dir> | |||
| </data-source> | |||
| <data-source source="LOCAL" name="postgres@localhost [2]" uuid="887ad5d3-a663-40c3-8bf7-5f0439ede4f9"> | |||
| <driver-ref>postgresql</driver-ref> | |||
| <synchronize>true</synchronize> | |||
| <jdbc-driver>org.postgresql.Driver</jdbc-driver> | |||
| <jdbc-url>jdbc:postgresql://localhost:5432/postgres</jdbc-url> | |||
| <working-dir>$ProjectFileDir$</working-dir> | |||
| </data-source> | |||
| </component> | |||
| </project> | |||
| @@ -7,15 +7,14 @@ edition = "2024" | |||
| actix-web = "4.11.0" | |||
| serde = { version = "1.0.219", features = ["derive"] } | |||
| serde_json = "1.0.142" | |||
| ipfs-api-backend-actix = "0.7.0" | |||
| uuid = { version = "1.17.0", features = ["v4"] } | |||
| diesel = { version = "2.2.12", features = ["postgres", "chrono", "r2d2"] } | |||
| dotenvy = "0.15.7" | |||
| chrono = { version = "0.4.41", features = ["serde"] } | |||
| futures = "0.3.31" | |||
| diesel-derive-enum = { version = "3.0.0-beta.1", features = ["postgres"] } | |||
| argon2 = "0.5" | |||
| base64 = "0.22" | |||
| ed25519-dalek = { version = "2.1", features = ["rand_core"] } | |||
| stellar-strkey = "0.0.13" | |||
| jsonwebtoken = "9" | |||
| argon2 = "0.6.0" | |||
| base64 = "0.23.1" | |||
| ed25519-dalek = { version = "3.0.0", features = ["rand_core"] } | |||
| stellar-strkey = "0.0.18" | |||
| jsonwebtoken = { version = "11.1.0", features = ["rust_crypto"] } | |||
| @@ -6,4 +6,4 @@ file = "src/schema.rs" | |||
| custom_type_derives = ["diesel::query_builder::QueryId", "Clone"] | |||
| [migrations_directory] | |||
| dir = "/Users/jaredbell/RustroverProjects/puffpastry-backend/migrations" | |||
| dir = "/home/jbell/RustroverProjects/puffpastry-wheat-backend/migrations" | |||
| @@ -1,11 +1,7 @@ | |||
| create table proposals ( | |||
| id serial primary key, | |||
| name text not null, | |||
| cid text not null, | |||
| summary text, | |||
| creator text not null, | |||
| is_current boolean not null, -- true for latest/current version | |||
| previous_cid text, -- the CID of the previous version, if any | |||
| created_at timestamp default now(), | |||
| created_at timestamp not null default now(), | |||
| updated_at timestamp | |||
| ) | |||
| @@ -8,12 +8,9 @@ create type amendment_status as enum ( | |||
| create table amendments ( | |||
| id serial primary key, | |||
| name text not null, | |||
| cid text not null, | |||
| summary text, | |||
| status amendment_status not null default 'proposed', | |||
| creator text not null, | |||
| is_current boolean not null default true, | |||
| previous_cid text, | |||
| created_at timestamptz not null default now(), | |||
| updated_at timestamptz not null default now(), | |||
| proposal_id integer not null references proposals(id) | |||
| @@ -3,10 +3,10 @@ | |||
| create table comments ( | |||
| id serial primary key, | |||
| cid text not null, | |||
| content text, | |||
| is_current bool not null default true, | |||
| previous_cid text, | |||
| proposal_id integer not null references proposals(id) on delete cascade, | |||
| created_at timestamptz not null default now(), | |||
| updated_at timestamptz not null default now() | |||
| updated_at timestamptz not null default now(), | |||
| user_id integer references users(id) on delete set null | |||
| ); | |||
| @@ -0,0 +1 @@ | |||
| -- This file should undo anything in `up.sql` | |||
| @@ -0,0 +1,6 @@ | |||
| create table proposal_paragraphs ( | |||
| id serial primary key, | |||
| paragraph_text text not null, | |||
| index integer not null, | |||
| proposal_id integer not null references proposals(id) | |||
| ) | |||
| @@ -1,82 +0,0 @@ | |||
| use crate::db::db::get_connection; | |||
| use crate::ipfs::ipfs::IpfsService; | |||
| use crate::schema::amendments; | |||
| use crate::types::amendment::{Amendment, AmendmentError, AmendmentFile}; | |||
| use crate::types::ipfs::IpfsResult; | |||
| use crate::utils::content::extract_summary; | |||
| use crate::utils::ipfs::{read_json_via_cat, upload_json_and_get_hash, DEFAULT_MAX_JSON_SIZE}; | |||
| use diesel::ExpressionMethods; | |||
| use diesel::QueryDsl; | |||
| use diesel::RunQueryDsl; | |||
| use ipfs_api_backend_actix::IpfsClient; | |||
| const STORAGE_DIR: &str = "/puffpastry/amendments"; | |||
| const FILE_EXTENSION: &str = "json"; | |||
| pub struct AmendmentService { | |||
| client: IpfsClient, | |||
| } | |||
| impl IpfsService<AmendmentFile> for AmendmentService { | |||
| type Err = AmendmentError; | |||
| async fn save(&mut self, item: AmendmentFile) -> IpfsResult<String, Self::Err> { | |||
| self.store_amendment_to_ipfs(item).await | |||
| } | |||
| async fn read(&mut self, hash: String) -> IpfsResult<AmendmentFile, Self::Err> { | |||
| self.read_amendment_file(&hash).await | |||
| } | |||
| } | |||
| impl AmendmentService { | |||
| pub fn new(client: IpfsClient) -> Self { | |||
| Self { client } | |||
| } | |||
| async fn store_amendment_to_ipfs(&self, amendment: AmendmentFile) -> IpfsResult<String, AmendmentError> { | |||
| let amendment_record = Amendment::new( | |||
| amendment.name.clone(), | |||
| Some(extract_summary(&amendment.content)), | |||
| amendment.creator.clone(), | |||
| amendment.proposal_id, | |||
| ); | |||
| let hash = self.upload_to_ipfs(&amendment).await?; | |||
| self.save_to_database(amendment_record.with_cid(hash.clone()).mark_as_current()).await?; | |||
| { | |||
| use crate::schema::amendments::dsl::{amendments as amendments_table, is_current as is_current_col, previous_cid as previous_cid_col}; | |||
| let mut conn = get_connection() | |||
| .map_err(|e| AmendmentError::DatabaseError(e))?; | |||
| // Ignore the count result; if no rows match, that's fine | |||
| let _ = diesel::update(amendments_table.filter(previous_cid_col.eq(hash.clone()))) | |||
| .set(is_current_col.eq(false)) | |||
| .execute(&mut conn) | |||
| .map_err(|e| AmendmentError::DatabaseError(Box::new(e)))?; | |||
| } | |||
| Ok(hash) | |||
| } | |||
| async fn upload_to_ipfs(&self, amendment: &AmendmentFile) -> IpfsResult<String, AmendmentError> { | |||
| upload_json_and_get_hash::<AmendmentFile, AmendmentError>(&self.client, STORAGE_DIR, FILE_EXTENSION, amendment).await | |||
| } | |||
| async fn save_to_database(&self, amendment: Amendment) -> IpfsResult<(), AmendmentError> { | |||
| let mut conn = get_connection() | |||
| .map_err(|e| AmendmentError::DatabaseError(e))?; | |||
| diesel::insert_into(amendments::table) | |||
| .values(&amendment) | |||
| .execute(&mut conn) | |||
| .map_err(|e| AmendmentError::DatabaseError(Box::new(e)))?; | |||
| Ok(()) | |||
| } | |||
| async fn read_amendment_file(&self, hash: &str) -> IpfsResult<AmendmentFile, AmendmentError> { | |||
| read_json_via_cat::<AmendmentFile, AmendmentError>(&self.client, hash, DEFAULT_MAX_JSON_SIZE).await | |||
| } | |||
| } | |||
| @@ -1,106 +0,0 @@ | |||
| use crate::ipfs::ipfs::IpfsService; | |||
| use crate::types::comment::{CommentError, CommentMetadata}; | |||
| use crate::types::ipfs::IpfsResult; | |||
| use crate::utils::ipfs::{ | |||
| create_file_path, | |||
| create_storage_directory, | |||
| read_json_via_cat, | |||
| retrieve_content_hash, | |||
| save_json_file, | |||
| DEFAULT_MAX_JSON_SIZE, | |||
| list_directory_file_hashes, | |||
| }; | |||
| use ipfs_api_backend_actix::IpfsClient; | |||
| const STORAGE_DIR: &str = "/puffpastry/comments"; | |||
| const FILE_EXTENSION: &str = "json"; | |||
| pub struct CommentService { | |||
| client: IpfsClient, | |||
| } | |||
| // Implement per-proposal subdirectory saving. The input is (proposal_cid, comments_batch) | |||
| impl IpfsService<(String, Vec<CommentMetadata>)> for CommentService { | |||
| type Err = CommentError; | |||
| async fn save(&mut self, item: (String, Vec<CommentMetadata>)) -> IpfsResult<String, Self::Err> { | |||
| let (proposal_cid, comments) = item; | |||
| // Allow batch save within the proposal's subdirectory. | |||
| if comments.is_empty() { | |||
| return Err(CommentError::from(std::io::Error::new( | |||
| std::io::ErrorKind::InvalidInput, | |||
| "Failed to store comment to IPFS: empty comments batch", | |||
| ))); | |||
| } | |||
| let mut last_cid: Option<String> = None; | |||
| for comment in comments { | |||
| let res = self | |||
| .store_comment_in_proposal_dir_and_publish(&proposal_cid, comment) | |||
| .await; | |||
| match res { | |||
| Ok(cid) => last_cid = Some(cid), | |||
| Err(e) => return Err(e), | |||
| } | |||
| } | |||
| match last_cid { | |||
| Some(cid) => Ok(cid), | |||
| None => Err(CommentError::from(std::io::Error::new( | |||
| std::io::ErrorKind::Other, | |||
| "Failed to store comment to IPFS", | |||
| ))), | |||
| } | |||
| } | |||
| async fn read(&mut self, hash: String) -> IpfsResult<(String, Vec<CommentMetadata>), Self::Err> { | |||
| // For reading, the caller should pass a directory hash (CID). We return only the comments vector here. | |||
| // Since IpfsService requires returning the same T, include an empty proposal id (unknown in read-by-hash). | |||
| // Alternatively, callers should not use the proposal id from this return value. | |||
| let comments = self.read_all_comments_in_dir(&hash).await?; | |||
| Ok((String::new(), comments)) | |||
| } | |||
| } | |||
| impl CommentService { | |||
| pub fn new(client: IpfsClient) -> Self { | |||
| Self { client } | |||
| } | |||
| // Writes a single comment JSON file into the MFS comments subdirectory for the proposal, | |||
| // then publishes the subdirectory snapshot (by retrieving the directory CID). Returns the comment file CID. | |||
| async fn store_comment_in_proposal_dir_and_publish( | |||
| &self, | |||
| proposal_cid: &str, | |||
| comment: CommentMetadata, | |||
| ) -> IpfsResult<String, CommentError> { | |||
| // Ensure the per-proposal storage subdirectory exists | |||
| let proposal_dir = format!("{}/{}", STORAGE_DIR, proposal_cid); | |||
| create_storage_directory::<CommentError>(&self.client, &proposal_dir).await?; | |||
| // Create a unique file path within the proposal subdirectory and save the JSON content | |||
| let file_path = create_file_path(&proposal_dir, FILE_EXTENSION); | |||
| save_json_file::<CommentMetadata, CommentError>(&self.client, &file_path, &comment).await?; | |||
| // Retrieve the file's CID to return to the caller | |||
| let file_cid = retrieve_content_hash::<CommentError>(&self.client, &file_path).await?; | |||
| // Publish a new snapshot of the proposal's comments subdirectory by retrieving its CID | |||
| let _dir_cid = retrieve_content_hash::<CommentError>(&self.client, &proposal_dir).await?; | |||
| Ok(file_cid) | |||
| } | |||
| async fn read_comment_file(&self, hash: &str) -> IpfsResult<CommentMetadata, CommentError> { | |||
| read_json_via_cat::<CommentMetadata, CommentError>(&self.client, hash, DEFAULT_MAX_JSON_SIZE).await | |||
| } | |||
| // Read all comment files within a directory identified by the given hash (directory CID) | |||
| async fn read_all_comments_in_dir(&self, dir_hash: &str) -> IpfsResult<Vec<CommentMetadata>, CommentError> { | |||
| let file_hashes = list_directory_file_hashes::<CommentError>(&self.client, dir_hash).await?; | |||
| let mut comments: Vec<CommentMetadata> = Vec::with_capacity(file_hashes.len()); | |||
| for fh in file_hashes { | |||
| let comment = self.read_comment_file(&fh).await?; | |||
| comments.push(comment); | |||
| } | |||
| Ok(comments) | |||
| } | |||
| } | |||
| @@ -1,7 +0,0 @@ | |||
| use crate::types::ipfs::IpfsResult; | |||
| pub trait IpfsService<T> { | |||
| type Err; | |||
| async fn save(&mut self, item: T) -> IpfsResult<String, Self::Err>; | |||
| async fn read(&mut self, hash: String) -> IpfsResult<T, Self::Err>; | |||
| } | |||
| @@ -1,4 +0,0 @@ | |||
| pub(crate) mod ipfs; | |||
| pub(crate) mod proposal; | |||
| pub(crate) mod amendment; | |||
| pub(crate) mod comment; | |||
| @@ -1,81 +0,0 @@ | |||
| use crate::db::db::get_connection; | |||
| use crate::ipfs::ipfs::IpfsService; | |||
| use crate::schema::proposals; | |||
| use crate::types::proposal::{Proposal, ProposalError, ProposalFile}; | |||
| use crate::utils::content::extract_summary; | |||
| use crate::utils::ipfs::{upload_json_and_get_hash, read_json_via_cat as util_read_json_via_cat, DEFAULT_MAX_JSON_SIZE}; | |||
| use diesel::{RunQueryDsl, QueryDsl, ExpressionMethods}; | |||
| use ipfs_api_backend_actix::IpfsClient; | |||
| use crate::types::ipfs::IpfsResult; | |||
| const STORAGE_DIR: &str = "/puffpastry/proposals"; | |||
| const FILE_EXTENSION: &str = "json"; | |||
| pub struct ProposalService { | |||
| client: IpfsClient, | |||
| } | |||
| impl IpfsService<ProposalFile> for ProposalService { | |||
| type Err = ProposalError; | |||
| async fn save(&mut self, proposal: ProposalFile) -> IpfsResult<String, Self::Err> { | |||
| self.store_proposal_to_ipfs(proposal).await | |||
| } | |||
| async fn read(&mut self, hash: String) -> IpfsResult<ProposalFile, Self::Err> { | |||
| self.read_proposal_file(hash).await | |||
| } | |||
| } | |||
| impl ProposalService { | |||
| pub fn new(client: IpfsClient) -> Self { | |||
| Self { client } | |||
| } | |||
| async fn store_proposal_to_ipfs(&self, proposal: ProposalFile) -> IpfsResult<String, ProposalError> { | |||
| let proposal_record = Proposal::new( | |||
| proposal.name.clone(), | |||
| Some(extract_summary(&proposal.content)), | |||
| proposal.creator.clone(), | |||
| ); | |||
| let hash = self.upload_to_ipfs(&proposal).await?; | |||
| self.save_to_database(proposal_record.with_cid(hash.clone()).mark_as_current()).await?; | |||
| // After saving the new proposal, set is_current = false for any proposal | |||
| // that has previous_cid equal to this new hash | |||
| { | |||
| use crate::schema::proposals::dsl::{proposals as proposals_table, previous_cid as previous_cid_col, is_current as is_current_col}; | |||
| let mut conn = get_connection() | |||
| .map_err(|e| ProposalError::DatabaseError(e))?; | |||
| // Ignore the count result; if no rows match, that's fine | |||
| let _ = diesel::update(proposals_table.filter(previous_cid_col.eq(hash.clone()))) | |||
| .set(is_current_col.eq(false)) | |||
| .execute(&mut conn) | |||
| .map_err(|e| ProposalError::DatabaseError(Box::new(e)))?; | |||
| } | |||
| Ok(hash) | |||
| } | |||
| async fn upload_to_ipfs(&self, data: &ProposalFile) -> IpfsResult<String, ProposalError> { | |||
| upload_json_and_get_hash::<ProposalFile, ProposalError>(&self.client, STORAGE_DIR, FILE_EXTENSION, data).await | |||
| } | |||
| async fn save_to_database(&self, proposal: Proposal) -> IpfsResult<(), ProposalError> { | |||
| let mut conn = get_connection() | |||
| .map_err(|e| ProposalError::DatabaseError(e))?; | |||
| diesel::insert_into(proposals::table) | |||
| .values(&proposal) | |||
| .execute(&mut conn) | |||
| .map_err(|e| ProposalError::DatabaseError(Box::new(e)))?; | |||
| Ok(()) | |||
| } | |||
| async fn read_proposal_file(&self, hash: String) -> IpfsResult<ProposalFile, ProposalError> { | |||
| util_read_json_via_cat::<ProposalFile, ProposalError>(&self.client, &hash, DEFAULT_MAX_JSON_SIZE).await | |||
| } | |||
| } | |||
| @@ -3,7 +3,6 @@ extern crate core; | |||
| mod routes; | |||
| mod schema; | |||
| mod db; | |||
| mod ipfs; | |||
| mod types; | |||
| mod repositories; | |||
| mod utils; | |||
| @@ -12,7 +11,6 @@ use crate::routes::proposal::{add_proposal, get_proposal, list_proposals, visit_ | |||
| use actix_web::{web, App, HttpServer}; | |||
| use crate::routes::amendment::{add_amendment, get_amendment, list_amendments}; | |||
| use crate::routes::auth::{register, login, login_freighter, get_nonce}; | |||
| use crate::routes::openapi::get_openapi; | |||
| use crate::utils::env::load_env; | |||
| #[actix_web::main] | |||
| @@ -26,7 +24,6 @@ async fn main() -> std::io::Result<()> { | |||
| println!("Starting server on http://{}:{}", ip, port); | |||
| HttpServer::new(|| { | |||
| App::new() | |||
| .service(get_openapi) | |||
| .service( | |||
| web::scope("/api/v1") | |||
| .service( | |||
| @@ -1,5 +1,4 @@ | |||
| use crate::db::db::establish_connection; | |||
| use crate::schema::amendments::dsl::amendments; | |||
| use crate::types::amendment::SelectableAmendment; | |||
| use diesel::{ExpressionMethods, QueryDsl, RunQueryDsl, SelectableHelper}; | |||
| @@ -27,23 +26,3 @@ pub fn get_amendments_for_proposal( | |||
| .order(crate::schema::amendments::id.asc()) | |||
| .load(&mut conn) | |||
| } | |||
| pub fn get_cid_for_amendment(amendment_id: i32) -> Result<String, diesel::result::Error> { | |||
| let pool = establish_connection() | |||
| .map_err(|_| diesel::result::Error::DatabaseError( | |||
| diesel::result::DatabaseErrorKind::Unknown, | |||
| Box::new("Failed to establish database connection".to_string()) | |||
| ))?; | |||
| let mut conn = pool.get() | |||
| .map_err(|_| diesel::result::Error::DatabaseError( | |||
| diesel::result::DatabaseErrorKind::Unknown, | |||
| Box::new("Failed to get database connection from pool".to_string()) | |||
| ))?; | |||
| use crate::schema::amendments::dsl::{id as id_col, cid as cid_col}; | |||
| amendments | |||
| .filter(id_col.eq(amendment_id)) | |||
| .select(cid_col) | |||
| .get_result(&mut conn) | |||
| } | |||
| @@ -1,76 +1,17 @@ | |||
| use diesel::{RunQueryDsl, ExpressionMethods, QueryDsl, Connection, BoolExpressionMethods}; | |||
| use crate::db::db::establish_connection; | |||
| use crate::db::db::get_connection; | |||
| use crate::schema::comments::dsl::*; | |||
| use crate::types::comment::NewComment; | |||
| use diesel::RunQueryDsl; | |||
| // Store the new CID for the updated comment thread file on IPFS | |||
| // Instead of updating an existing record, we: | |||
| // 1) Mark the previous record (by its CID) as not current (is_current=false) | |||
| // 2) Insert a new record with is_current=true and previous_cid set to the previous CID | |||
| pub fn store_new_thread_cid_for_proposal( | |||
| pid: i32, | |||
| new_cid_value: &str, | |||
| previous_cid_value: Option<&str>, | |||
| ) -> Result<(), diesel::result::Error> { | |||
| let pool = establish_connection() | |||
| .map_err(|_| diesel::result::Error::DatabaseError( | |||
| diesel::result::DatabaseErrorKind::Unknown, | |||
| Box::new("Failed to establish database connection".to_string()), | |||
| ))?; | |||
| let mut conn = pool.get().map_err(|_| { | |||
| diesel::result::Error::DatabaseError( | |||
| diesel::result::DatabaseErrorKind::Unknown, | |||
| Box::new("Failed to get database connection from pool".to_string()), | |||
| ) | |||
| })?; | |||
| conn.transaction(|conn| { | |||
| if let Some(prev) = previous_cid_value { | |||
| // Best-effort mark previous as not current; it's okay if 0 rows were updated (e.g., first insert) | |||
| let _ = diesel::update(comments.filter(cid.eq(prev))) | |||
| .set(is_current.eq(false)) | |||
| .execute(conn)?; | |||
| } | |||
| diesel::insert_into(crate::schema::comments::table) | |||
| .values(( | |||
| cid.eq(new_cid_value), | |||
| proposal_id.eq(pid), | |||
| is_current.eq(true), | |||
| previous_cid.eq(previous_cid_value), | |||
| )) | |||
| .execute(conn) | |||
| .map(|_| ()) | |||
| }) | |||
| } | |||
| // Retrieve the latest/current thread CID for a proposal | |||
| pub fn latest_thread_cid_for_proposal(pid_value: i32) -> Result<String, diesel::result::Error> { | |||
| let pool = establish_connection() | |||
| .map_err(|_| diesel::result::Error::DatabaseError( | |||
| diesel::result::DatabaseErrorKind::Unknown, | |||
| Box::new("Failed to establish database connection".to_string()), | |||
| ))?; | |||
| let mut conn = pool.get().map_err(|_| { | |||
| diesel::result::Error::DatabaseError( | |||
| diesel::result::DatabaseErrorKind::Unknown, | |||
| Box::new("Failed to get database connection from pool".to_string()), | |||
| ) | |||
| })?; | |||
| // Prefer the record marked as current; fallback to most recent if none are marked | |||
| match comments | |||
| .filter(proposal_id.eq(pid_value).and(is_current.eq(true))) | |||
| .order(created_at.desc()) | |||
| .select(cid) | |||
| .first::<String>(&mut conn) | |||
| { | |||
| Ok(current) => Ok(current), | |||
| Err(_) => comments | |||
| .filter(proposal_id.eq(pid_value)) | |||
| .order(created_at.desc()) | |||
| .select(cid) | |||
| .first::<String>(&mut conn), | |||
| } | |||
| } | |||
| pub(crate) fn create_comment(cont: String, prop_id: i32, usr_id: i32) { | |||
| let mut conn = get_connection().expect("Failed to establish connection"); | |||
| let new_comment = NewComment { | |||
| content: cont, | |||
| proposal_id: prop_id, | |||
| user_id: usr_id, | |||
| }; | |||
| diesel::insert_into(comments) | |||
| .values(&new_comment) | |||
| .execute(&mut conn) | |||
| .expect("Error saving new comment"); | |||
| } | |||
| @@ -1,48 +1,129 @@ | |||
| use diesel::ExpressionMethods; | |||
| use diesel::{QueryDsl, RunQueryDsl}; | |||
| use crate::db::db::establish_connection; | |||
| use crate::schema::proposals::dsl::*; | |||
| use crate::types::proposal::SelectableProposal; | |||
| pub fn paginate_proposals( | |||
| offset: i64, | |||
| limit: i64 | |||
| ) -> Result<Vec<SelectableProposal>, diesel::result::Error> { | |||
| use diesel::prelude::*; | |||
| use diesel::result::{DatabaseErrorKind, Error}; | |||
| use crate::db::db::{establish_connection, get_connection}; | |||
| use crate::schema::{proposal_paragraphs, proposals}; | |||
| use crate::types::proposal::{ | |||
| Proposal, ProposalParagraph, ProposalParagraphRequest, ProposalWithParagraphs, | |||
| ProposalWithSummary, SelectableProposal, | |||
| }; | |||
| fn db_error(message: impl Into<String>) -> Error { | |||
| Error::DatabaseError(DatabaseErrorKind::Unknown, Box::new(message.into())) | |||
| } | |||
| pub fn paginate_proposals(offset: i64, limit: i64) -> Result<Vec<ProposalWithSummary>, Error> { | |||
| let pool = establish_connection() | |||
| .map_err(|_| diesel::result::Error::DatabaseError( | |||
| diesel::result::DatabaseErrorKind::Unknown, | |||
| Box::new("Failed to establish database connection".to_string()) | |||
| ))?; | |||
| let mut conn = pool.get() | |||
| .map_err(|_| diesel::result::Error::DatabaseError( | |||
| diesel::result::DatabaseErrorKind::Unknown, | |||
| Box::new("Failed to get database connection from pool".to_string()) | |||
| ))?; | |||
| proposals | |||
| .filter(is_current.eq(true)) | |||
| .map_err(|_| db_error("Failed to establish database connection"))?; | |||
| let mut conn = pool | |||
| .get() | |||
| .map_err(|_| db_error("Failed to get database connection from pool"))?; | |||
| let props = proposals::table | |||
| .offset(offset) | |||
| .limit(limit) | |||
| .order(created_at.desc()) | |||
| .load::<SelectableProposal>(&mut conn) | |||
| .order(proposals::created_at.desc()) | |||
| .select(SelectableProposal::as_select()) | |||
| .load(&mut conn) | |||
| .unwrap_or_else(|_| { | |||
| eprintln!("Could not unwrap query"); | |||
| Vec::new() | |||
| }); | |||
| props | |||
| .into_iter() | |||
| .map(|p| { | |||
| println!("{}", p.id); | |||
| let summary = proposal_paragraphs::table | |||
| .filter(proposal_paragraphs::proposal_id.eq(p.id)) | |||
| .order(proposal_paragraphs::index.asc()) | |||
| .select(proposal_paragraphs::paragraph_text) | |||
| .first::<String>(&mut conn)?; | |||
| Ok(ProposalWithSummary { | |||
| id: p.id, | |||
| name: p.name, | |||
| creator: p.creator, | |||
| created_at: p.created_at, | |||
| updated_at: p.updated_at, | |||
| summary, | |||
| }) | |||
| }) | |||
| .collect() | |||
| } | |||
| pub fn get_cid_for_proposal(proposal_id: i32) -> Result<String, diesel::result::Error> { | |||
| let pool = establish_connection() | |||
| .map_err(|_| diesel::result::Error::DatabaseError( | |||
| diesel::result::DatabaseErrorKind::Unknown, | |||
| Box::new("Failed to establish database connection".to_string()) | |||
| ))?; | |||
| let mut conn = pool.get() | |||
| .map_err(|_| diesel::result::Error::DatabaseError( | |||
| diesel::result::DatabaseErrorKind::Unknown, | |||
| Box::new("Failed to get database connection from pool".to_string()) | |||
| ))?; | |||
| proposals. | |||
| filter(id.eq(proposal_id)) | |||
| .select(cid) | |||
| .get_result(&mut conn) | |||
| } | |||
| pub fn get_proposal(prop_id: i32) -> Result<ProposalWithParagraphs, Error> { | |||
| let mut conn = get_connection().expect("Could not get database connection"); | |||
| let proposal = proposals::table | |||
| .filter(proposals::id.eq(prop_id)) | |||
| .first::<SelectableProposal>(&mut conn) | |||
| .unwrap_or_else(|_| { | |||
| eprintln!("Could not unwrap proposal data"); | |||
| SelectableProposal { | |||
| id: 0, | |||
| name: String::new(), | |||
| creator: String::new(), | |||
| created_at: Default::default(), | |||
| updated_at: Default::default(), | |||
| } | |||
| }); | |||
| let paragraphs = proposal_paragraphs::table | |||
| .filter(proposal_paragraphs::proposal_id.eq(prop_id)) | |||
| .select(( | |||
| proposal_paragraphs::proposal_id, | |||
| proposal_paragraphs::paragraph_text, | |||
| proposal_paragraphs::index, | |||
| )) | |||
| .load::<ProposalParagraph>(&mut conn) | |||
| .unwrap_or_else(|_| { | |||
| eprintln!("Could not unwrap query"); | |||
| Vec::new() | |||
| }); | |||
| Ok(ProposalWithParagraphs { | |||
| id: prop_id, | |||
| name: proposal.name, | |||
| creator: proposal.creator, | |||
| created_at: proposal.created_at, | |||
| updated_at: proposal.updated_at, | |||
| paragraphs, | |||
| }) | |||
| } | |||
| pub fn save_proposal(proposal_name: String, proposal_creator: String) -> Result<i32, Error> { | |||
| let mut conn = get_connection().expect("Could not get database connection"); | |||
| let proposal = Proposal { | |||
| name: proposal_name, | |||
| creator: proposal_creator, | |||
| created_at: Default::default(), | |||
| updated_at: Default::default(), | |||
| }; | |||
| diesel::insert_into(proposals::table) | |||
| .values(&proposal) | |||
| .returning(proposals::id) | |||
| .get_result::<i32>(&mut conn) | |||
| .map_err(|e| db_error(format!("Failed to save proposal: {e}"))) | |||
| } | |||
| pub fn save_proposal_paragraphs(prop_id: i32, paragraphs: Vec<ProposalParagraphRequest>) { | |||
| let mut conn = get_connection().expect("Could not get database connection"); | |||
| let paragraph_items: Vec<ProposalParagraph> = paragraphs | |||
| .into_iter() | |||
| .map(|p| ProposalParagraph { | |||
| proposal_id: prop_id, | |||
| paragraph_text: p.paragraph_text, | |||
| index: p.index, | |||
| }) | |||
| .collect(); | |||
| diesel::insert_into(proposal_paragraphs::table) | |||
| .values(paragraph_items) | |||
| .execute(&mut conn) | |||
| .map_err(|e| db_error(format!("Failed to save proposal paragraphs: {e}"))) | |||
| .expect("Could not save proposal paragraphs"); | |||
| } | |||
| @@ -46,7 +46,7 @@ pub struct UserForAuth { | |||
| pub is_active: bool, | |||
| } | |||
| fn get_conn() -> Result<r2d2::PooledConnection<diesel::r2d2::ConnectionManager<PgConnection>>, diesel::result::Error> { | |||
| fn get_conn() -> Result<r2d2::PooledConnection<r2d2::ConnectionManager<PgConnection>>, diesel::result::Error> { | |||
| let pool = establish_connection().map_err(|_| { | |||
| diesel::result::Error::DatabaseError( | |||
| diesel::result::DatabaseErrorKind::Unknown, | |||
| @@ -1,11 +1,8 @@ | |||
| use actix_web::{get, post, web, HttpResponse}; | |||
| use ipfs_api_backend_actix::IpfsClient; | |||
| use crate::repositories::amendment::get_amendments_for_proposal; | |||
| use crate::types::amendment::AmendmentError; | |||
| use actix_web::{HttpResponse, get, post, web}; | |||
| use serde::{Deserialize, Serialize}; | |||
| use serde_json::json; | |||
| use crate::ipfs::amendment::AmendmentService; | |||
| use crate::ipfs::ipfs::IpfsService; | |||
| use crate::repositories::amendment::{get_amendments_for_proposal, get_cid_for_amendment}; | |||
| use crate::types::amendment::{AmendmentError, AmendmentFile}; | |||
| #[derive(Deserialize)] | |||
| #[serde(rename_all = "camelCase")] | |||
| @@ -70,33 +67,11 @@ pub struct GetAmendmentRequest { | |||
| } | |||
| #[get("get")] | |||
| async fn get_amendment(request: web::Query<GetAmendmentRequest>) -> Result<HttpResponse, actix_web::Error> { | |||
| let mut amendment_client = AmendmentService::new(IpfsClient::default()); | |||
| let cid = match get_cid_for_amendment(request.amendment_id) { | |||
| Ok(cid) => cid, | |||
| Err(err) => return Ok(HttpResponse::InternalServerError().json( | |||
| json!({"error": format!("Failed to get amendment: {}", err)}) | |||
| )) | |||
| }; | |||
| let item = amendment_client.read(cid).await; | |||
| match item { | |||
| Ok(item) => Ok(HttpResponse::Ok().json(json!({"amendment": item}))), | |||
| Err(e) => Ok(HttpResponse::InternalServerError().json(json!({"error": format!("Failed to read amendment: {}", e)}))) | |||
| } | |||
| // TODO: Implement actual logic to fetch an amendment by ID | |||
| Ok(HttpResponse::Ok().json(json!({"amendment": "placeholder"}))) | |||
| } | |||
| async fn create_and_store_amendment(request: AddAmendmentRequest) -> Result<String, AmendmentError> { | |||
| let amendment_data = AmendmentFile { | |||
| name: request.name.to_string(), | |||
| content: request.content.to_string(), | |||
| creator: request.creator.to_string(), | |||
| created_at: Default::default(), | |||
| updated_at: Default::default(), | |||
| proposal_id: request.proposal_id, | |||
| }; | |||
| let client = IpfsClient::default(); | |||
| let mut amendment_service = AmendmentService::new(client); | |||
| let cid = amendment_service.save(amendment_data).await?; | |||
| Ok(cid) | |||
| // TODO: Implement actual logic to create and store an amendment | |||
| Ok("placeholder".to_string()) | |||
| } | |||
| @@ -1,17 +1,13 @@ | |||
| use actix_web::{get, post, web, HttpResponse}; | |||
| use actix_web::{HttpResponse, get, post, web}; | |||
| use serde::Deserialize; | |||
| use serde_json::json; | |||
| use crate::repositories::user::{create_user, NewUserInsert, RegisteredUser, find_user_by_username, get_user_for_auth_by_stellar}; | |||
| use crate::repositories::user::{NewUserInsert, create_user, find_user_by_username, get_user_for_auth_by_stellar}; | |||
| const JWT_EXPIRATION_TIME: i64 = 240; | |||
| use argon2::{ | |||
| password_hash::{rand_core::OsRng, PasswordHasher, SaltString, PasswordHash, PasswordVerifier}, | |||
| Argon2, | |||
| }; | |||
| use crate::utils::auth::{verify_freighter_signature, generate_jwt, generate_freighter_nonce}; | |||
| use crate::utils::auth::{generate_freighter_nonce, generate_jwt}; | |||
| use argon2::{Argon2, PasswordHash, password_hash::{PasswordHasher, PasswordVerifier}}; | |||
| #[derive(Deserialize)] | |||
| #[serde(rename_all = "camelCase")] | |||
| @@ -39,9 +35,8 @@ pub async fn register(req: web::Json<RegisterRequest>) -> Result<HttpResponse, a | |||
| } | |||
| // Hash password with Argon2id and a random salt | |||
| let salt = SaltString::generate(&mut OsRng); | |||
| let argon2 = Argon2::default(); | |||
| let password_hash = match argon2.hash_password(req.password.as_bytes(), &salt) { | |||
| let password_hash = match argon2.hash_password(req.password.as_bytes()) { | |||
| Ok(ph) => ph.to_string(), | |||
| Err(e) => return Ok(HttpResponse::InternalServerError().json(json!({ | |||
| "error": format!("Failed to hash password: {}", e) | |||
| @@ -68,12 +63,6 @@ pub async fn register(req: web::Json<RegisterRequest>) -> Result<HttpResponse, a | |||
| } | |||
| } | |||
| #[derive(Deserialize)] | |||
| #[serde(rename_all = "camelCase")] | |||
| pub struct FreighterLoginRequest { | |||
| pub stellar_address: String, | |||
| } | |||
| #[post("login")] | |||
| pub async fn login(req: web::Json<LoginRequest>) -> Result<HttpResponse, actix_web::Error> { | |||
| @@ -145,6 +134,12 @@ pub async fn login(req: web::Json<LoginRequest>) -> Result<HttpResponse, actix_w | |||
| } | |||
| } | |||
| #[derive(Deserialize)] | |||
| #[serde(rename_all = "camelCase")] | |||
| pub struct FreighterLoginRequest { | |||
| pub stellar_address: String, | |||
| } | |||
| #[get("nonce")] | |||
| pub async fn get_nonce() -> HttpResponse { | |||
| let nonce = generate_freighter_nonce(); | |||
| @@ -1,10 +1,5 @@ | |||
| use actix_web::{post, web, HttpResponse}; | |||
| use chrono::Utc; | |||
| use ipfs_api_backend_actix::IpfsClient; | |||
| use serde::Deserialize; | |||
| use serde_json::json; | |||
| use crate::types::comment::{CommentError, CommentMetadata}; | |||
| #[derive(Deserialize)] | |||
| #[serde(rename_all = "camelCase")] | |||
| @@ -12,7 +7,4 @@ pub struct AppendLocalCommentRequest { | |||
| pub thread_id: String, | |||
| pub content: String, | |||
| pub creator: String, | |||
| pub proposal_cid: Option<String>, | |||
| pub amendment_cid: Option<String>, | |||
| pub parent_comment_cid: Option<String>, | |||
| } | |||
| @@ -1,5 +1,4 @@ | |||
| pub(crate) mod proposal; | |||
| pub(crate) mod amendment; | |||
| pub(crate) mod auth; | |||
| pub(crate) mod comment; | |||
| pub(crate) mod openapi; | |||
| pub(crate) mod comment; | |||
| @@ -1,251 +0,0 @@ | |||
| use actix_web::{get, HttpResponse}; | |||
| use serde_json::{json, Value}; | |||
| // Minimal OpenAPI 3.0 spec to enable frontend code generation. | |||
| // NOTE: This is intentionally lightweight. We can later replace the manual spec | |||
| // with an automatic generator (e.g., Apistos annotations) without changing the endpoint path. | |||
| fn openapi_spec() -> Value { | |||
| // Basic schemas for known request bodies | |||
| let add_proposal_request = json!({ | |||
| "type": "object", | |||
| "properties": { | |||
| "name": {"type": "string"}, | |||
| "description": {"type": "string"}, | |||
| "creator": {"type": "string"} | |||
| }, | |||
| "required": ["name", "description", "creator"], | |||
| "additionalProperties": false | |||
| }); | |||
| let visit_proposal_request = json!({ | |||
| "type": "object", | |||
| "properties": { | |||
| "proposalId": {"type": "integer", "format": "int32"}, | |||
| "path": {"type": "string"}, | |||
| "method": {"type": "string"}, | |||
| "queryString": {"type": "string"}, | |||
| "statusCode": {"type": "integer", "format": "int32"}, | |||
| "referrer": {"type": "string"} | |||
| }, | |||
| "required": ["proposalId"], | |||
| "additionalProperties": false | |||
| }); | |||
| let register_request = json!({ | |||
| "type": "object", | |||
| "properties": { | |||
| "username": {"type": "string"}, | |||
| "password": {"type": "string"}, | |||
| "stellarAddress": {"type": "string"}, | |||
| "email": {"type": ["string", "null"]} | |||
| }, | |||
| "required": ["username", "password", "stellarAddress"], | |||
| "additionalProperties": false | |||
| }); | |||
| let login_request = json!({ | |||
| "type": "object", | |||
| "properties": { | |||
| "username": {"type": "string"}, | |||
| "password": {"type": "string"} | |||
| }, | |||
| "required": ["username", "password"], | |||
| "additionalProperties": false | |||
| }); | |||
| let freighter_login_request = json!({ | |||
| "type": "object", | |||
| "properties": { | |||
| "stellarAddress": {"type": "string"} | |||
| }, | |||
| "required": ["stellarAddress"], | |||
| "additionalProperties": false | |||
| }); | |||
| let proposal_item = json!({ | |||
| "type": "object", | |||
| "properties": { | |||
| "id": {"type": "integer", "format": "int32"}, | |||
| "name": {"type": "string"}, | |||
| "cid": {"type": "string"}, | |||
| "summary": {"type": ["string", "null"]}, | |||
| "creator": {"type": "string"}, | |||
| "isCurrent": {"type": "boolean"}, | |||
| "previousCid": {"type": ["string", "null"]}, | |||
| "createdAt": {"type": ["string", "null"]}, | |||
| "updatedAt": {"type": ["string", "null"]} | |||
| }, | |||
| "required": ["id", "name", "cid", "creator", "isCurrent"], | |||
| "additionalProperties": false | |||
| }); | |||
| let proposal_list = json!({ | |||
| "type": "array", | |||
| "items": {"$ref": "#/components/schemas/ProposalItem"} | |||
| }); | |||
| let list_proposals_response = json!({ | |||
| "type": "object", | |||
| "properties": { | |||
| "proposals": {"$ref": "#/components/schemas/ProposalList"} | |||
| } | |||
| }); | |||
| // Generic JSON response schema when we don't have strict typing | |||
| let generic_object = json!({ | |||
| "type": "object", | |||
| "additionalProperties": true | |||
| }); | |||
| json!({ | |||
| "openapi": "3.0.3", | |||
| "info": { | |||
| "title": "PuffPastry API", | |||
| "version": "1.0.0" | |||
| }, | |||
| "servers": [ | |||
| {"url": "http://localhost:7300"} | |||
| ], | |||
| "paths": { | |||
| "/api/v1/proposals/list": { | |||
| "get": { | |||
| "summary": "List proposals", | |||
| "operationId": "listProposals", | |||
| "tags": ["Proposal"], | |||
| "parameters": [ | |||
| {"name": "offset", "in": "query", "required": true, "schema": {"type": "integer", "format": "int64"}}, | |||
| {"name": "limit", "in": "query", "required": true, "schema": {"type": "integer", "format": "int64"}} | |||
| ], | |||
| "responses": { | |||
| "200": {"description": "OK", "content": {"application/json": {"schema": {"$ref": "#/components/schemas/ListProposalsResponse"}}}} | |||
| } | |||
| } | |||
| }, | |||
| "/api/v1/proposal/add": { | |||
| "post": { | |||
| "summary": "Create a proposal", | |||
| "operationId": "createProposal", | |||
| "tags": ["Proposal"], | |||
| "requestBody": {"required": true, "content": {"application/json": {"schema": {"$ref": "#/components/schemas/AddProposalRequest"}}}}, | |||
| "responses": { | |||
| "200": {"description": "OK", "content": {"application/json": {"schema": generic_object}}}, | |||
| "400": {"description": "Bad Request"} | |||
| } | |||
| } | |||
| }, | |||
| "/api/v1/proposal/get": { | |||
| "get": { | |||
| "summary": "Get a proposal by id", | |||
| "operationId": "getProposal", | |||
| "tags": ["Proposal"], | |||
| "parameters": [ | |||
| {"name": "id", "in": "query", "required": true, "schema": {"type": "integer", "format": "int32"}} | |||
| ], | |||
| "responses": { | |||
| "200": {"description": "OK", "content": {"application/json": {"schema": generic_object}}} | |||
| } | |||
| } | |||
| }, | |||
| "/api/v1/proposal/visit": { | |||
| "post": { | |||
| "summary": "Record a proposal visit", | |||
| "operationId": "recordProposalVisit", | |||
| "tags": ["Proposal"], | |||
| "requestBody": {"required": true, "content": {"application/json": {"schema": {"$ref": "#/components/schemas/VisitProposalRequest"}}}}, | |||
| "responses": { | |||
| "201": {"description": "Created", "content": {"application/json": {"schema": generic_object}}} | |||
| } | |||
| } | |||
| }, | |||
| "/api/v1/amendment/add": { | |||
| "post": { | |||
| "summary": "Create an amendment", | |||
| "operationId": "createAmendment", | |||
| "tags": ["Amendment"], | |||
| "requestBody": {"required": true, "content": {"application/json": {"schema": generic_object}}}, | |||
| "responses": {"200": {"description": "OK"}} | |||
| } | |||
| }, | |||
| "/api/v1/amendment/get": { | |||
| "get": { | |||
| "summary": "Get an amendment", | |||
| "operationId": "getAmendment", | |||
| "tags": ["Amendment"], | |||
| "parameters": [ | |||
| {"name": "id", "in": "query", "required": true, "schema": {"type": "integer", "format": "int32"}} | |||
| ], | |||
| "responses": {"200": {"description": "OK", "content": {"application/json": {"schema": generic_object}}}} | |||
| } | |||
| }, | |||
| "/api/v1/amendments/list": { | |||
| "get": { | |||
| "summary": "List amendments", | |||
| "operationId": "listAmendments", | |||
| "tags": ["Amendment"], | |||
| "parameters": [ | |||
| {"name": "proposalId", "in": "query", "required": true, "schema": {"type": "integer", "format": "int32"}} | |||
| ], | |||
| "responses": {"200": {"description": "OK", "content": {"application/json": {"schema": generic_object}}}} | |||
| } | |||
| }, | |||
| "/api/v1/auth/register": { | |||
| "post": { | |||
| "summary": "Register a user", | |||
| "operationId": "registerUser", | |||
| "tags": ["Auth"], | |||
| "requestBody": {"required": true, "content": {"application/json": {"schema": {"$ref": "#/components/schemas/RegisterRequest"}}}}, | |||
| "responses": {"201": {"description": "Created", "content": {"application/json": {"schema": generic_object}}}, | |||
| "400": {"description": "Bad Request"}, | |||
| "409": {"description": "Conflict"}} | |||
| } | |||
| }, | |||
| "/api/v1/auth/login": { | |||
| "post": { | |||
| "summary": "Login with username/password", | |||
| "operationId": "loginUser", | |||
| "tags": ["Auth"], | |||
| "requestBody": {"required": true, "content": {"application/json": {"schema": {"$ref": "#/components/schemas/LoginRequest"}}}}, | |||
| "responses": {"200": {"description": "OK", "content": {"application/json": {"schema": generic_object}}}, | |||
| "401": {"description": "Unauthorized"}} | |||
| } | |||
| }, | |||
| "/api/v1/auth/nonce": { | |||
| "get": { | |||
| "summary": "Get freighter login nonce", | |||
| "operationId": "getFreighterNonce", | |||
| "tags": ["Auth"], | |||
| "responses": {"200": {"description": "OK", "content": {"application/json": {"schema": generic_object}}}} | |||
| } | |||
| }, | |||
| "/api/v1/auth/login/freighter": { | |||
| "post": { | |||
| "summary": "Login with Stellar (Freighter)", | |||
| "operationId": "loginFreighter", | |||
| "tags": ["Auth"], | |||
| "requestBody": {"required": true, "content": {"application/json": {"schema": {"$ref": "#/components/schemas/FreighterLoginRequest"}}}}, | |||
| "responses": {"200": {"description": "OK", "content": {"application/json": {"schema": generic_object}}}, | |||
| "401": {"description": "Unauthorized"}} | |||
| } | |||
| } | |||
| }, | |||
| "components": { | |||
| "schemas": { | |||
| "AddProposalRequest": add_proposal_request, | |||
| "VisitProposalRequest": visit_proposal_request, | |||
| "RegisterRequest": register_request, | |||
| "LoginRequest": login_request, | |||
| "FreighterLoginRequest": freighter_login_request, | |||
| "ListProposalsResponse": list_proposals_response, | |||
| "ProposalItem": proposal_item, | |||
| "ProposalList": proposal_list | |||
| } | |||
| } | |||
| }) | |||
| } | |||
| #[get("/openapi.json")] | |||
| pub async fn get_openapi() -> HttpResponse { | |||
| let spec = openapi_spec(); | |||
| HttpResponse::Ok().json(spec) | |||
| } | |||
| @@ -1,14 +1,11 @@ | |||
| use crate::ipfs::ipfs::IpfsService; | |||
| use crate::ipfs::proposal::ProposalService; | |||
| use crate::repositories::proposal::{get_cid_for_proposal, paginate_proposals}; | |||
| use crate::repositories::proposal::{paginate_proposals, save_proposal, save_proposal_paragraphs}; | |||
| use crate::repositories::visit::record_visit; | |||
| use crate::types::proposal::ProposalError; | |||
| use crate::types::proposal::ProposalFile; | |||
| use crate::types::proposal::{Proposal, ProposalParagraphRequest}; | |||
| use actix_web::{get, post, web, HttpResponse, HttpRequest}; | |||
| use chrono::Utc; | |||
| use ipfs_api_backend_actix::IpfsClient; | |||
| use serde::{Deserialize, Serialize}; | |||
| use serde_json::json; | |||
| use crate::repositories; | |||
| #[derive(Deserialize)] | |||
| struct PaginationParams { | |||
| @@ -25,45 +22,6 @@ async fn list_proposals(params: web::Query<PaginationParams>) -> Result<HttpResp | |||
| } | |||
| } | |||
| #[derive(Serialize, Deserialize)] | |||
| #[serde(rename_all = "camelCase")] | |||
| pub struct AddProposalRequest { | |||
| pub name: String, | |||
| pub description: String, | |||
| pub creator: String, | |||
| } | |||
| #[post("/add")] | |||
| async fn add_proposal( | |||
| request: web::Json<AddProposalRequest>, | |||
| ) -> Result<HttpResponse, actix_web::Error> { | |||
| let req = request.into_inner(); | |||
| // Basic input validation | |||
| if req.creator.trim().is_empty() { | |||
| return Ok(HttpResponse::BadRequest().json(json!({"error": "Creator is required"}))); | |||
| } | |||
| if req.name.trim().is_empty() { | |||
| return Ok(HttpResponse::BadRequest().json(json!({"error": "Name is required"}))); | |||
| } | |||
| if req.description.trim().is_empty() { | |||
| return Ok(HttpResponse::BadRequest().json(json!({"error": "Description is required"}))); | |||
| } | |||
| match create_and_store_proposal(&req).await { | |||
| Ok(cid) => Ok(HttpResponse::Ok().json(json!({ "cid": cid }))), | |||
| Err(e) => { | |||
| match e { | |||
| ProposalError::SerializationError(err) => Ok(HttpResponse::BadRequest().json(json!({ | |||
| "error": format!("Invalid proposal payload: {}", err) | |||
| }))), | |||
| _ => Ok(HttpResponse::InternalServerError().json(json!({ | |||
| "error": format!("Failed to create proposal: {}", e) | |||
| }))), | |||
| } | |||
| } | |||
| } | |||
| } | |||
| #[derive(Serialize, Deserialize)] | |||
| #[serde(rename_all = "camelCase")] | |||
| pub struct VisitProposalRequest { | |||
| @@ -122,36 +80,51 @@ struct GetProposalParams { | |||
| } | |||
| #[get("get")] | |||
| async fn get_proposal(params: web::Query<GetProposalParams>) -> Result<HttpResponse, actix_web::Error> { | |||
| let mut proposal_client = ProposalService::new(IpfsClient::default()); | |||
| let cid = match get_cid_for_proposal(params.id) { | |||
| Ok(cid) => cid, | |||
| Err(err) => return Ok(HttpResponse::InternalServerError().json(json!({ | |||
| "error": format!("Failed to fetch proposal CID: {}", err) | |||
| }))) | |||
| }; | |||
| let item = proposal_client.read(cid).await; | |||
| match item { | |||
| Ok(proposal) => Ok(HttpResponse::Ok().json(json!({"proposal": proposal}))), | |||
| Err(e) => Ok(HttpResponse::InternalServerError().json(json!({ | |||
| "error": format!("Failed to read proposal: {}", e) | |||
| }))) | |||
| let proposal = repositories::proposal::get_proposal(params.id); | |||
| match proposal { | |||
| Ok(p) => Ok(HttpResponse::Ok().json(json!({"proposal": p}))), | |||
| Err(e) => | |||
| Err(actix_web::error::ErrorInternalServerError(format!("Failed to get proposal: {}", e.to_string()))) | |||
| } | |||
| } | |||
| async fn create_and_store_proposal( | |||
| request: &AddProposalRequest, | |||
| ) -> Result<String, ProposalError> { | |||
| let proposal_data = ProposalFile { | |||
| name: request.name.to_string(), | |||
| content: request.description.to_string(), | |||
| creator: request.creator.to_string(), | |||
| #[derive(Serialize, Deserialize)] | |||
| #[serde(rename_all = "camelCase")] | |||
| pub struct AddProposalRequest { | |||
| pub name: String, | |||
| pub creator: String, | |||
| pub paragraphs: Vec<ProposalParagraphRequest>, | |||
| } | |||
| #[post("add")] | |||
| async fn add_proposal( | |||
| request: web::Json<AddProposalRequest>, | |||
| ) -> Result<HttpResponse, actix_web::Error> { | |||
| let req = request.into_inner(); | |||
| // Basic input validation | |||
| if req.creator.trim().is_empty() { | |||
| return Ok(HttpResponse::BadRequest().json(json!({"error": "Creator is required"}))); | |||
| } | |||
| if req.name.trim().is_empty() { | |||
| return Ok(HttpResponse::BadRequest().json(json!({"error": "Name is required"}))); | |||
| } | |||
| if req.paragraphs.is_empty() { | |||
| return Ok(HttpResponse::BadRequest().json(json!({"error": "Paragraphs are required"}))); | |||
| } | |||
| let proposal_data = Proposal { | |||
| name: req.name.to_string(), | |||
| creator: req.creator.to_string(), | |||
| created_at: Utc::now().naive_utc(), | |||
| updated_at: Utc::now().naive_utc(), | |||
| }; | |||
| let client = IpfsClient::default(); | |||
| let mut proposal_service = ProposalService::new(client); | |||
| let proposal_id = save_proposal( | |||
| proposal_data.clone().name, | |||
| proposal_data.clone().creator | |||
| ).expect("Could not save proposal"); | |||
| save_proposal_paragraphs(proposal_id, req.paragraphs); | |||
| let cid = proposal_service.save(proposal_data).await?; | |||
| Ok(cid) | |||
| Ok(HttpResponse::Ok().json(proposal_data)) | |||
| } | |||
| @@ -13,12 +13,9 @@ diesel::table! { | |||
| amendments (id) { | |||
| id -> Int4, | |||
| name -> Text, | |||
| cid -> Text, | |||
| summary -> Nullable<Text>, | |||
| status -> AmendmentStatus, | |||
| creator -> Text, | |||
| is_current -> Bool, | |||
| previous_cid -> Nullable<Text>, | |||
| created_at -> Timestamptz, | |||
| updated_at -> Timestamptz, | |||
| proposal_id -> Int4, | |||
| @@ -28,12 +25,21 @@ diesel::table! { | |||
| diesel::table! { | |||
| comments (id) { | |||
| id -> Int4, | |||
| cid -> Text, | |||
| content -> Nullable<Text>, | |||
| is_current -> Bool, | |||
| previous_cid -> Nullable<Text>, | |||
| proposal_id -> Int4, | |||
| created_at -> Timestamptz, | |||
| updated_at -> Timestamptz, | |||
| user_id -> Nullable<Int4>, | |||
| } | |||
| } | |||
| diesel::table! { | |||
| proposal_paragraphs (id) { | |||
| id -> Int4, | |||
| paragraph_text -> Text, | |||
| index -> Int4, | |||
| proposal_id -> Int4, | |||
| } | |||
| } | |||
| @@ -41,12 +47,8 @@ diesel::table! { | |||
| proposals (id) { | |||
| id -> Int4, | |||
| name -> Text, | |||
| cid -> Text, | |||
| summary -> Nullable<Text>, | |||
| creator -> Text, | |||
| is_current -> Bool, | |||
| previous_cid -> Nullable<Text>, | |||
| created_at -> Nullable<Timestamp>, | |||
| created_at -> Timestamp, | |||
| updated_at -> Nullable<Timestamp>, | |||
| } | |||
| } | |||
| @@ -89,11 +91,14 @@ diesel::table! { | |||
| diesel::joinable!(amendments -> proposals (proposal_id)); | |||
| diesel::joinable!(comments -> proposals (proposal_id)); | |||
| diesel::joinable!(comments -> users (user_id)); | |||
| diesel::joinable!(proposal_paragraphs -> proposals (proposal_id)); | |||
| diesel::joinable!(visits -> users (user_id)); | |||
| diesel::allow_tables_to_appear_in_same_query!( | |||
| amendments, | |||
| comments, | |||
| proposal_paragraphs, | |||
| proposals, | |||
| users, | |||
| visits, | |||
| @@ -48,12 +48,9 @@ pub enum AmendmentStatus { | |||
| pub struct SelectableAmendment { | |||
| pub id: i32, | |||
| pub name: String, | |||
| pub cid: String, | |||
| pub summary: Option<String>, | |||
| pub status: AmendmentStatus, | |||
| pub creator: String, | |||
| pub is_current: bool, | |||
| pub previous_cid: Option<String>, | |||
| pub created_at: NaiveDateTime, | |||
| pub updated_at: NaiveDateTime, | |||
| pub proposal_id: i32, | |||
| @@ -65,38 +62,27 @@ pub struct SelectableAmendment { | |||
| #[serde(rename_all = "camelCase")] | |||
| pub struct Amendment { | |||
| pub name: String, | |||
| pub cid: Option<String>, | |||
| pub summary: Option<String>, | |||
| pub status: AmendmentStatus, | |||
| pub creator: String, | |||
| pub is_current: bool, | |||
| pub previous_cid: Option<String>, | |||
| pub created_at: NaiveDateTime, | |||
| pub updated_at: NaiveDateTime, | |||
| pub proposal_id: i32, | |||
| } | |||
| impl Amendment { | |||
| pub fn new(name: String, summary: Option<String>, creator: String, proposal_id: i32) -> Self { | |||
| pub fn new(name: String, creator: String, proposal_id: i32) -> Self { | |||
| Amendment { | |||
| name, | |||
| summary, | |||
| status: AmendmentStatus::Proposed, | |||
| creator, | |||
| is_current: true, | |||
| previous_cid: None, | |||
| created_at: chrono::Utc::now().naive_utc(), | |||
| updated_at: chrono::Utc::now().naive_utc(), | |||
| proposal_id, | |||
| cid: None, | |||
| } | |||
| } | |||
| pub fn with_cid(mut self, cid: String) -> Self { | |||
| self.cid = Some(cid); | |||
| self | |||
| } | |||
| pub fn mark_as_current(mut self) -> Self { | |||
| self.is_current = true; | |||
| self.updated_at = chrono::Utc::now().naive_utc(); | |||
| @@ -1,6 +1,9 @@ | |||
| use chrono::NaiveDateTime; | |||
| use serde::{Deserialize, Serialize}; | |||
| use std::fmt; | |||
| use diesel::Insertable; | |||
| use crate::schema::comments; | |||
| #[derive(Debug)] | |||
| pub enum CommentError { | |||
| @@ -26,17 +29,6 @@ pub struct CommentMetadata { | |||
| pub updated_at: NaiveDateTime, | |||
| } | |||
| #[derive(Clone, Serialize, Deserialize)] | |||
| #[serde(rename_all = "camelCase")] | |||
| pub struct CommentFile { | |||
| pub comments: Vec<CommentMetadata>, | |||
| // Optional association fields to link this comment to a target entity | |||
| pub proposal_cid: Option<String>, | |||
| pub amendment_cid: Option<String>, | |||
| // Optional parent comment for threads | |||
| pub parent_comment_cid: Option<String>, | |||
| } | |||
| impl From<serde_json::Error> for CommentError { | |||
| fn from(e: serde_json::Error) -> Self { | |||
| CommentError::SerializationError(e) | |||
| @@ -48,3 +40,21 @@ impl From<std::io::Error> for CommentError { | |||
| CommentError::IpfsError(Box::new(e)) | |||
| } | |||
| } | |||
| #[derive(Clone, Serialize, Deserialize, Insertable)] | |||
| #[diesel(table_name = comments)] | |||
| pub struct NewComment { | |||
| pub content: String, | |||
| pub proposal_id: i32, | |||
| pub user_id: i32, | |||
| } | |||
| impl NewComment { | |||
| pub fn new(content: String, proposal_id: i32, user_id: i32) -> Self { | |||
| NewComment { | |||
| content, | |||
| proposal_id, | |||
| user_id, | |||
| } | |||
| } | |||
| } | |||
| @@ -1,50 +1,51 @@ | |||
| use crate::schema::proposals; | |||
| use std::error::Error; | |||
| use std::fmt; | |||
| use chrono::NaiveDateTime; | |||
| use diesel::Insertable; | |||
| use diesel::Queryable; | |||
| use diesel::Selectable; | |||
| use serde::de::StdError; | |||
| use diesel::{Insertable, Queryable, Selectable}; | |||
| use serde::{Deserialize, Serialize}; | |||
| use std::fmt; | |||
| use crate::schema::{proposal_paragraphs, proposals}; | |||
| #[derive(Debug)] | |||
| pub enum ProposalError { | |||
| ProposalNotFound, | |||
| DatabaseError(Box<dyn std::error::Error>), | |||
| IpfsError(Box<dyn std::error::Error>), | |||
| DatabaseError(Box<dyn Error>), | |||
| IpfsError(Box<dyn Error>), | |||
| SerializationError(serde_json::Error), | |||
| } | |||
| impl fmt::Display for ProposalError { | |||
| fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result { | |||
| fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { | |||
| match self { | |||
| ProposalError::ProposalNotFound => write!(f, "Proposal not found"), | |||
| ProposalError::DatabaseError(e) => write!(f, "Database error: {}", e), | |||
| ProposalError::IpfsError(e) => write!(f, "IPFS error: {}", e), | |||
| ProposalError::SerializationError(e) => write!(f, "Serialization error: {}", e), | |||
| Self::ProposalNotFound => write!(f, "Proposal not found"), | |||
| Self::DatabaseError(e) => write!(f, "Database error: {e}"), | |||
| Self::IpfsError(e) => write!(f, "IPFS error: {e}"), | |||
| Self::SerializationError(e) => write!(f, "Serialization error: {e}"), | |||
| } | |||
| } | |||
| } | |||
| impl StdError for ProposalError { | |||
| fn source(&self) -> Option<&(dyn StdError + 'static)> { | |||
| impl Error for ProposalError { | |||
| fn source(&self) -> Option<&(dyn Error + 'static)> { | |||
| match self { | |||
| ProposalError::ProposalNotFound => None, | |||
| ProposalError::DatabaseError(e) => Some(e.as_ref()), | |||
| ProposalError::IpfsError(e) => Some(e.as_ref()), | |||
| ProposalError::SerializationError(e) => Some(e), | |||
| Self::ProposalNotFound => None, | |||
| Self::DatabaseError(e) | Self::IpfsError(e) => Some(e.as_ref()), | |||
| Self::SerializationError(e) => Some(e), | |||
| } | |||
| } | |||
| } | |||
| #[derive(Serialize, Deserialize, Clone)] | |||
| #[serde(rename_all = "camelCase")] | |||
| pub struct ProposalFile { | |||
| pub name: String, | |||
| pub content: String, | |||
| pub creator: String, | |||
| pub created_at: NaiveDateTime, | |||
| pub updated_at: NaiveDateTime, | |||
| impl From<serde_json::Error> for ProposalError { | |||
| fn from(e: serde_json::Error) -> Self { | |||
| Self::SerializationError(e) | |||
| } | |||
| } | |||
| impl From<std::io::Error> for ProposalError { | |||
| fn from(e: std::io::Error) -> Self { | |||
| Self::IpfsError(Box::new(e)) | |||
| } | |||
| } | |||
| #[derive(Queryable, Selectable, Clone, Serialize, Deserialize)] | |||
| @@ -54,12 +55,8 @@ pub struct ProposalFile { | |||
| pub struct SelectableProposal { | |||
| pub id: i32, | |||
| pub name: String, | |||
| pub cid: String, | |||
| pub summary: Option<String>, | |||
| pub creator: String, | |||
| pub is_current: bool, | |||
| pub previous_cid: Option<String>, | |||
| pub created_at: Option<NaiveDateTime>, | |||
| pub created_at: NaiveDateTime, | |||
| pub updated_at: Option<NaiveDateTime>, | |||
| } | |||
| @@ -69,59 +66,57 @@ pub struct SelectableProposal { | |||
| #[serde(rename_all = "camelCase")] | |||
| pub struct Proposal { | |||
| pub name: String, | |||
| pub cid: Option<String>, | |||
| pub summary: Option<String>, | |||
| pub creator: String, | |||
| pub is_current: bool, | |||
| pub previous_cid: Option<String>, | |||
| pub created_at: NaiveDateTime, | |||
| pub updated_at: NaiveDateTime, | |||
| } | |||
| impl Proposal { | |||
| pub fn new(name: String, summary: Option<String>, creator: String) -> Self { | |||
| pub fn new(name: String, creator: String) -> Self { | |||
| let now = chrono::Local::now().naive_local(); | |||
| Self { | |||
| name, | |||
| cid: None, | |||
| creator, | |||
| summary, | |||
| is_current: false, | |||
| previous_cid: None, | |||
| created_at: chrono::Local::now().naive_local(), | |||
| updated_at: chrono::Local::now().naive_local(), | |||
| created_at: now, | |||
| updated_at: now, | |||
| } | |||
| } | |||
| } | |||
| pub fn with_cid(mut self, cid: String) -> Self { | |||
| self.cid = Some(cid); | |||
| self | |||
| } | |||
| pub fn with_summary(mut self, summary: String) -> Self { | |||
| self.summary = Some(summary); | |||
| self | |||
| } | |||
| pub fn with_previous_cid(mut self, previous_cid: String) -> Self { | |||
| self.previous_cid = Some(previous_cid); | |||
| self | |||
| } | |||
| #[derive(Serialize, Deserialize)] | |||
| pub struct ProposalParagraphRequest { | |||
| pub paragraph_text: String, | |||
| pub index: i32, | |||
| } | |||
| pub fn mark_as_current(mut self) -> Self { | |||
| self.is_current = true; | |||
| self.updated_at = chrono::Local::now().naive_local(); | |||
| self | |||
| } | |||
| #[derive(Queryable, Insertable, Clone, Serialize, Deserialize)] | |||
| #[diesel(table_name = proposal_paragraphs)] | |||
| #[diesel(check_for_backend(diesel::pg::Pg))] | |||
| #[serde(rename_all = "camelCase")] | |||
| pub struct ProposalParagraph { | |||
| pub proposal_id: i32, | |||
| pub paragraph_text: String, | |||
| pub index: i32, | |||
| } | |||
| impl From<serde_json::Error> for ProposalError { | |||
| fn from(e: serde_json::Error) -> Self { | |||
| ProposalError::SerializationError(e) | |||
| } | |||
| #[derive(Clone, Serialize, Deserialize)] | |||
| #[serde(rename_all = "camelCase")] | |||
| pub struct ProposalWithSummary { | |||
| pub id: i32, | |||
| pub name: String, | |||
| pub creator: String, | |||
| pub created_at: NaiveDateTime, | |||
| pub updated_at: Option<NaiveDateTime>, | |||
| pub summary: String, | |||
| } | |||
| impl From<std::io::Error> for ProposalError { | |||
| fn from(e: std::io::Error) -> Self { | |||
| ProposalError::IpfsError(Box::new(e)) | |||
| } | |||
| } | |||
| #[derive(Clone, Serialize, Deserialize)] | |||
| #[serde(rename_all = "camelCase")] | |||
| pub struct ProposalWithParagraphs { | |||
| pub id: i32, | |||
| pub name: String, | |||
| pub creator: String, | |||
| pub created_at: NaiveDateTime, | |||
| pub updated_at: Option<NaiveDateTime>, | |||
| pub paragraphs: Vec<ProposalParagraph>, | |||
| } | |||
| @@ -1,148 +0,0 @@ | |||
| use ipfs_api_backend_actix::{IpfsApi, IpfsClient}; | |||
| use serde::{de::DeserializeOwned, Serialize}; | |||
| use std::io::{Cursor, Error as IoError, ErrorKind}; | |||
| use std::path::Path; | |||
| use crate::types::ipfs::IpfsResult; | |||
| pub fn create_file_path(base_dir: &str, extension: &str) -> String { | |||
| let filename = format!("{}.{}", uuid::Uuid::new_v4(), extension); | |||
| let path = Path::new(base_dir).join(filename); | |||
| path.to_string_lossy().into_owned() | |||
| } | |||
| pub async fn create_storage_directory<E>(client: &IpfsClient, dir: &str) -> IpfsResult<(), E> | |||
| where | |||
| E: From<IoError>, | |||
| { | |||
| client | |||
| .files_mkdir(dir, true) | |||
| .await | |||
| .map_err(|e| E::from(IoError::new(ErrorKind::Other, format!("IPFS mkdir error: {}", e))))?; | |||
| Ok(()) | |||
| } | |||
| pub async fn save_json_file<T, E>( | |||
| client: &IpfsClient, | |||
| path: &str, | |||
| data: &T, | |||
| ) -> IpfsResult<(), E> | |||
| where | |||
| T: Serialize, | |||
| E: From<serde_json::Error> + From<IoError>, | |||
| { | |||
| let json = serde_json::to_string::<T>(data)?; | |||
| let file_content = Cursor::new(json.into_bytes()); | |||
| client | |||
| .files_write(path, true, true, file_content) | |||
| .await | |||
| .map_err(|e| E::from(IoError::new(ErrorKind::Other, format!("IPFS write error: {}", e))))?; | |||
| Ok(()) | |||
| } | |||
| pub async fn read_json_via_cat<T, E>( | |||
| client: &IpfsClient, | |||
| hash: &str, | |||
| max_size: usize, | |||
| ) -> IpfsResult<T, E> | |||
| where | |||
| T: DeserializeOwned, | |||
| E: From<serde_json::Error> + From<IoError>, | |||
| { | |||
| let stream = client.cat(hash); | |||
| let mut content = Vec::with_capacity(1024); | |||
| let mut total_size = 0; | |||
| futures::pin_mut!(stream); | |||
| while let Some(chunk) = futures::StreamExt::next(&mut stream).await { | |||
| let chunk = chunk.map_err(|e| E::from(IoError::new( | |||
| ErrorKind::Other, | |||
| format!("Failed to read IPFS chunk: {}", e), | |||
| )))?; | |||
| total_size += chunk.len(); | |||
| if total_size > max_size { | |||
| return Err(E::from(IoError::new( | |||
| ErrorKind::Other, | |||
| "File exceeds maximum allowed size", | |||
| ))); | |||
| } | |||
| content.extend_from_slice(&chunk); | |||
| } | |||
| if content.is_empty() { | |||
| return Err(E::from(IoError::new( | |||
| ErrorKind::Other, | |||
| "Empty response from IPFS", | |||
| ))); | |||
| } | |||
| let value = serde_json::from_slice::<T>(&content)?; | |||
| Ok(value) | |||
| } | |||
| pub async fn retrieve_content_hash<E>(client: &IpfsClient, path: &str) -> IpfsResult<String, E> | |||
| where | |||
| E: From<IoError>, | |||
| { | |||
| let stat = client | |||
| .files_stat(path) | |||
| .await | |||
| .map_err(|e| E::from(IoError::new(ErrorKind::Other, format!("IPFS stat error: {}", e))))?; | |||
| Ok(stat.hash) | |||
| } | |||
| pub const DEFAULT_MAX_JSON_SIZE: usize = 10 * 1024 * 1024; | |||
| /// List the child entries of an IPFS directory given its CID/hash and | |||
| /// return the CIDs of child files. | |||
| pub async fn list_directory_file_hashes<E>(client: &IpfsClient, dir_hash: &str) -> IpfsResult<Vec<String>, E> | |||
| where | |||
| E: From<IoError>, | |||
| { | |||
| // Use `ls` which works with CIDs (not only MFS paths) | |||
| let resp = client | |||
| .ls(dir_hash) | |||
| .await | |||
| .map_err(|e| E::from(IoError::new( | |||
| ErrorKind::Other, | |||
| format!("IPFS ls error: {}", e), | |||
| )))?; | |||
| // Collect links from the first (and usually only) object | |||
| let mut file_hashes: Vec<String> = Vec::new(); | |||
| for obj in resp.objects { | |||
| for link in obj.links { | |||
| // The response type typically has `typ` as a numeric code (2 for file) or a string. | |||
| // To remain compatible without depending on exact type semantics, we accept all links | |||
| // and rely on subsequent JSON parsing to validate. Optionally, filter by name extension. | |||
| if link.name.ends_with(".json") { | |||
| file_hashes.push(link.hash.clone()); | |||
| } else { | |||
| // Still include; some JSON files may not follow extension naming in IPFS. | |||
| file_hashes.push(link.hash.clone()); | |||
| } | |||
| } | |||
| } | |||
| Ok(file_hashes) | |||
| } | |||
| pub async fn upload_json_and_get_hash<T, E>( | |||
| client: &IpfsClient, | |||
| storage_dir: &str, | |||
| file_extension: &str, | |||
| data: &T, | |||
| ) -> IpfsResult<String, E> | |||
| where | |||
| T: Serialize, | |||
| E: From<serde_json::Error> + From<IoError>, | |||
| { | |||
| let file_path = create_file_path(storage_dir, file_extension); | |||
| create_storage_directory::<E>(client, storage_dir).await?; | |||
| save_json_file::<T, E>(client, &file_path, data).await?; | |||
| retrieve_content_hash::<E>(client, &file_path).await | |||
| } | |||
| @@ -1,4 +1,3 @@ | |||
| pub(crate) mod content; | |||
| pub(crate) mod ipfs; | |||
| pub(crate) mod auth; | |||
| pub(crate) mod env; | |||