diff --git a/src/infrastructure/client/brews.rs b/src/infrastructure/client/brews.rs new file mode 100644 index 0000000..17a5480 --- /dev/null +++ b/src/infrastructure/client/brews.rs @@ -0,0 +1,94 @@ +use anyhow::{Context, Result}; + +use crate::domain::brews::BrewWithDetails; +use crate::domain::ids::{BagId, BrewId, GearId}; + +use super::BrewlogClient; + +pub struct BrewsClient<'a> { + inner: &'a BrewlogClient, +} + +impl<'a> BrewsClient<'a> { + pub(crate) fn new(inner: &'a BrewlogClient) -> Self { + Self { inner } + } + + #[allow(clippy::too_many_arguments)] + pub async fn create( + &self, + bag_id: BagId, + coffee_weight: f64, + grinder_id: GearId, + grind_setting: f64, + brewer_id: GearId, + water_volume: i32, + water_temp: f64, + ) -> Result { + let url = self.inner.endpoint("api/v1/brews")?; + let payload = serde_json::json!({ + "bag_id": bag_id, + "coffee_weight": coffee_weight, + "grinder_id": grinder_id, + "grind_setting": grind_setting, + "brewer_id": brewer_id, + "water_volume": water_volume, + "water_temp": water_temp, + }); + + let response = self + .inner + .request(reqwest::Method::POST, url) + .json(&payload) + .send() + .await + .context("failed to issue create brew request")?; + + self.inner.handle_response(response).await + } + + pub async fn list(&self, bag_id: Option) -> Result> { + let mut url = self.inner.endpoint("api/v1/brews")?; + if let Some(bag_id) = bag_id { + url.query_pairs_mut() + .append_pair("bag_id", &bag_id.to_string()); + } + + let response = self + .inner + .request(reqwest::Method::GET, url) + .send() + .await + .context("failed to issue list brews request")?; + + self.inner.handle_response(response).await + } + + pub async fn get(&self, id: BrewId) -> Result { + let url = self.inner.endpoint(&format!("api/v1/brews/{id}"))?; + let response = self + .inner + .request(reqwest::Method::GET, url) + .send() + .await + .context("failed to issue get brew request")?; + + self.inner.handle_response(response).await + } + + pub async fn delete(&self, id: BrewId) -> Result<()> { + let url = self.inner.endpoint(&format!("api/v1/brews/{id}"))?; + let response = self + .inner + .request(reqwest::Method::DELETE, url) + .send() + .await + .context("failed to issue delete brew request")?; + + if response.status().is_success() { + Ok(()) + } else { + Err(self.inner.response_error(response).await) + } + } +} diff --git a/src/infrastructure/client/mod.rs b/src/infrastructure/client/mod.rs index 43cacf0..962b898 100644 --- a/src/infrastructure/client/mod.rs +++ b/src/infrastructure/client/mod.rs @@ -1,4 +1,5 @@ pub mod bags; +pub mod brews; pub mod gear; pub mod roasters; pub mod roasts; @@ -61,6 +62,10 @@ impl BrewlogClient { gear::GearClient::new(self) } + pub fn brews(&self) -> brews::BrewsClient<'_> { + brews::BrewsClient::new(self) + } + pub(crate) fn endpoint(&self, path: &str) -> Result { self.base_url .join(path) diff --git a/src/infrastructure/repositories/brews.rs b/src/infrastructure/repositories/brews.rs new file mode 100644 index 0000000..aadcbde --- /dev/null +++ b/src/infrastructure/repositories/brews.rs @@ -0,0 +1,269 @@ +use async_trait::async_trait; +use chrono::{DateTime, Utc}; +use sqlx::query_as; + +use crate::domain::RepositoryError; +use crate::domain::brews::{Brew, BrewFilter, BrewSortKey, BrewWithDetails, NewBrew}; +use crate::domain::ids::{BagId, BrewId, GearId}; +use crate::domain::listing::{ListRequest, Page, SortDirection}; +use crate::domain::repositories::BrewRepository; +use crate::infrastructure::database::DatabasePool; + +const BASE_SELECT: &str = r" + SELECT + br.id, br.bag_id, br.coffee_weight, br.grinder_id, br.grind_setting, + br.brewer_id, br.water_volume, br.water_temp, br.created_at, br.updated_at, + r.name as roast_name, r.slug as roast_slug, + rr.name as roaster_name, rr.slug as roaster_slug, + (g_grinder.make || ' ' || g_grinder.model) as grinder_name, + (g_brewer.make || ' ' || g_brewer.model) as brewer_name + FROM brews br + JOIN bags b ON br.bag_id = b.id + JOIN roasts r ON b.roast_id = r.id + JOIN roasters rr ON r.roaster_id = rr.id + JOIN gear g_grinder ON br.grinder_id = g_grinder.id + JOIN gear g_brewer ON br.brewer_id = g_brewer.id +"; + +#[derive(Clone)] +pub struct SqlBrewRepository { + pool: DatabasePool, +} + +impl SqlBrewRepository { + pub fn new(pool: DatabasePool) -> Self { + Self { pool } + } + + fn order_clause(request: &ListRequest) -> String { + let dir_sql = match request.sort_direction() { + SortDirection::Asc => "ASC", + SortDirection::Desc => "DESC", + }; + + match request.sort_key() { + BrewSortKey::CreatedAt => format!("br.created_at {dir_sql}, br.id DESC"), + BrewSortKey::CoffeeWeight => format!("br.coffee_weight {dir_sql}, br.created_at DESC"), + BrewSortKey::WaterVolume => format!("br.water_volume {dir_sql}, br.created_at DESC"), + } + } + + fn to_domain(record: BrewRecord) -> Brew { + Brew { + id: BrewId::new(record.id), + bag_id: BagId::new(record.bag_id), + coffee_weight: record.coffee_weight, + grinder_id: GearId::new(record.grinder_id), + grind_setting: record.grind_setting, + brewer_id: GearId::new(record.brewer_id), + water_volume: record.water_volume, + water_temp: record.water_temp, + created_at: record.created_at, + updated_at: record.updated_at, + } + } + + fn to_domain_with_details(record: BrewWithDetailsRecord) -> BrewWithDetails { + BrewWithDetails { + brew: Brew { + id: BrewId::new(record.id), + bag_id: BagId::new(record.bag_id), + coffee_weight: record.coffee_weight, + grinder_id: GearId::new(record.grinder_id), + grind_setting: record.grind_setting, + brewer_id: GearId::new(record.brewer_id), + water_volume: record.water_volume, + water_temp: record.water_temp, + created_at: record.created_at, + updated_at: record.updated_at, + }, + roast_name: record.roast_name, + roaster_name: record.roaster_name, + roast_slug: record.roast_slug, + roaster_slug: record.roaster_slug, + grinder_name: record.grinder_name, + brewer_name: record.brewer_name, + } + } + + fn build_where_clause(filter: &BrewFilter) -> Option { + // SAFETY: Direct interpolation is safe here because `bag_id` is an i64 from a typed wrapper. + filter + .bag_id + .map(|bag_id| format!("br.bag_id = {}", bag_id.into_inner())) + } +} + +#[async_trait] +impl BrewRepository for SqlBrewRepository { + async fn insert(&self, brew: NewBrew) -> Result { + // Use a transaction to atomically: + // 1. Deduct coffee_weight from bag's remaining + // 2. Insert the brew + let mut tx = self + .pool + .begin() + .await + .map_err(|err| RepositoryError::unexpected(err.to_string()))?; + + // Deduct coffee weight from bag's remaining amount + let update_bag_query = r" + UPDATE bags + SET remaining = remaining - ?, updated_at = CURRENT_TIMESTAMP + WHERE id = ? AND remaining >= ? AND closed = FALSE + "; + + let result = sqlx::query(update_bag_query) + .bind(brew.coffee_weight) + .bind(brew.bag_id.into_inner()) + .bind(brew.coffee_weight) + .execute(&mut *tx) + .await + .map_err(|err| RepositoryError::unexpected(err.to_string()))?; + + if result.rows_affected() == 0 { + return Err(RepositoryError::conflict( + "Insufficient coffee remaining in bag or bag is closed", + )); + } + + // Insert the brew + let insert_query = r" + INSERT INTO brews (bag_id, coffee_weight, grinder_id, grind_setting, brewer_id, water_volume, water_temp) + VALUES (?, ?, ?, ?, ?, ?, ?) + RETURNING id, bag_id, coffee_weight, grinder_id, grind_setting, brewer_id, water_volume, water_temp, created_at, updated_at + "; + + let record = query_as::<_, BrewRecord>(insert_query) + .bind(brew.bag_id.into_inner()) + .bind(brew.coffee_weight) + .bind(brew.grinder_id.into_inner()) + .bind(brew.grind_setting) + .bind(brew.brewer_id.into_inner()) + .bind(brew.water_volume) + .bind(brew.water_temp) + .fetch_one(&mut *tx) + .await + .map_err(|err| RepositoryError::unexpected(err.to_string()))?; + + tx.commit() + .await + .map_err(|err| RepositoryError::unexpected(err.to_string()))?; + + Ok(Self::to_domain(record)) + } + + async fn get(&self, id: BrewId) -> Result { + let query = r" + SELECT id, bag_id, coffee_weight, grinder_id, grind_setting, brewer_id, water_volume, water_temp, created_at, updated_at + FROM brews + WHERE id = ? + "; + + let record = query_as::<_, BrewRecord>(query) + .bind(id.into_inner()) + .fetch_optional(&self.pool) + .await + .map_err(|err| RepositoryError::unexpected(err.to_string()))? + .ok_or(RepositoryError::NotFound)?; + + Ok(Self::to_domain(record)) + } + + async fn get_with_details(&self, id: BrewId) -> Result { + let query = format!("{BASE_SELECT} WHERE br.id = ?"); + + let record = query_as::<_, BrewWithDetailsRecord>(&query) + .bind(id.into_inner()) + .fetch_optional(&self.pool) + .await + .map_err(|err| RepositoryError::unexpected(err.to_string()))? + .ok_or(RepositoryError::NotFound)?; + + Ok(Self::to_domain_with_details(record)) + } + + async fn list( + &self, + filter: BrewFilter, + request: &ListRequest, + ) -> Result, RepositoryError> { + let order_clause = Self::order_clause(request); + let where_clause = Self::build_where_clause(&filter); + + let base_query = match &where_clause { + Some(w) => format!("{BASE_SELECT} WHERE {w}"), + None => BASE_SELECT.to_string(), + }; + + let count_base = r" + SELECT COUNT(*) FROM brews br + JOIN bags b ON br.bag_id = b.id + "; + + let count_query = match &where_clause { + Some(w) => format!("{count_base} WHERE {w}"), + None => count_base.to_string(), + }; + + crate::infrastructure::repositories::pagination::paginate( + &self.pool, + request, + &base_query, + &count_query, + &order_clause, + |record| Ok(Self::to_domain_with_details(record)), + ) + .await + } + + async fn delete(&self, id: BrewId) -> Result<(), RepositoryError> { + let query = "DELETE FROM brews WHERE id = ?"; + + let result = sqlx::query(query) + .bind(id.into_inner()) + .execute(&self.pool) + .await + .map_err(|err| RepositoryError::unexpected(err.to_string()))?; + + if result.rows_affected() == 0 { + return Err(RepositoryError::NotFound); + } + + Ok(()) + } +} + +#[derive(sqlx::FromRow)] +struct BrewRecord { + id: i64, + bag_id: i64, + coffee_weight: f64, + grinder_id: i64, + grind_setting: f64, + brewer_id: i64, + water_volume: i32, + water_temp: f64, + created_at: DateTime, + updated_at: DateTime, +} + +#[derive(sqlx::FromRow)] +struct BrewWithDetailsRecord { + id: i64, + bag_id: i64, + coffee_weight: f64, + grinder_id: i64, + grind_setting: f64, + brewer_id: i64, + water_volume: i32, + water_temp: f64, + created_at: DateTime, + updated_at: DateTime, + roast_name: String, + roast_slug: String, + roaster_name: String, + roaster_slug: String, + grinder_name: String, + brewer_name: String, +} diff --git a/src/infrastructure/repositories/mod.rs b/src/infrastructure/repositories/mod.rs index e2c7f90..e771d0a 100644 --- a/src/infrastructure/repositories/mod.rs +++ b/src/infrastructure/repositories/mod.rs @@ -1,4 +1,5 @@ pub mod bags; +pub mod brews; pub mod gear; mod macros; pub mod pagination;