feat: sync timeline events with entity edits
Add background timeline rebuild task (mirroring stats cache pattern) that refreshes denormalized timeline event snapshots when entities are updated. Includes cascade logic so editing a roaster refreshes timeline events for its roasts, bags, brews, and cups. - Add update_by_entity/delete_by_entity/delete_all to TimelineEventRepository - Add TimelineInvalidator with debounced background rebuild task - Add invalidate() calls to all 7 entity update handlers - Add delete_by_entity cleanup to define_delete_handler! macro - Add gear_id filter to BrewFilter for cascade traversal - Add `brewlog timeline rebuild` CLI command for full rebuild - Add 5 integration tests for timeline sync behavior
This commit is contained in:
parent
ae016c6437
commit
0e3847ade9
31 changed files with 1221 additions and 15 deletions
|
|
@ -221,6 +221,9 @@ pub(crate) async fn update_bag(
|
||||||
|
|
||||||
info!(%id, closed = ?update.closed, "bag updated");
|
info!(%id, closed = ?update.closed, "bag updated");
|
||||||
state.stats_invalidator.invalidate();
|
state.stats_invalidator.invalidate();
|
||||||
|
state
|
||||||
|
.timeline_invalidator
|
||||||
|
.invalidate(EntityType::Bag, i64::from(bag.id));
|
||||||
|
|
||||||
save_deferred_image(
|
save_deferred_image(
|
||||||
&state,
|
&state,
|
||||||
|
|
@ -265,7 +268,8 @@ define_delete_handler!(
|
||||||
bag_repo,
|
bag_repo,
|
||||||
render_bag_list_fragment,
|
render_bag_list_fragment,
|
||||||
"type=bags",
|
"type=bags",
|
||||||
"/data?type=bags"
|
"/data?type=bags",
|
||||||
|
entity_type: crate::domain::entity_type::EntityType::Bag
|
||||||
);
|
);
|
||||||
|
|
||||||
#[derive(Debug, Deserialize)]
|
#[derive(Debug, Deserialize)]
|
||||||
|
|
|
||||||
|
|
@ -417,6 +417,9 @@ pub(crate) async fn update_brew(
|
||||||
|
|
||||||
info!(%id, "brew updated");
|
info!(%id, "brew updated");
|
||||||
state.stats_invalidator.invalidate();
|
state.stats_invalidator.invalidate();
|
||||||
|
state
|
||||||
|
.timeline_invalidator
|
||||||
|
.invalidate(EntityType::Brew, i64::from(id));
|
||||||
|
|
||||||
save_deferred_image(
|
save_deferred_image(
|
||||||
&state,
|
&state,
|
||||||
|
|
|
||||||
|
|
@ -204,6 +204,9 @@ pub(crate) async fn update_cafe(
|
||||||
.map_err(AppError::from)?;
|
.map_err(AppError::from)?;
|
||||||
info!(%id, "cafe updated");
|
info!(%id, "cafe updated");
|
||||||
state.stats_invalidator.invalidate();
|
state.stats_invalidator.invalidate();
|
||||||
|
state
|
||||||
|
.timeline_invalidator
|
||||||
|
.invalidate(EntityType::Cafe, i64::from(cafe.id));
|
||||||
|
|
||||||
save_deferred_image(
|
save_deferred_image(
|
||||||
&state,
|
&state,
|
||||||
|
|
|
||||||
|
|
@ -143,6 +143,9 @@ pub(crate) async fn update_cup(
|
||||||
|
|
||||||
info!(%id, "cup updated");
|
info!(%id, "cup updated");
|
||||||
state.stats_invalidator.invalidate();
|
state.stats_invalidator.invalidate();
|
||||||
|
state
|
||||||
|
.timeline_invalidator
|
||||||
|
.invalidate(EntityType::Cup, i64::from(cup.id));
|
||||||
|
|
||||||
save_deferred_image(
|
save_deferred_image(
|
||||||
&state,
|
&state,
|
||||||
|
|
|
||||||
|
|
@ -174,6 +174,9 @@ pub(crate) async fn update_gear(
|
||||||
|
|
||||||
info!(%id, "gear updated");
|
info!(%id, "gear updated");
|
||||||
state.stats_invalidator.invalidate();
|
state.stats_invalidator.invalidate();
|
||||||
|
state
|
||||||
|
.timeline_invalidator
|
||||||
|
.invalidate(EntityType::Gear, i64::from(gear.id));
|
||||||
|
|
||||||
save_deferred_image(
|
save_deferred_image(
|
||||||
&state,
|
&state,
|
||||||
|
|
|
||||||
|
|
@ -195,6 +195,9 @@ pub(crate) async fn update_roaster(
|
||||||
.map_err(AppError::from)?;
|
.map_err(AppError::from)?;
|
||||||
info!(%id, "roaster updated");
|
info!(%id, "roaster updated");
|
||||||
state.stats_invalidator.invalidate();
|
state.stats_invalidator.invalidate();
|
||||||
|
state
|
||||||
|
.timeline_invalidator
|
||||||
|
.invalidate(EntityType::Roaster, i64::from(roaster.id));
|
||||||
|
|
||||||
save_deferred_image(
|
save_deferred_image(
|
||||||
&state,
|
&state,
|
||||||
|
|
|
||||||
|
|
@ -249,6 +249,9 @@ pub(crate) async fn update_roast(
|
||||||
|
|
||||||
info!(%id, "roast updated");
|
info!(%id, "roast updated");
|
||||||
state.stats_invalidator.invalidate();
|
state.stats_invalidator.invalidate();
|
||||||
|
state
|
||||||
|
.timeline_invalidator
|
||||||
|
.invalidate(EntityType::Roast, i64::from(id));
|
||||||
|
|
||||||
save_deferred_image(
|
save_deferred_image(
|
||||||
&state,
|
&state,
|
||||||
|
|
|
||||||
|
|
@ -86,13 +86,13 @@ macro_rules! define_enriched_get_handler {
|
||||||
/// );
|
/// );
|
||||||
/// ```
|
/// ```
|
||||||
macro_rules! define_delete_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) => {
|
($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, None);
|
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) => {
|
($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))]
|
#[tracing::instrument(skip(state, _auth_user, headers, query))]
|
||||||
pub(crate) async fn $fn_name(
|
pub(crate) async fn $fn_name(
|
||||||
axum::extract::State(state): axum::extract::State<crate::application::state::AppState>,
|
axum::extract::State(state): axum::extract::State<crate::application::state::AppState>,
|
||||||
|
|
@ -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");
|
tracing::info!(%id, "entity deleted");
|
||||||
state.stats_invalidator.invalidate();
|
state.stats_invalidator.invalidate();
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -9,7 +9,7 @@ pub(crate) mod system;
|
||||||
pub(crate) use analytics::stats;
|
pub(crate) use analytics::stats;
|
||||||
pub(crate) use auth::{tokens, webauthn};
|
pub(crate) use auth::{tokens, webauthn};
|
||||||
pub(crate) use coffee::{bags, brews, cafes, checkin, cups, gear, roasters, roasts, scan};
|
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::extract::DefaultBodyLimit;
|
||||||
use axum::routing::{get, post};
|
use axum::routing::{get, post};
|
||||||
|
|
@ -103,6 +103,7 @@ pub(super) fn router() -> axum::Router<AppState> {
|
||||||
)
|
)
|
||||||
.route("/backup/reset", post(backup::reset_database))
|
.route("/backup/reset", post(backup::reset_database))
|
||||||
.route("/stats/recompute", post(stats::recompute_stats))
|
.route("/stats/recompute", post(stats::recompute_stats))
|
||||||
|
.route("/timeline/rebuild", post(timeline::rebuild_timeline))
|
||||||
.route(
|
.route(
|
||||||
"/{entity_type}/{id}/image",
|
"/{entity_type}/{id}/image",
|
||||||
get(images::get_image)
|
get(images::get_image)
|
||||||
|
|
|
||||||
|
|
@ -1,2 +1,3 @@
|
||||||
pub(crate) mod admin;
|
pub(crate) mod admin;
|
||||||
pub(crate) mod backup;
|
pub(crate) mod backup;
|
||||||
|
pub(crate) mod timeline;
|
||||||
|
|
|
||||||
18
src/application/routes/api/system/timeline.rs
Normal file
18
src/application/routes/api/system/timeline.rs
Normal file
|
|
@ -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<AppState>,
|
||||||
|
_auth_user: AuthenticatedUser,
|
||||||
|
) -> Result<Response, ApiError> {
|
||||||
|
info!("timeline rebuild requested");
|
||||||
|
state.timeline_invalidator.rebuild_all();
|
||||||
|
Ok(StatusCode::NO_CONTENT.into_response())
|
||||||
|
}
|
||||||
|
|
@ -9,8 +9,9 @@ use tracing::info;
|
||||||
use webauthn_rs::prelude::*;
|
use webauthn_rs::prelude::*;
|
||||||
|
|
||||||
use crate::application::routes::app_router;
|
use crate::application::routes::app_router;
|
||||||
use crate::application::services::StatsInvalidator;
|
|
||||||
use crate::application::services::stats::stats_recomputation_task;
|
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::application::state::{AppState, AppStateConfig};
|
||||||
use crate::domain::registration_tokens::NewRegistrationToken;
|
use crate::domain::registration_tokens::NewRegistrationToken;
|
||||||
use crate::domain::repositories::{RegistrationTokenRepository, UserRepository};
|
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_tx, stats_rx) = tokio::sync::mpsc::channel::<()>(32);
|
||||||
let stats_invalidator = StatsInvalidator::new(stats_tx);
|
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(
|
let state = AppState::from_database(
|
||||||
&database,
|
&database,
|
||||||
AppStateConfig {
|
AppStateConfig {
|
||||||
|
|
@ -56,6 +62,7 @@ pub async fn serve(config: ServerConfig) -> anyhow::Result<()> {
|
||||||
openrouter_api_key: config.openrouter_api_key,
|
openrouter_api_key: config.openrouter_api_key,
|
||||||
openrouter_model: config.openrouter_model,
|
openrouter_model: config.openrouter_model,
|
||||||
stats_invalidator: stats_invalidator.clone(),
|
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),
|
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
|
// Seed the stats cache on startup
|
||||||
stats_invalidator.invalidate();
|
stats_invalidator.invalidate();
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -3,12 +3,14 @@ mod brews;
|
||||||
mod cups;
|
mod cups;
|
||||||
mod roasts;
|
mod roasts;
|
||||||
pub mod stats;
|
pub mod stats;
|
||||||
|
pub mod timeline_refresh;
|
||||||
|
|
||||||
pub use bags::BagService;
|
pub use bags::BagService;
|
||||||
pub use brews::BrewService;
|
pub use brews::BrewService;
|
||||||
pub use cups::CupService;
|
pub use cups::CupService;
|
||||||
pub use roasts::RoastService;
|
pub use roasts::RoastService;
|
||||||
pub use stats::StatsInvalidator;
|
pub use stats::StatsInvalidator;
|
||||||
|
pub use timeline_refresh::TimelineInvalidator;
|
||||||
|
|
||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
|
|
||||||
|
|
|
||||||
478
src/application/services/timeline_refresh.rs
Normal file
478
src/application/services/timeline_refresh.rs
Normal file
|
|
@ -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<TimelineInvalidation>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl TimelineInvalidator {
|
||||||
|
pub fn new(tx: mpsc::Sender<TimelineInvalidation>) -> 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<dyn TimelineEventRepository>,
|
||||||
|
pub roaster_repo: Arc<dyn RoasterRepository>,
|
||||||
|
pub roast_repo: Arc<dyn RoastRepository>,
|
||||||
|
pub bag_repo: Arc<dyn BagRepository>,
|
||||||
|
pub brew_repo: Arc<dyn BrewRepository>,
|
||||||
|
pub cup_repo: Arc<dyn CupRepository>,
|
||||||
|
pub gear_repo: Arc<dyn GearRepository>,
|
||||||
|
pub cafe_repo: Arc<dyn CafeRepository>,
|
||||||
|
}
|
||||||
|
|
||||||
|
/// 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<TimelineInvalidation>,
|
||||||
|
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::<crate::domain::brews::BrewSortKey>(),
|
||||||
|
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::<crate::domain::cups::CupSortKey>(),
|
||||||
|
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::<crate::domain::bags::BagSortKey>(),
|
||||||
|
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::<crate::domain::cups::CupSortKey>(),
|
||||||
|
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::<crate::domain::brews::BrewSortKey>(),
|
||||||
|
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<S: crate::domain::listing::SortKey>() -> crate::domain::listing::ListRequest<S>
|
||||||
|
{
|
||||||
|
let sort_key = S::default();
|
||||||
|
crate::domain::listing::ListRequest::show_all(sort_key, sort_key.default_direction())
|
||||||
|
}
|
||||||
|
|
@ -4,7 +4,7 @@ use webauthn_rs::prelude::*;
|
||||||
|
|
||||||
use crate::application::services::{
|
use crate::application::services::{
|
||||||
BagService, BrewService, CafeService, CupService, GearService, RoastService, RoasterService,
|
BagService, BrewService, CafeService, CupService, GearService, RoastService, RoasterService,
|
||||||
StatsInvalidator,
|
StatsInvalidator, TimelineInvalidator,
|
||||||
};
|
};
|
||||||
use crate::domain::repositories::{
|
use crate::domain::repositories::{
|
||||||
AiUsageRepository, BagRepository, BrewRepository, CafeRepository, CupRepository,
|
AiUsageRepository, BagRepository, BrewRepository, CafeRepository, CupRepository,
|
||||||
|
|
@ -44,6 +44,7 @@ pub struct AppStateConfig {
|
||||||
pub openrouter_api_key: String,
|
pub openrouter_api_key: String,
|
||||||
pub openrouter_model: String,
|
pub openrouter_model: String,
|
||||||
pub stats_invalidator: StatsInvalidator,
|
pub stats_invalidator: StatsInvalidator,
|
||||||
|
pub timeline_invalidator: TimelineInvalidator,
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Clone)]
|
#[derive(Clone)]
|
||||||
|
|
@ -82,6 +83,7 @@ pub struct AppState {
|
||||||
pub cup_service: CupService,
|
pub cup_service: CupService,
|
||||||
pub insecure_cookies: bool,
|
pub insecure_cookies: bool,
|
||||||
pub stats_invalidator: StatsInvalidator,
|
pub stats_invalidator: StatsInvalidator,
|
||||||
|
pub timeline_invalidator: TimelineInvalidator,
|
||||||
pub image_semaphore: Arc<tokio::sync::Semaphore>,
|
pub image_semaphore: Arc<tokio::sync::Semaphore>,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -173,6 +175,7 @@ impl AppState {
|
||||||
cup_service,
|
cup_service,
|
||||||
insecure_cookies: config.insecure_cookies,
|
insecure_cookies: config.insecure_cookies,
|
||||||
stats_invalidator: config.stats_invalidator,
|
stats_invalidator: config.stats_invalidator,
|
||||||
|
timeline_invalidator: config.timeline_invalidator,
|
||||||
image_semaphore: Arc::new(tokio::sync::Semaphore::new(4)),
|
image_semaphore: Arc::new(tokio::sync::Semaphore::new(4)),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -256,6 +256,7 @@ pub struct UpdateBrew {
|
||||||
#[derive(Debug, Default, Clone)]
|
#[derive(Debug, Default, Clone)]
|
||||||
pub struct BrewFilter {
|
pub struct BrewFilter {
|
||||||
pub bag_id: Option<BagId>,
|
pub bag_id: Option<BagId>,
|
||||||
|
pub gear_id: Option<GearId>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl BrewFilter {
|
impl BrewFilter {
|
||||||
|
|
@ -268,6 +269,15 @@ impl BrewFilter {
|
||||||
pub fn for_bag(bag_id: BagId) -> Self {
|
pub fn for_bag(bag_id: BagId) -> Self {
|
||||||
Self {
|
Self {
|
||||||
bag_id: Some(bag_id),
|
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()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -99,6 +99,21 @@ pub trait TimelineEventRepository: Send + Sync {
|
||||||
request: &ListRequest<TimelineSortKey>,
|
request: &ListRequest<TimelineSortKey>,
|
||||||
) -> Result<Page<TimelineEvent>, RepositoryError>;
|
) -> Result<Page<TimelineEvent>, 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<Vec<TimelineEvent>, RepositoryError> {
|
async fn list_all(&self) -> Result<Vec<TimelineEvent>, RepositoryError> {
|
||||||
let sort_key = <TimelineSortKey as SortKey>::default();
|
let sort_key = <TimelineSortKey as SortKey>::default();
|
||||||
let request =
|
let request =
|
||||||
|
|
@ -150,6 +165,13 @@ pub trait BagRepository: Send + Sync {
|
||||||
) -> Result<Page<BagWithRoast>, RepositoryError>;
|
) -> Result<Page<BagWithRoast>, RepositoryError>;
|
||||||
async fn update(&self, id: BagId, changes: UpdateBag) -> Result<Bag, RepositoryError>;
|
async fn update(&self, id: BagId, changes: UpdateBag) -> Result<Bag, RepositoryError>;
|
||||||
async fn delete(&self, id: BagId) -> Result<(), RepositoryError>;
|
async fn delete(&self, id: BagId) -> Result<(), RepositoryError>;
|
||||||
|
|
||||||
|
async fn list_all(&self) -> Result<Vec<BagWithRoast>, RepositoryError> {
|
||||||
|
let sort_key = <BagSortKey as SortKey>::default();
|
||||||
|
let request = ListRequest::<BagSortKey>::show_all(sort_key, sort_key.default_direction());
|
||||||
|
let page = self.list(BagFilter::default(), &request, None).await?;
|
||||||
|
Ok(page.items)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
#[async_trait]
|
#[async_trait]
|
||||||
|
|
@ -164,6 +186,13 @@ pub trait GearRepository: Send + Sync {
|
||||||
) -> Result<Page<Gear>, RepositoryError>;
|
) -> Result<Page<Gear>, RepositoryError>;
|
||||||
async fn update(&self, id: GearId, changes: UpdateGear) -> Result<Gear, RepositoryError>;
|
async fn update(&self, id: GearId, changes: UpdateGear) -> Result<Gear, RepositoryError>;
|
||||||
async fn delete(&self, id: GearId) -> Result<(), RepositoryError>;
|
async fn delete(&self, id: GearId) -> Result<(), RepositoryError>;
|
||||||
|
|
||||||
|
async fn list_all(&self) -> Result<Vec<Gear>, RepositoryError> {
|
||||||
|
let sort_key = <GearSortKey as SortKey>::default();
|
||||||
|
let request = ListRequest::<GearSortKey>::show_all(sort_key, sort_key.default_direction());
|
||||||
|
let page = self.list(GearFilter::default(), &request, None).await?;
|
||||||
|
Ok(page.items)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
#[async_trait]
|
#[async_trait]
|
||||||
|
|
@ -181,6 +210,13 @@ pub trait BrewRepository: Send + Sync {
|
||||||
) -> Result<Page<BrewWithDetails>, RepositoryError>;
|
) -> Result<Page<BrewWithDetails>, RepositoryError>;
|
||||||
async fn update(&self, id: BrewId, changes: UpdateBrew) -> Result<Brew, RepositoryError>;
|
async fn update(&self, id: BrewId, changes: UpdateBrew) -> Result<Brew, RepositoryError>;
|
||||||
async fn delete(&self, id: BrewId) -> Result<(), RepositoryError>;
|
async fn delete(&self, id: BrewId) -> Result<(), RepositoryError>;
|
||||||
|
|
||||||
|
async fn list_all(&self) -> Result<Vec<BrewWithDetails>, RepositoryError> {
|
||||||
|
let sort_key = <BrewSortKey as SortKey>::default();
|
||||||
|
let request = ListRequest::<BrewSortKey>::show_all(sort_key, sort_key.default_direction());
|
||||||
|
let page = self.list(BrewFilter::default(), &request, None).await?;
|
||||||
|
Ok(page.items)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
#[async_trait]
|
#[async_trait]
|
||||||
|
|
@ -227,6 +263,13 @@ pub trait CupRepository: Send + Sync {
|
||||||
) -> Result<Page<CupWithDetails>, RepositoryError>;
|
) -> Result<Page<CupWithDetails>, RepositoryError>;
|
||||||
async fn update(&self, id: CupId, changes: UpdateCup) -> Result<Cup, RepositoryError>;
|
async fn update(&self, id: CupId, changes: UpdateCup) -> Result<Cup, RepositoryError>;
|
||||||
async fn delete(&self, id: CupId) -> Result<(), RepositoryError>;
|
async fn delete(&self, id: CupId) -> Result<(), RepositoryError>;
|
||||||
|
|
||||||
|
async fn list_all(&self) -> Result<Vec<CupWithDetails>, RepositoryError> {
|
||||||
|
let sort_key = <CupSortKey as SortKey>::default();
|
||||||
|
let request = ListRequest::<CupSortKey>::show_all(sort_key, sort_key.default_direction());
|
||||||
|
let page = self.list(CupFilter::default(), &request, None).await?;
|
||||||
|
Ok(page.items)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
#[async_trait]
|
#[async_trait]
|
||||||
|
|
|
||||||
|
|
@ -6,6 +6,7 @@ pub mod cups;
|
||||||
pub mod gear;
|
pub mod gear;
|
||||||
pub mod roasters;
|
pub mod roasters;
|
||||||
pub mod roasts;
|
pub mod roasts;
|
||||||
|
pub mod timeline;
|
||||||
pub mod tokens;
|
pub mod tokens;
|
||||||
|
|
||||||
use anyhow::{Context, Result, anyhow};
|
use anyhow::{Context, Result, anyhow};
|
||||||
|
|
@ -81,6 +82,10 @@ impl BrewlogClient {
|
||||||
cups::CupsClient::new(self)
|
cups::CupsClient::new(self)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub fn timeline(&self) -> timeline::TimelineClient<'_> {
|
||||||
|
timeline::TimelineClient::new(self)
|
||||||
|
}
|
||||||
|
|
||||||
pub(crate) fn endpoint(&self, path: &str) -> Result<Url> {
|
pub(crate) fn endpoint(&self, path: &str) -> Result<Url> {
|
||||||
self.base_url
|
self.base_url
|
||||||
.join(path)
|
.join(path)
|
||||||
|
|
|
||||||
29
src/infrastructure/client/timeline.rs
Normal file
29
src/infrastructure/client/timeline.rs
Normal file
|
|
@ -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),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -68,6 +68,74 @@ impl TimelineEventRepository for SqlTimelineEventRepository {
|
||||||
record.into_domain()
|
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(
|
async fn list(
|
||||||
&self,
|
&self,
|
||||||
request: &ListRequest<TimelineSortKey>,
|
request: &ListRequest<TimelineSortKey>,
|
||||||
|
|
|
||||||
|
|
@ -83,10 +83,24 @@ impl SqlBrewRepository {
|
||||||
}
|
}
|
||||||
|
|
||||||
fn build_where_clause(filter: &BrewFilter) -> Option<String> {
|
fn build_where_clause(filter: &BrewFilter) -> Option<String> {
|
||||||
// SAFETY: Direct interpolation is safe here because `bag_id` is an i64 from a typed wrapper.
|
let mut conditions = Vec::new();
|
||||||
filter
|
|
||||||
.bag_id
|
// SAFETY: Direct interpolation is safe here because IDs are i64 from typed wrappers.
|
||||||
.map(|bag_id| format!("br.bag_id = {}", bag_id.into_inner()))
|
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 "))
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
12
src/main.rs
12
src/main.rs
|
|
@ -3,7 +3,7 @@ use brewlog::application::{ServerConfig, serve};
|
||||||
use brewlog::infrastructure::backup::BackupData;
|
use brewlog::infrastructure::backup::BackupData;
|
||||||
use brewlog::infrastructure::client::BrewlogClient;
|
use brewlog::infrastructure::client::BrewlogClient;
|
||||||
use brewlog::presentation::cli::{
|
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 clap::Parser;
|
||||||
use tracing_subscriber::{EnvFilter, layer::SubscriberExt, util::SubscriberInitExt};
|
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)?;
|
let client = BrewlogClient::from_base_url(&cli.api_url)?;
|
||||||
tokens::run(&client, command).await
|
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) => {
|
Commands::Backup(_cmd) => {
|
||||||
let client = BrewlogClient::from_base_url(&cli.api_url)?;
|
let client = BrewlogClient::from_base_url(&cli.api_url)?;
|
||||||
let data = client.backup().export().await?;
|
let data = client.backup().export().await?;
|
||||||
|
|
|
||||||
|
|
@ -7,6 +7,7 @@ pub mod gear;
|
||||||
mod macros;
|
mod macros;
|
||||||
pub mod roasters;
|
pub mod roasters;
|
||||||
pub mod roasts;
|
pub mod roasts;
|
||||||
|
pub mod timeline;
|
||||||
pub mod tokens;
|
pub mod tokens;
|
||||||
|
|
||||||
use std::net::SocketAddr;
|
use std::net::SocketAddr;
|
||||||
|
|
@ -22,6 +23,7 @@ use cups::CupCommands;
|
||||||
use gear::GearCommands;
|
use gear::GearCommands;
|
||||||
use roasters::RoasterCommands;
|
use roasters::RoasterCommands;
|
||||||
use roasts::RoastCommands;
|
use roasts::RoastCommands;
|
||||||
|
use timeline::TimelineCommands;
|
||||||
use tokens::TokenCommands;
|
use tokens::TokenCommands;
|
||||||
|
|
||||||
#[derive(Debug, Parser)]
|
#[derive(Debug, Parser)]
|
||||||
|
|
@ -92,6 +94,12 @@ pub enum Commands {
|
||||||
command: TokenCommands,
|
command: TokenCommands,
|
||||||
},
|
},
|
||||||
|
|
||||||
|
/// Manage timeline events
|
||||||
|
Timeline {
|
||||||
|
#[command(subcommand)]
|
||||||
|
command: TimelineCommands,
|
||||||
|
},
|
||||||
|
|
||||||
/// Back up all coffee data to JSON (stdout)
|
/// Back up all coffee data to JSON (stdout)
|
||||||
Backup(BackupCommand),
|
Backup(BackupCommand),
|
||||||
|
|
||||||
|
|
|
||||||
7
src/presentation/cli/timeline.rs
Normal file
7
src/presentation/cli/timeline.rs
Normal file
|
|
@ -0,0 +1,7 @@
|
||||||
|
use clap::Subcommand;
|
||||||
|
|
||||||
|
#[derive(Debug, Subcommand)]
|
||||||
|
pub enum TimelineCommands {
|
||||||
|
/// Rebuild all timeline events from current entity data
|
||||||
|
Rebuild,
|
||||||
|
}
|
||||||
|
|
@ -86,6 +86,7 @@ fn ensure_server_started() -> Result<(String, String), String> {
|
||||||
.expect("Failed to connect to test database");
|
.expect("Failed to connect to test database");
|
||||||
|
|
||||||
let (stats_tx, _stats_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(32);
|
||||||
let state = AppState::from_database(
|
let state = AppState::from_database(
|
||||||
&database,
|
&database,
|
||||||
AppStateConfig {
|
AppStateConfig {
|
||||||
|
|
@ -100,9 +101,31 @@ fn ensure_server_started() -> Result<(String, String), String> {
|
||||||
stats_invalidator: brewlog::application::services::StatsInvalidator::new(
|
stats_invalidator: brewlog::application::services::StatsInvalidator::new(
|
||||||
stats_tx,
|
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);
|
let app = app_router(state);
|
||||||
|
|
||||||
#[allow(clippy::expect_used)]
|
#[allow(clippy::expect_used)]
|
||||||
|
|
|
||||||
|
|
@ -8,4 +8,5 @@ pub mod helpers;
|
||||||
pub mod roasters_cli;
|
pub mod roasters_cli;
|
||||||
pub mod roasts_cli;
|
pub mod roasts_cli;
|
||||||
pub mod test_macros;
|
pub mod test_macros;
|
||||||
|
pub mod timeline_cli;
|
||||||
pub mod tokens_cli;
|
pub mod tokens_cli;
|
||||||
|
|
|
||||||
75
tests/cli/timeline_cli.rs
Normal file
75
tests/cli/timeline_cli.rs
Normal file
|
|
@ -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}"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
@ -13,3 +13,4 @@ mod image_tests;
|
||||||
mod navigation_tests;
|
mod navigation_tests;
|
||||||
mod roaster_tests;
|
mod roaster_tests;
|
||||||
mod scan_tests;
|
mod scan_tests;
|
||||||
|
mod timeline_tests;
|
||||||
|
|
|
||||||
74
tests/e2e/timeline_tests.rs
Normal file
74
tests/e2e/timeline_tests.rs
Normal file
|
|
@ -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;
|
||||||
|
}
|
||||||
|
|
@ -77,7 +77,8 @@ pub async fn spawn_app() -> TestApp {
|
||||||
}
|
}
|
||||||
|
|
||||||
fn test_state_config() -> AppStateConfig {
|
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 {
|
AppStateConfig {
|
||||||
webauthn: test_webauthn(),
|
webauthn: test_webauthn(),
|
||||||
insecure_cookies: true,
|
insecure_cookies: true,
|
||||||
|
|
@ -86,7 +87,8 @@ fn test_state_config() -> AppStateConfig {
|
||||||
openrouter_url: brewlog::infrastructure::ai::OPENROUTER_URL.to_string(),
|
openrouter_url: brewlog::infrastructure::ai::OPENROUTER_URL.to_string(),
|
||||||
openrouter_api_key: String::new(),
|
openrouter_api_key: String::new(),
|
||||||
openrouter_model: "openrouter/free".to_string(),
|
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<wiremock::MockServer>,
|
mock_server: Option<wiremock::MockServer>,
|
||||||
) -> TestApp {
|
) -> TestApp {
|
||||||
let state = AppState::from_database(&database, config);
|
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<wiremock::MockServer>,
|
||||||
|
) -> TestApp {
|
||||||
// Clone repos we need for TestApp before consuming state in the router
|
// Clone repos we need for TestApp before consuming state in the router
|
||||||
let roaster_repo = state.roaster_repo.clone();
|
let roaster_repo = state.roaster_repo.clone();
|
||||||
let roast_repo = state.roast_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
|
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 {
|
pub async fn spawn_app_with_foursquare_mock() -> TestApp {
|
||||||
let mock_server = wiremock::MockServer::start().await;
|
let mock_server = wiremock::MockServer::start().await;
|
||||||
let foursquare_url = format!("{}/places/search", mock_server.uri());
|
let foursquare_url = format!("{}/places/search", mock_server.uri());
|
||||||
|
|
|
||||||
|
|
@ -1,6 +1,7 @@
|
||||||
use crate::helpers::{
|
use crate::helpers::{
|
||||||
create_cafe_with_payload, create_default_bag, create_default_cafe, create_default_gear,
|
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,
|
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::brews::NewBrew;
|
||||||
use brewlog::domain::cafes::NewCafe;
|
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}"
|
"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}"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue