diff --git a/src/application/routes/api/coffee/bags.rs b/src/application/routes/api/coffee/bags.rs index e85c900..ff555da 100644 --- a/src/application/routes/api/coffee/bags.rs +++ b/src/application/routes/api/coffee/bags.rs @@ -221,6 +221,9 @@ pub(crate) async fn update_bag( info!(%id, closed = ?update.closed, "bag updated"); state.stats_invalidator.invalidate(); + state + .timeline_invalidator + .invalidate(EntityType::Bag, i64::from(bag.id)); save_deferred_image( &state, @@ -265,7 +268,8 @@ define_delete_handler!( bag_repo, render_bag_list_fragment, "type=bags", - "/data?type=bags" + "/data?type=bags", + entity_type: crate::domain::entity_type::EntityType::Bag ); #[derive(Debug, Deserialize)] diff --git a/src/application/routes/api/coffee/brews.rs b/src/application/routes/api/coffee/brews.rs index af09b26..ac19f9a 100644 --- a/src/application/routes/api/coffee/brews.rs +++ b/src/application/routes/api/coffee/brews.rs @@ -417,6 +417,9 @@ pub(crate) async fn update_brew( info!(%id, "brew updated"); state.stats_invalidator.invalidate(); + state + .timeline_invalidator + .invalidate(EntityType::Brew, i64::from(id)); save_deferred_image( &state, diff --git a/src/application/routes/api/coffee/cafes.rs b/src/application/routes/api/coffee/cafes.rs index f6b70e3..c1544ce 100644 --- a/src/application/routes/api/coffee/cafes.rs +++ b/src/application/routes/api/coffee/cafes.rs @@ -204,6 +204,9 @@ pub(crate) async fn update_cafe( .map_err(AppError::from)?; info!(%id, "cafe updated"); state.stats_invalidator.invalidate(); + state + .timeline_invalidator + .invalidate(EntityType::Cafe, i64::from(cafe.id)); save_deferred_image( &state, diff --git a/src/application/routes/api/coffee/cups.rs b/src/application/routes/api/coffee/cups.rs index a8e8e0c..c2e3a8d 100644 --- a/src/application/routes/api/coffee/cups.rs +++ b/src/application/routes/api/coffee/cups.rs @@ -143,6 +143,9 @@ pub(crate) async fn update_cup( info!(%id, "cup updated"); state.stats_invalidator.invalidate(); + state + .timeline_invalidator + .invalidate(EntityType::Cup, i64::from(cup.id)); save_deferred_image( &state, diff --git a/src/application/routes/api/coffee/gear.rs b/src/application/routes/api/coffee/gear.rs index 0c38ce3..5332a59 100644 --- a/src/application/routes/api/coffee/gear.rs +++ b/src/application/routes/api/coffee/gear.rs @@ -174,6 +174,9 @@ pub(crate) async fn update_gear( info!(%id, "gear updated"); state.stats_invalidator.invalidate(); + state + .timeline_invalidator + .invalidate(EntityType::Gear, i64::from(gear.id)); save_deferred_image( &state, diff --git a/src/application/routes/api/coffee/roasters.rs b/src/application/routes/api/coffee/roasters.rs index 608e6b9..ec86b98 100644 --- a/src/application/routes/api/coffee/roasters.rs +++ b/src/application/routes/api/coffee/roasters.rs @@ -195,6 +195,9 @@ pub(crate) async fn update_roaster( .map_err(AppError::from)?; info!(%id, "roaster updated"); state.stats_invalidator.invalidate(); + state + .timeline_invalidator + .invalidate(EntityType::Roaster, i64::from(roaster.id)); save_deferred_image( &state, diff --git a/src/application/routes/api/coffee/roasts.rs b/src/application/routes/api/coffee/roasts.rs index c9c6592..f084d98 100644 --- a/src/application/routes/api/coffee/roasts.rs +++ b/src/application/routes/api/coffee/roasts.rs @@ -249,6 +249,9 @@ pub(crate) async fn update_roast( info!(%id, "roast updated"); state.stats_invalidator.invalidate(); + state + .timeline_invalidator + .invalidate(EntityType::Roast, i64::from(id)); save_deferred_image( &state, diff --git a/src/application/routes/api/macros.rs b/src/application/routes/api/macros.rs index 85bcd98..e9942f2 100644 --- a/src/application/routes/api/macros.rs +++ b/src/application/routes/api/macros.rs @@ -86,13 +86,13 @@ macro_rules! define_enriched_get_handler { /// ); /// ``` macro_rules! define_delete_handler { - ($fn_name:ident, $id_type:ty, $sort_key:ty, $repo_field:ident, $render_fragment:path, $referer_match:literal, $redirect_url:literal) => { - define_delete_handler!(@inner $fn_name, $id_type, $sort_key, $repo_field, $render_fragment, $referer_match, $redirect_url, None); + ($fn_name:ident, $id_type:ty, $sort_key:ty, $repo_field:ident, $render_fragment:path, $referer_match:literal, $redirect_url:literal, entity_type: $entity_type:expr) => { + define_delete_handler!(@inner $fn_name, $id_type, $sort_key, $repo_field, $render_fragment, $referer_match, $redirect_url, $entity_type, None); }; ($fn_name:ident, $id_type:ty, $sort_key:ty, $repo_field:ident, $render_fragment:path, $referer_match:literal, $redirect_url:literal, image_type: $image_type:expr) => { - define_delete_handler!(@inner $fn_name, $id_type, $sort_key, $repo_field, $render_fragment, $referer_match, $redirect_url, Some($image_type)); + define_delete_handler!(@inner $fn_name, $id_type, $sort_key, $repo_field, $render_fragment, $referer_match, $redirect_url, $image_type, Some($image_type)); }; - (@inner $fn_name:ident, $id_type:ty, $sort_key:ty, $repo_field:ident, $render_fragment:path, $referer_match:literal, $redirect_url:literal, $image_type:expr) => { + (@inner $fn_name:ident, $id_type:ty, $sort_key:ty, $repo_field:ident, $render_fragment:path, $referer_match:literal, $redirect_url:literal, $entity_type:expr, $image_type:expr) => { #[tracing::instrument(skip(state, _auth_user, headers, query))] pub(crate) async fn $fn_name( axum::extract::State(state): axum::extract::State, @@ -116,6 +116,10 @@ macro_rules! define_delete_handler { } } + if let Err(err) = state.timeline_repo.delete_by_entity($entity_type, i64::from(id)).await { + tracing::warn!(%id, error = %err, "failed to delete timeline event"); + } + tracing::info!(%id, "entity deleted"); state.stats_invalidator.invalidate(); diff --git a/src/application/routes/api/mod.rs b/src/application/routes/api/mod.rs index 82e0887..298d014 100644 --- a/src/application/routes/api/mod.rs +++ b/src/application/routes/api/mod.rs @@ -9,7 +9,7 @@ pub(crate) mod system; pub(crate) use analytics::stats; pub(crate) use auth::{tokens, webauthn}; pub(crate) use coffee::{bags, brews, cafes, checkin, cups, gear, roasters, roasts, scan}; -pub(crate) use system::{admin, backup}; +pub(crate) use system::{admin, backup, timeline}; use axum::extract::DefaultBodyLimit; use axum::routing::{get, post}; @@ -103,6 +103,7 @@ pub(super) fn router() -> axum::Router { ) .route("/backup/reset", post(backup::reset_database)) .route("/stats/recompute", post(stats::recompute_stats)) + .route("/timeline/rebuild", post(timeline::rebuild_timeline)) .route( "/{entity_type}/{id}/image", get(images::get_image) diff --git a/src/application/routes/api/system/mod.rs b/src/application/routes/api/system/mod.rs index e8355e8..71918d5 100644 --- a/src/application/routes/api/system/mod.rs +++ b/src/application/routes/api/system/mod.rs @@ -1,2 +1,3 @@ pub(crate) mod admin; pub(crate) mod backup; +pub(crate) mod timeline; diff --git a/src/application/routes/api/system/timeline.rs b/src/application/routes/api/system/timeline.rs new file mode 100644 index 0000000..c551cce --- /dev/null +++ b/src/application/routes/api/system/timeline.rs @@ -0,0 +1,18 @@ +use axum::extract::State; +use axum::http::StatusCode; +use axum::response::{IntoResponse, Response}; +use tracing::info; + +use crate::application::auth::AuthenticatedUser; +use crate::application::errors::ApiError; +use crate::application::state::AppState; + +#[tracing::instrument(skip(state, _auth_user))] +pub(crate) async fn rebuild_timeline( + State(state): State, + _auth_user: AuthenticatedUser, +) -> Result { + info!("timeline rebuild requested"); + state.timeline_invalidator.rebuild_all(); + Ok(StatusCode::NO_CONTENT.into_response()) +} diff --git a/src/application/server.rs b/src/application/server.rs index 6379b54..2012ce6 100644 --- a/src/application/server.rs +++ b/src/application/server.rs @@ -9,8 +9,9 @@ use tracing::info; use webauthn_rs::prelude::*; use crate::application::routes::app_router; -use crate::application::services::StatsInvalidator; use crate::application::services::stats::stats_recomputation_task; +use crate::application::services::timeline_refresh::{TimelineRebuilder, timeline_rebuild_task}; +use crate::application::services::{StatsInvalidator, TimelineInvalidator}; use crate::application::state::{AppState, AppStateConfig}; use crate::domain::registration_tokens::NewRegistrationToken; use crate::domain::repositories::{RegistrationTokenRepository, UserRepository}; @@ -45,6 +46,11 @@ pub async fn serve(config: ServerConfig) -> anyhow::Result<()> { let (stats_tx, stats_rx) = tokio::sync::mpsc::channel::<()>(32); let stats_invalidator = StatsInvalidator::new(stats_tx); + let (timeline_tx, timeline_rx) = tokio::sync::mpsc::channel::< + crate::application::services::timeline_refresh::TimelineInvalidation, + >(32); + let timeline_invalidator = TimelineInvalidator::new(timeline_tx); + let state = AppState::from_database( &database, AppStateConfig { @@ -56,6 +62,7 @@ pub async fn serve(config: ServerConfig) -> anyhow::Result<()> { openrouter_api_key: config.openrouter_api_key, openrouter_model: config.openrouter_model, stats_invalidator: stats_invalidator.clone(), + timeline_invalidator, }, ); @@ -67,6 +74,23 @@ pub async fn serve(config: ServerConfig) -> anyhow::Result<()> { std::time::Duration::from_secs(2), )); + // Spawn background timeline rebuild task + let rebuilder = TimelineRebuilder { + timeline_repo: Arc::clone(&state.timeline_repo), + roaster_repo: Arc::clone(&state.roaster_repo), + roast_repo: Arc::clone(&state.roast_repo), + bag_repo: Arc::clone(&state.bag_repo), + brew_repo: Arc::clone(&state.brew_repo), + cup_repo: Arc::clone(&state.cup_repo), + gear_repo: Arc::clone(&state.gear_repo), + cafe_repo: Arc::clone(&state.cafe_repo), + }; + tokio::spawn(timeline_rebuild_task( + timeline_rx, + rebuilder, + std::time::Duration::from_secs(2), + )); + // Seed the stats cache on startup stats_invalidator.invalidate(); diff --git a/src/application/services/mod.rs b/src/application/services/mod.rs index 7fe9bee..a9a8fb7 100644 --- a/src/application/services/mod.rs +++ b/src/application/services/mod.rs @@ -3,12 +3,14 @@ mod brews; mod cups; mod roasts; pub mod stats; +pub mod timeline_refresh; pub use bags::BagService; pub use brews::BrewService; pub use cups::CupService; pub use roasts::RoastService; pub use stats::StatsInvalidator; +pub use timeline_refresh::TimelineInvalidator; use std::sync::Arc; diff --git a/src/application/services/timeline_refresh.rs b/src/application/services/timeline_refresh.rs new file mode 100644 index 0000000..b0a8a31 --- /dev/null +++ b/src/application/services/timeline_refresh.rs @@ -0,0 +1,478 @@ +use std::collections::HashSet; +use std::sync::Arc; +use std::time::Duration; + +use tokio::sync::mpsc; +use tracing::{error, info, warn}; + +use crate::domain::bags::bag_timeline_event; +use crate::domain::entity_type::EntityType; +use crate::domain::ids::{BagId, BrewId, CafeId, CupId, GearId, RoastId, RoasterId}; +use crate::domain::repositories::{ + BagRepository, BrewRepository, CafeRepository, CupRepository, GearRepository, RoastRepository, + RoasterRepository, TimelineEventRepository, +}; +use crate::domain::roasts::roast_timeline_event; + +/// Invalidation signal sent by HTTP handlers to the background rebuild task. +#[derive(Debug, Clone, PartialEq, Eq, Hash)] +pub enum TimelineInvalidation { + /// Refresh a single entity's timeline event(s) and cascade to downstream entities. + Entity { + entity_type: EntityType, + entity_id: i64, + }, + /// Delete all timeline events and rebuild from scratch. + Full, +} + +/// Sends invalidation signals to the background timeline rebuild task. +/// Non-blocking and fire-and-forget — safe to call from any handler. +#[derive(Clone)] +pub struct TimelineInvalidator { + tx: mpsc::Sender, +} + +impl TimelineInvalidator { + pub fn new(tx: mpsc::Sender) -> Self { + Self { tx } + } + + /// Signal that a specific entity's timeline event(s) need refreshing. + pub fn invalidate(&self, entity_type: EntityType, entity_id: i64) { + let _ = self.tx.try_send(TimelineInvalidation::Entity { + entity_type, + entity_id, + }); + } + + /// Signal a full rebuild of all timeline events. + pub fn rebuild_all(&self) { + let _ = self.tx.try_send(TimelineInvalidation::Full); + } +} + +/// All repository dependencies needed to rebuild timeline events. +pub struct TimelineRebuilder { + pub timeline_repo: Arc, + pub roaster_repo: Arc, + pub roast_repo: Arc, + pub bag_repo: Arc, + pub brew_repo: Arc, + pub cup_repo: Arc, + pub gear_repo: Arc, + pub cafe_repo: Arc, +} + +/// Listens for invalidation signals, debounces, and rebuilds affected timeline events. +/// Runs as a long-lived background task — spawn with `tokio::spawn`. +pub async fn timeline_rebuild_task( + mut rx: mpsc::Receiver, + rebuilder: TimelineRebuilder, + debounce: Duration, +) { + loop { + let Some(first) = rx.recv().await else { + break; + }; + + // Debounce: wait then drain any accumulated signals + tokio::time::sleep(debounce).await; + + let mut dirty = HashSet::new(); + let mut full_rebuild = matches!(first, TimelineInvalidation::Full); + if let TimelineInvalidation::Entity { + entity_type, + entity_id, + } = first + { + dirty.insert((entity_type, entity_id)); + } + + while let Ok(signal) = rx.try_recv() { + match signal { + TimelineInvalidation::Full => full_rebuild = true, + TimelineInvalidation::Entity { + entity_type, + entity_id, + } => { + dirty.insert((entity_type, entity_id)); + } + } + } + + if full_rebuild { + if let Err(err) = rebuild_all(&rebuilder).await { + error!(error = %err, "timeline full rebuild failed"); + } + } else { + // Expand cascade targets before refreshing + let mut all_targets = HashSet::new(); + for (entity_type, entity_id) in &dirty { + collect_cascade_targets(&rebuilder, *entity_type, *entity_id, &mut all_targets) + .await; + } + + for (entity_type, entity_id) in &all_targets { + if let Err(err) = refresh_entity_event(&rebuilder, *entity_type, *entity_id).await { + warn!( + error = %err, + entity_type = entity_type.as_str(), + entity_id, + "failed to refresh timeline event" + ); + } + } + + if !all_targets.is_empty() { + info!(count = all_targets.len(), "timeline events refreshed"); + } + } + } +} + +/// Collect the edited entity plus all downstream cascade targets into `targets`. +async fn collect_cascade_targets( + rebuilder: &TimelineRebuilder, + entity_type: EntityType, + entity_id: i64, + targets: &mut HashSet<(EntityType, i64)>, +) { + targets.insert((entity_type, entity_id)); + + match entity_type { + EntityType::Roaster => { + // Roaster → Roasts → Bags → Brews, and Roasts → Cups + let roasts = rebuilder + .roast_repo + .list_by_roaster(RoasterId::new(entity_id)) + .await + .unwrap_or_default(); + for rwr in &roasts { + let roast_id = rwr.roast.id.into_inner(); + targets.insert((EntityType::Roast, roast_id)); + collect_roast_downstream(rebuilder, rwr.roast.id, targets).await; + } + } + EntityType::Roast => { + collect_roast_downstream(rebuilder, RoastId::new(entity_id), targets).await; + } + EntityType::Bag => { + collect_bag_downstream(rebuilder, BagId::new(entity_id), targets).await; + } + EntityType::Gear => { + // Gear → Brews that use this gear piece + let filter = crate::domain::brews::BrewFilter::for_gear(GearId::new(entity_id)); + if let Ok(page) = rebuilder + .brew_repo + .list( + filter, + &show_all_request::(), + None, + ) + .await + { + for bwd in &page.items { + targets.insert((EntityType::Brew, bwd.brew.id.into_inner())); + } + } + } + EntityType::Cafe => { + // Cafe → Cups at this cafe + let filter = crate::domain::cups::CupFilter { + cafe_id: Some(CafeId::new(entity_id)), + ..Default::default() + }; + if let Ok(page) = rebuilder + .cup_repo + .list( + filter, + &show_all_request::(), + None, + ) + .await + { + for cwd in &page.items { + targets.insert((EntityType::Cup, cwd.cup.id.into_inner())); + } + } + } + EntityType::Brew | EntityType::Cup => { + // Leaf entities — no downstream cascade + } + } +} + +/// Collect bags, brews, and cups downstream of a roast. +async fn collect_roast_downstream( + rebuilder: &TimelineRebuilder, + roast_id: RoastId, + targets: &mut HashSet<(EntityType, i64)>, +) { + // Bags for this roast + let bag_filter = crate::domain::bags::BagFilter { + roast_id: Some(roast_id), + ..Default::default() + }; + if let Ok(page) = rebuilder + .bag_repo + .list( + bag_filter, + &show_all_request::(), + None, + ) + .await + { + for bwr in &page.items { + targets.insert((EntityType::Bag, bwr.bag.id.into_inner())); + collect_bag_downstream(rebuilder, bwr.bag.id, targets).await; + } + } + + // Cups for this roast + let cup_filter = crate::domain::cups::CupFilter { + roast_id: Some(roast_id), + ..Default::default() + }; + if let Ok(page) = rebuilder + .cup_repo + .list( + cup_filter, + &show_all_request::(), + None, + ) + .await + { + for cwd in &page.items { + targets.insert((EntityType::Cup, cwd.cup.id.into_inner())); + } + } +} + +/// Collect brews downstream of a bag. +async fn collect_bag_downstream( + rebuilder: &TimelineRebuilder, + bag_id: BagId, + targets: &mut HashSet<(EntityType, i64)>, +) { + let filter = crate::domain::brews::BrewFilter::for_bag(bag_id); + if let Ok(page) = rebuilder + .brew_repo + .list( + filter, + &show_all_request::(), + None, + ) + .await + { + for bwd in &page.items { + targets.insert((EntityType::Brew, bwd.brew.id.into_inner())); + } + } +} + +/// Refresh a single entity's timeline event(s) by regenerating from current data. +async fn refresh_entity_event( + rebuilder: &TimelineRebuilder, + entity_type: EntityType, + entity_id: i64, +) -> Result<(), crate::domain::RepositoryError> { + let event = match entity_type { + EntityType::Roaster => { + let roaster = rebuilder + .roaster_repo + .get(RoasterId::new(entity_id)) + .await?; + roaster.to_timeline_event() + } + EntityType::Roast => { + let rwr = rebuilder + .roast_repo + .get_with_roaster(RoastId::new(entity_id)) + .await?; + let roaster = rebuilder.roaster_repo.get(rwr.roast.roaster_id).await?; + roast_timeline_event(&rwr.roast, &roaster) + } + EntityType::Bag => { + let bwr = rebuilder + .bag_repo + .get_with_roast(BagId::new(entity_id)) + .await?; + let roast = rebuilder.roast_repo.get(bwr.bag.roast_id).await?; + let roaster = rebuilder.roaster_repo.get(roast.roaster_id).await?; + // update_by_entity updates all events for this entity, + // preserving each event's original action ("added" or "finished") + bag_timeline_event(&bwr.bag, "added", &roast, &roaster) + } + EntityType::Brew => { + let enriched = rebuilder + .brew_repo + .get_with_details(BrewId::new(entity_id)) + .await?; + enriched.to_timeline_event() + } + EntityType::Cup => { + let enriched = rebuilder + .cup_repo + .get_with_details(CupId::new(entity_id)) + .await?; + enriched.to_timeline_event() + } + EntityType::Gear => { + let gear = rebuilder.gear_repo.get(GearId::new(entity_id)).await?; + gear.to_timeline_event() + } + EntityType::Cafe => { + let cafe = rebuilder.cafe_repo.get(CafeId::new(entity_id)).await?; + cafe.to_timeline_event() + } + }; + + rebuilder + .timeline_repo + .update_by_entity(entity_type, entity_id, event) + .await +} + +/// Delete all timeline events and rebuild from current entity data. +pub async fn rebuild_all( + rebuilder: &TimelineRebuilder, +) -> Result<(), crate::domain::RepositoryError> { + let start = std::time::Instant::now(); + + rebuilder.timeline_repo.delete_all().await?; + + // Roasters + let roasters = rebuilder.roaster_repo.list_all().await?; + for roaster in &roasters { + if let Err(err) = rebuilder + .timeline_repo + .insert(roaster.to_timeline_event()) + .await + { + warn!(error = %err, id = %roaster.id, "failed to rebuild roaster timeline event"); + } + } + + // Cafes + let cafes = rebuilder.cafe_repo.list_all().await?; + for cafe in &cafes { + if let Err(err) = rebuilder + .timeline_repo + .insert(cafe.to_timeline_event()) + .await + { + warn!(error = %err, id = %cafe.id, "failed to rebuild cafe timeline event"); + } + } + + // Gear + let gear_list = rebuilder.gear_repo.list_all().await?; + for gear in &gear_list { + if let Err(err) = rebuilder + .timeline_repo + .insert(gear.to_timeline_event()) + .await + { + warn!(error = %err, id = %gear.id, "failed to rebuild gear timeline event"); + } + } + + // Roasts (need roaster for each) + let roasts = rebuilder.roast_repo.list_all().await?; + for rwr in &roasts { + let roaster = match rebuilder.roaster_repo.get(rwr.roast.roaster_id).await { + Ok(r) => r, + Err(err) => { + warn!(error = %err, id = %rwr.roast.id, "failed to fetch roaster for roast timeline rebuild"); + continue; + } + }; + if let Err(err) = rebuilder + .timeline_repo + .insert(roast_timeline_event(&rwr.roast, &roaster)) + .await + { + warn!(error = %err, id = %rwr.roast.id, "failed to rebuild roast timeline event"); + } + } + + rebuild_bag_events(rebuilder).await?; + + // Brews (get_with_details for enrichment) + let brews = rebuilder.brew_repo.list_all().await?; + for bwd in &brews { + if let Err(err) = rebuilder + .timeline_repo + .insert(bwd.to_timeline_event()) + .await + { + warn!(error = %err, id = %bwd.brew.id, "failed to rebuild brew timeline event"); + } + } + + // Cups (list_all returns CupWithDetails) + let cups = rebuilder.cup_repo.list_all().await?; + for cwd in &cups { + if let Err(err) = rebuilder + .timeline_repo + .insert(cwd.to_timeline_event()) + .await + { + warn!(error = %err, id = %cwd.cup.id, "failed to rebuild cup timeline event"); + } + } + + info!( + duration_ms = start.elapsed().as_millis(), + "timeline events rebuilt" + ); + Ok(()) +} + +/// Rebuild timeline events for all bags (need roast + roaster lookups for each). +async fn rebuild_bag_events( + rebuilder: &TimelineRebuilder, +) -> Result<(), crate::domain::RepositoryError> { + let bags = rebuilder.bag_repo.list_all().await?; + for bwr in &bags { + let roast = match rebuilder.roast_repo.get(bwr.bag.roast_id).await { + Ok(r) => r, + Err(err) => { + warn!(error = %err, id = %bwr.bag.id, "failed to fetch roast for bag timeline rebuild"); + continue; + } + }; + let roaster = match rebuilder.roaster_repo.get(roast.roaster_id).await { + Ok(r) => r, + Err(err) => { + warn!(error = %err, id = %bwr.bag.id, "failed to fetch roaster for bag timeline rebuild"); + continue; + } + }; + if let Err(err) = rebuilder + .timeline_repo + .insert(bag_timeline_event(&bwr.bag, "added", &roast, &roaster)) + .await + { + warn!(error = %err, id = %bwr.bag.id, "failed to rebuild bag 'added' timeline event"); + } + if bwr.bag.closed { + let mut finished_event = bag_timeline_event(&bwr.bag, "finished", &roast, &roaster); + if let Some(finished_at) = bwr.bag.finished_at { + finished_event.occurred_at = finished_at.and_time(chrono::NaiveTime::MIN).and_utc(); + } + if let Err(err) = rebuilder.timeline_repo.insert(finished_event).await { + warn!(error = %err, id = %bwr.bag.id, "failed to rebuild bag 'finished' timeline event"); + } + } + } + Ok(()) +} + +/// Helper to create a show-all `ListRequest` for any `SortKey` type. +fn show_all_request() -> crate::domain::listing::ListRequest +{ + let sort_key = S::default(); + crate::domain::listing::ListRequest::show_all(sort_key, sort_key.default_direction()) +} diff --git a/src/application/state.rs b/src/application/state.rs index 42726ef..f0df2e1 100644 --- a/src/application/state.rs +++ b/src/application/state.rs @@ -4,7 +4,7 @@ use webauthn_rs::prelude::*; use crate::application::services::{ BagService, BrewService, CafeService, CupService, GearService, RoastService, RoasterService, - StatsInvalidator, + StatsInvalidator, TimelineInvalidator, }; use crate::domain::repositories::{ AiUsageRepository, BagRepository, BrewRepository, CafeRepository, CupRepository, @@ -44,6 +44,7 @@ pub struct AppStateConfig { pub openrouter_api_key: String, pub openrouter_model: String, pub stats_invalidator: StatsInvalidator, + pub timeline_invalidator: TimelineInvalidator, } #[derive(Clone)] @@ -82,6 +83,7 @@ pub struct AppState { pub cup_service: CupService, pub insecure_cookies: bool, pub stats_invalidator: StatsInvalidator, + pub timeline_invalidator: TimelineInvalidator, pub image_semaphore: Arc, } @@ -173,6 +175,7 @@ impl AppState { cup_service, insecure_cookies: config.insecure_cookies, stats_invalidator: config.stats_invalidator, + timeline_invalidator: config.timeline_invalidator, image_semaphore: Arc::new(tokio::sync::Semaphore::new(4)), } } diff --git a/src/domain/coffee/brews.rs b/src/domain/coffee/brews.rs index 685ba29..e62e7ef 100644 --- a/src/domain/coffee/brews.rs +++ b/src/domain/coffee/brews.rs @@ -256,6 +256,7 @@ pub struct UpdateBrew { #[derive(Debug, Default, Clone)] pub struct BrewFilter { pub bag_id: Option, + pub gear_id: Option, } impl BrewFilter { @@ -268,6 +269,15 @@ impl BrewFilter { pub fn for_bag(bag_id: BagId) -> Self { Self { bag_id: Some(bag_id), + ..Self::default() + } + } + + /// Filter for brews using a specific gear piece (grinder, brewer, or filter paper). + pub fn for_gear(gear_id: GearId) -> Self { + Self { + gear_id: Some(gear_id), + ..Self::default() } } } diff --git a/src/domain/repositories.rs b/src/domain/repositories.rs index 8bd8a7e..8aa98a1 100644 --- a/src/domain/repositories.rs +++ b/src/domain/repositories.rs @@ -99,6 +99,21 @@ pub trait TimelineEventRepository: Send + Sync { request: &ListRequest, ) -> Result, RepositoryError>; + async fn update_by_entity( + &self, + entity_type: EntityType, + entity_id: i64, + event: NewTimelineEvent, + ) -> Result<(), RepositoryError>; + + async fn delete_by_entity( + &self, + entity_type: EntityType, + entity_id: i64, + ) -> Result<(), RepositoryError>; + + async fn delete_all(&self) -> Result<(), RepositoryError>; + async fn list_all(&self) -> Result, RepositoryError> { let sort_key = ::default(); let request = @@ -150,6 +165,13 @@ pub trait BagRepository: Send + Sync { ) -> Result, RepositoryError>; async fn update(&self, id: BagId, changes: UpdateBag) -> Result; async fn delete(&self, id: BagId) -> Result<(), RepositoryError>; + + async fn list_all(&self) -> Result, RepositoryError> { + let sort_key = ::default(); + let request = ListRequest::::show_all(sort_key, sort_key.default_direction()); + let page = self.list(BagFilter::default(), &request, None).await?; + Ok(page.items) + } } #[async_trait] @@ -164,6 +186,13 @@ pub trait GearRepository: Send + Sync { ) -> Result, RepositoryError>; async fn update(&self, id: GearId, changes: UpdateGear) -> Result; async fn delete(&self, id: GearId) -> Result<(), RepositoryError>; + + async fn list_all(&self) -> Result, RepositoryError> { + let sort_key = ::default(); + let request = ListRequest::::show_all(sort_key, sort_key.default_direction()); + let page = self.list(GearFilter::default(), &request, None).await?; + Ok(page.items) + } } #[async_trait] @@ -181,6 +210,13 @@ pub trait BrewRepository: Send + Sync { ) -> Result, RepositoryError>; async fn update(&self, id: BrewId, changes: UpdateBrew) -> Result; async fn delete(&self, id: BrewId) -> Result<(), RepositoryError>; + + async fn list_all(&self) -> Result, RepositoryError> { + let sort_key = ::default(); + let request = ListRequest::::show_all(sort_key, sort_key.default_direction()); + let page = self.list(BrewFilter::default(), &request, None).await?; + Ok(page.items) + } } #[async_trait] @@ -227,6 +263,13 @@ pub trait CupRepository: Send + Sync { ) -> Result, RepositoryError>; async fn update(&self, id: CupId, changes: UpdateCup) -> Result; async fn delete(&self, id: CupId) -> Result<(), RepositoryError>; + + async fn list_all(&self) -> Result, RepositoryError> { + let sort_key = ::default(); + let request = ListRequest::::show_all(sort_key, sort_key.default_direction()); + let page = self.list(CupFilter::default(), &request, None).await?; + Ok(page.items) + } } #[async_trait] diff --git a/src/infrastructure/client/mod.rs b/src/infrastructure/client/mod.rs index 7db87b0..4babba5 100644 --- a/src/infrastructure/client/mod.rs +++ b/src/infrastructure/client/mod.rs @@ -6,6 +6,7 @@ pub mod cups; pub mod gear; pub mod roasters; pub mod roasts; +pub mod timeline; pub mod tokens; use anyhow::{Context, Result, anyhow}; @@ -81,6 +82,10 @@ impl BrewlogClient { cups::CupsClient::new(self) } + pub fn timeline(&self) -> timeline::TimelineClient<'_> { + timeline::TimelineClient::new(self) + } + pub(crate) fn endpoint(&self, path: &str) -> Result { self.base_url .join(path) diff --git a/src/infrastructure/client/timeline.rs b/src/infrastructure/client/timeline.rs new file mode 100644 index 0000000..3e2252d --- /dev/null +++ b/src/infrastructure/client/timeline.rs @@ -0,0 +1,29 @@ +use anyhow::{Context, Result}; +use reqwest::StatusCode; + +use super::BrewlogClient; + +pub struct TimelineClient<'a> { + inner: &'a BrewlogClient, +} + +impl<'a> TimelineClient<'a> { + pub(crate) fn new(inner: &'a BrewlogClient) -> Self { + Self { inner } + } + + pub async fn rebuild(&self) -> Result<()> { + let url = self.inner.endpoint("api/v1/timeline/rebuild")?; + let response = self + .inner + .request(reqwest::Method::POST, url) + .send() + .await + .context("failed to issue timeline rebuild request")?; + + match response.status() { + StatusCode::NO_CONTENT => Ok(()), + _ => Err(self.inner.response_error(response).await), + } + } +} diff --git a/src/infrastructure/repositories/analytics/timeline_events.rs b/src/infrastructure/repositories/analytics/timeline_events.rs index f5552a9..4686780 100644 --- a/src/infrastructure/repositories/analytics/timeline_events.rs +++ b/src/infrastructure/repositories/analytics/timeline_events.rs @@ -68,6 +68,74 @@ impl TimelineEventRepository for SqlTimelineEventRepository { record.into_domain() } + async fn update_by_entity( + &self, + entity_type: EntityType, + entity_id: i64, + event: NewTimelineEvent, + ) -> Result<(), RepositoryError> { + let details_json = serde_json::to_string(&event.details).map_err(|err| { + RepositoryError::unexpected(format!("failed to encode timeline event details: {err}")) + })?; + + let tasting_notes_json = serde_json::to_string(&event.tasting_notes).map_err(|err| { + RepositoryError::unexpected(format!( + "failed to encode timeline event tasting notes: {err}" + )) + })?; + + let brew_data_json = event + .brew_data + .as_ref() + .map(serde_json::to_string) + .transpose() + .map_err(|err| { + RepositoryError::unexpected(format!("failed to encode brew data: {err}")) + })?; + + sqlx::query( + r"UPDATE timeline_events + SET title = ?, details_json = ?, tasting_notes_json = ?, + slug = ?, roaster_slug = ?, brew_data_json = ? + WHERE entity_type = ? AND entity_id = ?", + ) + .bind(event.title) + .bind(details_json) + .bind(tasting_notes_json) + .bind(event.slug) + .bind(event.roaster_slug) + .bind(brew_data_json) + .bind(entity_type.as_str()) + .bind(entity_id) + .execute(&self.pool) + .await + .map_err(|err| RepositoryError::unexpected(err.to_string()))?; + + Ok(()) + } + + async fn delete_by_entity( + &self, + entity_type: EntityType, + entity_id: i64, + ) -> Result<(), RepositoryError> { + sqlx::query("DELETE FROM timeline_events WHERE entity_type = ? AND entity_id = ?") + .bind(entity_type.as_str()) + .bind(entity_id) + .execute(&self.pool) + .await + .map_err(|err| RepositoryError::unexpected(err.to_string()))?; + Ok(()) + } + + async fn delete_all(&self) -> Result<(), RepositoryError> { + sqlx::query("DELETE FROM timeline_events") + .execute(&self.pool) + .await + .map_err(|err| RepositoryError::unexpected(err.to_string()))?; + Ok(()) + } + async fn list( &self, request: &ListRequest, diff --git a/src/infrastructure/repositories/coffee/brews.rs b/src/infrastructure/repositories/coffee/brews.rs index 95f81af..38e3804 100644 --- a/src/infrastructure/repositories/coffee/brews.rs +++ b/src/infrastructure/repositories/coffee/brews.rs @@ -83,10 +83,24 @@ impl SqlBrewRepository { } 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())) + let mut conditions = Vec::new(); + + // SAFETY: Direct interpolation is safe here because IDs are i64 from typed wrappers. + if let Some(bag_id) = filter.bag_id { + conditions.push(format!("br.bag_id = {}", bag_id.into_inner())); + } + if let Some(gear_id) = filter.gear_id { + let id = gear_id.into_inner(); + conditions.push(format!( + "(br.grinder_id = {id} OR br.brewer_id = {id} OR br.filter_paper_id = {id})" + )); + } + + if conditions.is_empty() { + None + } else { + Some(conditions.join(" AND ")) + } } } diff --git a/src/main.rs b/src/main.rs index 9da8be9..1ab0159 100644 --- a/src/main.rs +++ b/src/main.rs @@ -3,7 +3,7 @@ use brewlog::application::{ServerConfig, serve}; use brewlog::infrastructure::backup::BackupData; use brewlog::infrastructure::client::BrewlogClient; use brewlog::presentation::cli::{ - Cli, Commands, ServeCommand, bags, brews, cafes, cups, gear, roasters, roasts, tokens, + Cli, Commands, ServeCommand, bags, brews, cafes, cups, gear, roasters, roasts, timeline, tokens, }; use clap::Parser; use tracing_subscriber::{EnvFilter, layer::SubscriberExt, util::SubscriberInitExt}; @@ -51,6 +51,16 @@ async fn main() -> Result<()> { let client = BrewlogClient::from_base_url(&cli.api_url)?; tokens::run(&client, command).await } + Commands::Timeline { command } => { + let client = BrewlogClient::from_base_url(&cli.api_url)?; + match command { + timeline::TimelineCommands::Rebuild => { + client.timeline().rebuild().await?; + eprintln!("Timeline rebuild initiated."); + Ok(()) + } + } + } Commands::Backup(_cmd) => { let client = BrewlogClient::from_base_url(&cli.api_url)?; let data = client.backup().export().await?; diff --git a/src/presentation/cli/mod.rs b/src/presentation/cli/mod.rs index eab7311..e5be939 100644 --- a/src/presentation/cli/mod.rs +++ b/src/presentation/cli/mod.rs @@ -7,6 +7,7 @@ pub mod gear; mod macros; pub mod roasters; pub mod roasts; +pub mod timeline; pub mod tokens; use std::net::SocketAddr; @@ -22,6 +23,7 @@ use cups::CupCommands; use gear::GearCommands; use roasters::RoasterCommands; use roasts::RoastCommands; +use timeline::TimelineCommands; use tokens::TokenCommands; #[derive(Debug, Parser)] @@ -92,6 +94,12 @@ pub enum Commands { command: TokenCommands, }, + /// Manage timeline events + Timeline { + #[command(subcommand)] + command: TimelineCommands, + }, + /// Back up all coffee data to JSON (stdout) Backup(BackupCommand), diff --git a/src/presentation/cli/timeline.rs b/src/presentation/cli/timeline.rs new file mode 100644 index 0000000..e9ee08e --- /dev/null +++ b/src/presentation/cli/timeline.rs @@ -0,0 +1,7 @@ +use clap::Subcommand; + +#[derive(Debug, Subcommand)] +pub enum TimelineCommands { + /// Rebuild all timeline events from current entity data + Rebuild, +} diff --git a/tests/cli/helpers.rs b/tests/cli/helpers.rs index 6ebaf2c..87743a6 100644 --- a/tests/cli/helpers.rs +++ b/tests/cli/helpers.rs @@ -86,6 +86,7 @@ fn ensure_server_started() -> Result<(String, String), String> { .expect("Failed to connect to test database"); let (stats_tx, _stats_rx) = tokio::sync::mpsc::channel(1); + let (timeline_tx, timeline_rx) = tokio::sync::mpsc::channel(32); let state = AppState::from_database( &database, AppStateConfig { @@ -100,9 +101,31 @@ fn ensure_server_started() -> Result<(String, String), String> { stats_invalidator: brewlog::application::services::StatsInvalidator::new( stats_tx, ), + timeline_invalidator: + brewlog::application::services::TimelineInvalidator::new(timeline_tx), }, ); + // Spawn background timeline rebuild task + use brewlog::application::services::timeline_refresh::{ + TimelineRebuilder, timeline_rebuild_task, + }; + let rebuilder = TimelineRebuilder { + timeline_repo: state.timeline_repo.clone(), + roaster_repo: state.roaster_repo.clone(), + roast_repo: state.roast_repo.clone(), + bag_repo: state.bag_repo.clone(), + brew_repo: state.brew_repo.clone(), + cup_repo: state.cup_repo.clone(), + gear_repo: state.gear_repo.clone(), + cafe_repo: state.cafe_repo.clone(), + }; + tokio::spawn(timeline_rebuild_task( + timeline_rx, + rebuilder, + std::time::Duration::from_millis(100), + )); + let app = app_router(state); #[allow(clippy::expect_used)] diff --git a/tests/cli/main.rs b/tests/cli/main.rs index e19073f..2734a8c 100644 --- a/tests/cli/main.rs +++ b/tests/cli/main.rs @@ -8,4 +8,5 @@ pub mod helpers; pub mod roasters_cli; pub mod roasts_cli; pub mod test_macros; +pub mod timeline_cli; pub mod tokens_cli; diff --git a/tests/cli/timeline_cli.rs b/tests/cli/timeline_cli.rs new file mode 100644 index 0000000..bb8bd4e --- /dev/null +++ b/tests/cli/timeline_cli.rs @@ -0,0 +1,75 @@ +use crate::helpers::{create_roast, create_roaster, create_token, run_brewlog, server_info}; + +#[test] +fn test_editing_roaster_updates_timeline_events_after_rebuild() { + let token = create_token("test-timeline-rebuild"); + let (address, _) = server_info(); + + // Create a roaster and roast + let original_name = "Timeline CLI Roasters"; + let roaster_id = create_roaster(original_name, &token); + let _roast_id = create_roast(&roaster_id, "Timeline CLI Roast", &token); + + // Allow timeline events to be created + std::thread::sleep(std::time::Duration::from_millis(50)); + + // Verify original roaster name appears in timeline HTML + let client = reqwest::blocking::Client::new(); + let body = client + .get(format!("{address}/timeline")) + .send() + .expect("failed to fetch timeline") + .text() + .expect("failed to read body"); + assert!( + body.contains(original_name), + "Expected original roaster name in timeline before edit" + ); + + // Update the roaster name via CLI + let updated_name = "Renamed CLI Roasters"; + let output = run_brewlog( + &[ + "roaster", + "update", + "--id", + &roaster_id, + "--name", + updated_name, + ], + &[("BREWLOG_TOKEN", &token)], + ); + assert!( + output.status.success(), + "roaster update should succeed: {}", + String::from_utf8_lossy(&output.stderr) + ); + + // Trigger timeline rebuild via CLI + let output = run_brewlog(&["timeline", "rebuild"], &[("BREWLOG_TOKEN", &token)]); + assert!( + output.status.success(), + "timeline rebuild should succeed: {}", + String::from_utf8_lossy(&output.stderr) + ); + + // Wait for background rebuild task to process + std::thread::sleep(std::time::Duration::from_millis(500)); + + // Verify both timeline events now show the updated roaster name + let body = client + .get(format!("{address}/timeline")) + .send() + .expect("failed to fetch timeline after rebuild") + .text() + .expect("failed to read body"); + + assert!( + body.contains(updated_name), + "Expected updated roaster name '{updated_name}' in timeline after rebuild, got: {body}" + ); + assert!( + !body.contains(original_name), + "Expected original roaster name '{original_name}' to no longer appear in timeline, got: {body}" + ); +} diff --git a/tests/e2e/main.rs b/tests/e2e/main.rs index 538f0f1..3f07a09 100644 --- a/tests/e2e/main.rs +++ b/tests/e2e/main.rs @@ -13,3 +13,4 @@ mod image_tests; mod navigation_tests; mod roaster_tests; mod scan_tests; +mod timeline_tests; diff --git a/tests/e2e/timeline_tests.rs b/tests/e2e/timeline_tests.rs new file mode 100644 index 0000000..b04a99f --- /dev/null +++ b/tests/e2e/timeline_tests.rs @@ -0,0 +1,74 @@ +use thirtyfour::prelude::*; + +use crate::helpers::auth::authenticate_browser; +use crate::helpers::browser::BrowserSession; +use crate::helpers::forms::{fill_input, submit_visible_form}; +use crate::helpers::server_helpers::{ + create_default_roast, create_default_roaster, spawn_app_with_timeline_sync, +}; +use crate::helpers::wait::{wait_for_url_not_contains, wait_for_visible}; + +#[tokio::test] +async fn editing_roaster_updates_timeline_events_in_browser() { + let app = spawn_app_with_timeline_sync().await; + let roaster = create_default_roaster(&app).await; + let _roast = create_default_roast(&app, roaster.id).await; + + let session = BrowserSession::new(&app.address).await.unwrap(); + authenticate_browser(&session, &app).await.unwrap(); + + // Verify original roaster name appears in the timeline + session.goto("/timeline").await.unwrap(); + wait_for_visible(&session.driver, "#timeline-events") + .await + .unwrap(); + + let body = session.driver.find(By::Css("body")).await.unwrap(); + let text = body.text().await.unwrap(); + assert!( + text.contains("Test Roasters"), + "Expected original roaster name 'Test Roasters' in timeline, got: {text}" + ); + + // Navigate to the roaster edit page and change the name + session + .goto(&format!("/roasters/{}/edit", roaster.id)) + .await + .unwrap(); + wait_for_visible(&session.driver, "input[name='name']") + .await + .unwrap(); + + fill_input(&session.driver, "name", "E2E Renamed Roasters") + .await + .unwrap(); + submit_visible_form(&session.driver).await.unwrap(); + + // Wait for redirect back to the detail page + wait_for_url_not_contains(&session.driver, "/edit") + .await + .unwrap(); + + // Wait for the background timeline rebuild task to process + tokio::time::sleep(std::time::Duration::from_millis(300)).await; + + // Go back to the timeline and verify both events reflect the new name + session.goto("/timeline").await.unwrap(); + wait_for_visible(&session.driver, "#timeline-events") + .await + .unwrap(); + + let body = session.driver.find(By::Css("body")).await.unwrap(); + let text = body.text().await.unwrap(); + + assert!( + text.contains("E2E Renamed Roasters"), + "Expected updated roaster name 'E2E Renamed Roasters' in timeline after edit, got: {text}" + ); + assert!( + !text.contains("Test Roasters"), + "Expected original name 'Test Roasters' to no longer appear in timeline, got: {text}" + ); + + session.quit().await; +} diff --git a/tests/server/helpers.rs b/tests/server/helpers.rs index bbc0595..38e13b0 100644 --- a/tests/server/helpers.rs +++ b/tests/server/helpers.rs @@ -77,7 +77,8 @@ pub async fn spawn_app() -> TestApp { } fn test_state_config() -> AppStateConfig { - let (tx, _rx) = tokio::sync::mpsc::channel(1); + let (stats_tx, _stats_rx) = tokio::sync::mpsc::channel(1); + let (timeline_tx, _timeline_rx) = tokio::sync::mpsc::channel(1); AppStateConfig { webauthn: test_webauthn(), insecure_cookies: true, @@ -86,7 +87,8 @@ fn test_state_config() -> AppStateConfig { openrouter_url: brewlog::infrastructure::ai::OPENROUTER_URL.to_string(), openrouter_api_key: String::new(), openrouter_model: "openrouter/free".to_string(), - stats_invalidator: brewlog::application::services::StatsInvalidator::new(tx), + stats_invalidator: brewlog::application::services::StatsInvalidator::new(stats_tx), + timeline_invalidator: brewlog::application::services::TimelineInvalidator::new(timeline_tx), } } @@ -96,7 +98,13 @@ async fn spawn_app_inner( mock_server: Option, ) -> TestApp { let state = AppState::from_database(&database, config); + spawn_app_inner_from_state(state, mock_server).await +} +async fn spawn_app_inner_from_state( + state: AppState, + mock_server: Option, +) -> TestApp { // Clone repos we need for TestApp before consuming state in the router let roaster_repo = state.roaster_repo.clone(); let roast_repo = state.roast_repo.clone(); @@ -142,6 +150,56 @@ pub async fn spawn_app_with_auth() -> TestApp { add_auth_to_app(app).await } +/// Spawn a test app with the timeline background rebuild task running. +/// Uses a short debounce (50ms) so tests don't have to wait long. +pub async fn spawn_app_with_timeline_sync() -> TestApp { + use brewlog::application::services::timeline_refresh::{ + TimelineRebuilder, timeline_rebuild_task, + }; + + let database = Database::connect("sqlite::memory:") + .await + .expect("Failed to connect to in-memory database"); + + let (stats_tx, _stats_rx) = tokio::sync::mpsc::channel(1); + let (timeline_tx, timeline_rx) = tokio::sync::mpsc::channel(32); + + let config = AppStateConfig { + webauthn: test_webauthn(), + insecure_cookies: true, + foursquare_url: brewlog::infrastructure::foursquare::FOURSQUARE_SEARCH_URL.to_string(), + foursquare_api_key: String::new(), + openrouter_url: brewlog::infrastructure::ai::OPENROUTER_URL.to_string(), + openrouter_api_key: String::new(), + openrouter_model: "openrouter/free".to_string(), + stats_invalidator: brewlog::application::services::StatsInvalidator::new(stats_tx), + timeline_invalidator: brewlog::application::services::TimelineInvalidator::new(timeline_tx), + }; + + let state = AppState::from_database(&database, config); + + // Clone repos for the rebuilder before consuming state + let rebuilder = TimelineRebuilder { + timeline_repo: state.timeline_repo.clone(), + roaster_repo: state.roaster_repo.clone(), + roast_repo: state.roast_repo.clone(), + bag_repo: state.bag_repo.clone(), + brew_repo: state.brew_repo.clone(), + cup_repo: state.cup_repo.clone(), + gear_repo: state.gear_repo.clone(), + cafe_repo: state.cafe_repo.clone(), + }; + + tokio::spawn(timeline_rebuild_task( + timeline_rx, + rebuilder, + std::time::Duration::from_millis(50), + )); + + let app = spawn_app_inner_from_state(state, None).await; + add_auth_to_app(app).await +} + pub async fn spawn_app_with_foursquare_mock() -> TestApp { let mock_server = wiremock::MockServer::start().await; let foursquare_url = format!("{}/places/search", mock_server.uri()); diff --git a/tests/server/timeline.rs b/tests/server/timeline.rs index faf1a91..ae86806 100644 --- a/tests/server/timeline.rs +++ b/tests/server/timeline.rs @@ -1,6 +1,7 @@ use crate::helpers::{ create_cafe_with_payload, create_default_bag, create_default_cafe, create_default_gear, create_default_roast, create_default_roaster, create_roaster_with_payload, spawn_app_with_auth, + spawn_app_with_timeline_sync, }; use brewlog::domain::brews::NewBrew; use brewlog::domain::cafes::NewCafe; @@ -667,3 +668,229 @@ async fn creating_a_brew_surfaces_on_the_timeline() { "Expected roast name to appear in brew timeline event, got: {body}" ); } + +#[tokio::test] +async fn editing_a_roaster_updates_its_timeline_event() { + let app = spawn_app_with_timeline_sync().await; + let client = Client::new(); + + let original_name = "Original Roasters"; + let roaster = create_roaster_with_payload( + &app, + NewRoaster { + name: original_name.to_string(), + country: "UK".to_string(), + city: None, + homepage: None, + created_at: None, + }, + ) + .await; + + sleep(Duration::from_millis(10)).await; + + // Verify original name appears + let body = client + .get(format!("{}/timeline", app.address)) + .send() + .await + .expect("failed to fetch timeline") + .text() + .await + .expect("failed to read body"); + assert!( + body.contains(original_name), + "Expected original name in timeline, got: {body}" + ); + + // Update the roaster name + let updated_name = "Renamed Roasters"; + let response = client + .put(app.api_url(&format!("/roasters/{}", roaster.id))) + .bearer_auth(app.auth_token.as_ref().unwrap()) + .json(&serde_json::json!({ "name": updated_name })) + .send() + .await + .expect("failed to update roaster"); + assert_eq!(response.status(), 200); + + // Wait for the background task to process (debounce + processing) + sleep(Duration::from_millis(200)).await; + + // Verify timeline now shows the updated name + let body = client + .get(format!("{}/timeline", app.address)) + .send() + .await + .expect("failed to fetch timeline after update") + .text() + .await + .expect("failed to read body"); + assert!( + body.contains(updated_name), + "Expected updated name '{updated_name}' in timeline after edit, got: {body}" + ); + assert!( + !body.contains(original_name), + "Expected original name '{original_name}' to be replaced in timeline, got: {body}" + ); +} + +#[tokio::test] +async fn deleting_a_roaster_removes_its_timeline_event() { + let app = spawn_app_with_timeline_sync().await; + let client = Client::new(); + + let roaster_name = "Deletable Roasters"; + let roaster = create_roaster_with_payload( + &app, + NewRoaster { + name: roaster_name.to_string(), + country: "UK".to_string(), + city: None, + homepage: None, + created_at: None, + }, + ) + .await; + + sleep(Duration::from_millis(10)).await; + + // Verify event exists + let body = client + .get(format!("{}/timeline", app.address)) + .send() + .await + .expect("failed to fetch timeline") + .text() + .await + .expect("failed to read body"); + assert!( + body.contains(roaster_name), + "Expected roaster name in timeline before delete, got: {body}" + ); + + // Delete the roaster + let response = client + .delete(app.api_url(&format!("/roasters/{}", roaster.id))) + .bearer_auth(app.auth_token.as_ref().unwrap()) + .send() + .await + .expect("failed to delete roaster"); + assert!( + response.status().is_success(), + "Expected successful delete, got: {}", + response.status() + ); + + sleep(Duration::from_millis(10)).await; + + // Verify timeline no longer shows the roaster + let body = client + .get(format!("{}/timeline", app.address)) + .send() + .await + .expect("failed to fetch timeline after delete") + .text() + .await + .expect("failed to read body"); + assert!( + !body.contains(roaster_name), + "Expected roaster name to be removed from timeline after delete, got: {body}" + ); +} + +#[tokio::test] +async fn timeline_rebuild_endpoint_requires_auth() { + let app = spawn_app_with_auth().await; + let client = Client::new(); + + // Without auth token → 401 + let response = client + .post(app.api_url("/timeline/rebuild")) + .send() + .await + .expect("failed to send rebuild request"); + assert_eq!(response.status(), 401); +} + +#[tokio::test] +async fn timeline_rebuild_endpoint_returns_204() { + let app = spawn_app_with_auth().await; + let client = Client::new(); + + let response = client + .post(app.api_url("/timeline/rebuild")) + .bearer_auth(app.auth_token.as_ref().unwrap()) + .send() + .await + .expect("failed to send rebuild request"); + assert_eq!(response.status(), 204); +} + +#[tokio::test] +async fn editing_a_roaster_cascades_to_roast_timeline_event() { + let app = spawn_app_with_timeline_sync().await; + let client = Client::new(); + + let original_name = "Cascade Roasters"; + let roaster = create_roaster_with_payload( + &app, + NewRoaster { + name: original_name.to_string(), + country: "UK".to_string(), + city: None, + homepage: None, + created_at: None, + }, + ) + .await; + + sleep(Duration::from_millis(5)).await; + let roast_name = "Cascade Roast"; + create_roast(&app, roaster.id, roast_name).await; + + sleep(Duration::from_millis(10)).await; + + // Verify roast timeline event shows original roaster name + let body = client + .get(format!("{}/timeline", app.address)) + .send() + .await + .expect("failed to fetch timeline") + .text() + .await + .expect("failed to read body"); + assert!( + body.contains(original_name), + "Expected original roaster name in roast timeline event" + ); + + // Rename the roaster + let updated_name = "Cascade Renamed"; + let response = client + .put(app.api_url(&format!("/roasters/{}", roaster.id))) + .bearer_auth(app.auth_token.as_ref().unwrap()) + .json(&serde_json::json!({ "name": updated_name })) + .send() + .await + .expect("failed to update roaster"); + assert_eq!(response.status(), 200); + + // Wait for background cascade + sleep(Duration::from_millis(300)).await; + + // Verify roast timeline event now shows the updated roaster name + let body = client + .get(format!("{}/timeline", app.address)) + .send() + .await + .expect("failed to fetch timeline after cascade") + .text() + .await + .expect("failed to read body"); + assert!( + body.contains(updated_name), + "Expected cascaded roaster name '{updated_name}' in roast timeline event, got: {body}" + ); +}