diff --git a/src/adapters/mcp.rs b/src/adapters/mcp.rs index 59d5679d..c8bfc785 100644 --- a/src/adapters/mcp.rs +++ b/src/adapters/mcp.rs @@ -5,9 +5,8 @@ //! templates so a host can read a specific entry or review without going //! through a tool call. //! -//! Every incoming MCP request is scoped to a single, fixed `UserContext` -//! configured at startup. MCP serving is intended for local clients and the -//! CLI rejects non-loopback bind addresses when MCP is enabled. +//! Every incoming MCP request is served from one tenant's isolated database; +//! the bearer token presented to the HTTP listener selects that tenant. use std::{borrow::Cow, sync::Arc}; @@ -25,7 +24,7 @@ use rmcp::{ }; use crate::journal::analyzer::{ - AnalyzerMcpComponents, UserContext, + AnalyzerMcpComponents, journal::JournalReadService, review::ReviewReadService, tools::{ToolError, ToolRegistry}, @@ -43,7 +42,6 @@ pub struct AnalyzerMcpServer { registry: Arc, journal_service: Arc, review_service: Arc, - user: UserContext, server_info: ServerInfo, tools: Arc<[Tool]>, resource_templates: Arc<[rmcp::model::ResourceTemplate]>, @@ -51,7 +49,7 @@ pub struct AnalyzerMcpServer { } impl AnalyzerMcpServer { - pub fn new(components: AnalyzerMcpComponents, user: UserContext) -> Self { + pub fn new(components: AnalyzerMcpComponents) -> Self { let AnalyzerMcpComponents { registry, journal_service, @@ -99,7 +97,6 @@ impl AnalyzerMcpServer { registry, journal_service, review_service, - user, server_info, tools, resource_templates, @@ -135,7 +132,7 @@ impl ServerHandler for AnalyzerMcpServer { let result = self .registry - .dispatch(request.name.as_ref(), &self.user, arguments) + .dispatch(request.name.as_ref(), arguments) .await; map_dispatch_result(result) } @@ -199,7 +196,7 @@ impl AnalyzerMcpServer { ResourceRef::JournalEntry(id) => { let entry = self .journal_service - .get_by_id(&self.user, &id) + .get_by_id(&id) .await .map_err(analyzer_to_mcp)? .ok_or_else(|| { @@ -219,7 +216,7 @@ impl AnalyzerMcpServer { ResourceRef::DailyReview(date) => { let review = self .review_service - .get_daily_review(&self.user, date) + .get_daily_review(date) .await .map_err(analyzer_to_mcp)? .ok_or_else(|| { @@ -239,7 +236,7 @@ impl AnalyzerMcpServer { ResourceRef::WeeklyReview(week_start) => { let review = self .review_service - .get_weekly_review(&self.user, week_start) + .get_weekly_review(week_start) .await .map_err(analyzer_to_mcp)? .ok_or_else(|| { @@ -489,13 +486,13 @@ mod tests { self.name } fn description(&self) -> &'static str { - "echoes its arguments and the user_id" + "echoes its arguments" } fn input_schema(&self) -> Value { json!({"type": "object", "properties": {"x": {"type": "integer"}}}) } - async fn dispatch(&self, ctx: &UserContext, args: Value) -> Result { - Ok(json!({"user": ctx.user_id, "args": args})) + async fn dispatch(&self, args: Value) -> Result { + Ok(json!({"args": args})) } } @@ -512,7 +509,7 @@ mod tests { fn input_schema(&self) -> Value { json!({"type": "object"}) } - async fn dispatch(&self, _ctx: &UserContext, _args: Value) -> Result { + async fn dispatch(&self, _args: Value) -> Result { Err(ToolError::Analyzer(AnalyzerError::InvalidArgument( "limit must be > 0".into(), ))) @@ -529,30 +526,23 @@ mod tests { impl JournalReadService for StubJournalService { async fn get_recent( &self, - _ctx: &UserContext, _request: GetRecentRequest, ) -> Result, AnalyzerError> { Ok(Vec::new()) } async fn search_text( &self, - _ctx: &UserContext, _request: SearchTextRequest, ) -> Result, AnalyzerError> { Ok(Vec::new()) } async fn search_semantic( &self, - _ctx: &UserContext, _request: SearchSemanticRequest, ) -> Result, AnalyzerError> { Ok(Vec::new()) } - async fn get_by_id( - &self, - _ctx: &UserContext, - id: &str, - ) -> Result, AnalyzerError> { + async fn get_by_id(&self, id: &str) -> Result, AnalyzerError> { *self.last_id.lock().unwrap() = Some(id.to_string()); Ok(self.get_by_id_response.lock().unwrap().clone()) } @@ -570,21 +560,18 @@ mod tests { impl ReviewReadService for StubReviewService { async fn get_daily_reviews( &self, - _ctx: &UserContext, _request: GetReviewsRequest, ) -> Result, AnalyzerError> { Ok(Vec::new()) } async fn get_weekly_reviews( &self, - _ctx: &UserContext, _request: GetReviewsRequest, ) -> Result, AnalyzerError> { Ok(Vec::new()) } async fn get_daily_review( &self, - _ctx: &UserContext, review_date: NaiveDate, ) -> Result, AnalyzerError> { *self.last_daily_date.lock().unwrap() = Some(review_date); @@ -592,7 +579,6 @@ mod tests { } async fn get_weekly_review( &self, - _ctx: &UserContext, week_start: NaiveDate, ) -> Result, AnalyzerError> { *self.last_weekly_week.lock().unwrap() = Some(week_start); @@ -623,13 +609,10 @@ mod tests { } fn server() -> AnalyzerMcpServer { - AnalyzerMcpServer::new( - components( - Arc::new(StubJournalService::default()), - Arc::new(StubReviewService::default()), - ), - UserContext::new("u-123"), - ) + AnalyzerMcpServer::new(components( + Arc::new(StubJournalService::default()), + Arc::new(StubReviewService::default()), + )) } fn at(h: u32) -> DateTime { @@ -651,10 +634,7 @@ mod tests { .iter() .find(|t| t.name == "echo") .expect("echo tool present"); - assert_eq!( - echo.description.as_deref(), - Some("echoes its arguments and the user_id") - ); + assert_eq!(echo.description.as_deref(), Some("echoes its arguments")); let schema = serde_json::Value::Object((*echo.input_schema).clone()); assert_eq!(schema["type"], "object"); assert!(schema["properties"]["x"].is_object()); @@ -829,7 +809,7 @@ mod tests { journal: Arc, review: Arc, ) -> AnalyzerMcpServer { - AnalyzerMcpServer::new(components(journal, review), UserContext::new("u-123")) + AnalyzerMcpServer::new(components(journal, review)) } #[tokio::test] @@ -943,17 +923,16 @@ mod tests { } #[tokio::test] - async fn dispatch_routes_to_registry_with_fixed_user() { + async fn dispatch_routes_arguments_to_registry() { let server = server(); let mut args = serde_json::Map::new(); args.insert("x".to_string(), json!(7)); let result = server .registry - .dispatch("echo", &server.user, Value::Object(args.clone())) + .dispatch("echo", Value::Object(args.clone())) .await .expect("echo dispatch ok"); - assert_eq!(result["user"], "u-123"); assert_eq!(result["args"]["x"], 7); } @@ -962,7 +941,7 @@ mod tests { let server = server(); let err = server .registry - .dispatch("missing", &server.user, json!({})) + .dispatch("missing", json!({})) .await .unwrap_err(); assert!(matches!(err, ToolError::UnknownTool(name) if name == "missing")); @@ -973,7 +952,7 @@ mod tests { let server = server(); let err = server .registry - .dispatch("failing", &server.user, json!({})) + .dispatch("failing", json!({})) .await .unwrap_err(); assert!(matches!( diff --git a/src/adapters/telegram.rs b/src/adapters/telegram.rs index 42af56c5..cf530b26 100644 --- a/src/adapters/telegram.rs +++ b/src/adapters/telegram.rs @@ -12,7 +12,7 @@ use crate::{ handler::MessageHandler, journal::command::{DEFAULT_RECENT_LIMIT, JournalCommand, JournalCommandRequest}, journal::transfer::{TransferError, TransferService}, - messages::{IncomingMessage, MessageSource, SINGLE_USER_ID}, + messages::{IncomingMessage, MessageSource}, tokens::TokenIssuer, }; @@ -272,7 +272,6 @@ async fn handle_message( let request = JournalCommandRequest { source: MessageSource::Telegram, source_conversation_id: message.chat.id.to_string(), - user_id: SINGLE_USER_ID.to_string(), received_at: message.date, command, }; @@ -360,7 +359,6 @@ fn incoming_from_text_message(message: &Message) -> IncomingMessage { source: MessageSource::Telegram, source_conversation_id: message.chat.id.to_string(), source_message_id: message.id.to_string(), - user_id: SINGLE_USER_ID.to_string(), text: message.text().unwrap_or_default().to_string(), received_at: message.date, } @@ -577,7 +575,6 @@ mod tests { assert_eq!(incoming.source, MessageSource::Telegram); assert_eq!(incoming.source_conversation_id, "42"); assert_eq!(incoming.source_message_id, "100"); - assert_eq!(incoming.user_id, SINGLE_USER_ID); assert_eq!(incoming.text, "hello froid"); assert_eq!( incoming.received_at, diff --git a/src/http.rs b/src/http.rs index 2acd08d5..4185ce0e 100644 --- a/src/http.rs +++ b/src/http.rs @@ -27,7 +27,7 @@ use crate::{ adapters::mcp::AnalyzerMcpServer, auth::{AuthenticatedTenant, TokenResolver, require_user_bearer}, journal::{ - analyzer::{DefaultSemanticJournalSearcher, UserContext, build_analyzer_mcp_components}, + analyzer::{DefaultSemanticJournalSearcher, build_analyzer_mcp_components}, embedding::{EmbeddingConfig, RigOpenAiEmbedder, SqliteEmbeddingRepository}, registry::JournalServiceRegistry, repository::JournalRepository, @@ -59,8 +59,7 @@ fn build_tenant_router(pool: &SqlitePool, config: &TenantRouterConfig) -> Result )); let components = build_analyzer_mcp_components(pool.clone(), semantic); - let user = UserContext::new(crate::messages::SINGLE_USER_ID); - let server = AnalyzerMcpServer::new(components, user); + let server = AnalyzerMcpServer::new(components); let service = StreamableHttpService::new( { diff --git a/src/journal/analyzer/journal.rs b/src/journal/analyzer/journal.rs index 7c1b41f2..39ff88b7 100644 --- a/src/journal/analyzer/journal.rs +++ b/src/journal/analyzer/journal.rs @@ -7,7 +7,7 @@ use crate::journal::repository::JournalRepository; use super::semantic::SemanticJournalSearcher; use super::types::{ AnalyzerError, GetRecentRequest, JournalEntryView, MAX_RECENT_LIMIT, MAX_SEMANTIC_LIMIT, - MAX_TEXT_SEARCH_LIMIT, SearchSemanticRequest, SearchTextRequest, SemanticHit, UserContext, + MAX_TEXT_SEARCH_LIMIT, SearchSemanticRequest, SearchTextRequest, SemanticHit, }; use super::validation::{validate_limit, validate_optional_range}; @@ -15,29 +15,22 @@ use super::validation::{validate_limit, validate_optional_range}; pub trait JournalReadService: Send + Sync { async fn get_recent( &self, - ctx: &UserContext, request: GetRecentRequest, ) -> Result, AnalyzerError>; async fn search_text( &self, - ctx: &UserContext, request: SearchTextRequest, ) -> Result, AnalyzerError>; async fn search_semantic( &self, - ctx: &UserContext, request: SearchSemanticRequest, ) -> Result, AnalyzerError>; /// Return the journal entry with the given id, or `None` if it does not /// exist (or does not belong to this user). - async fn get_by_id( - &self, - ctx: &UserContext, - id: &str, - ) -> Result, AnalyzerError>; + async fn get_by_id(&self, id: &str) -> Result, AnalyzerError>; } #[derive(Clone)] @@ -63,7 +56,6 @@ fn map_storage_error(err: sqlx::Error) -> AnalyzerError { impl JournalReadService for DefaultJournalReadService { async fn get_recent( &self, - _ctx: &UserContext, request: GetRecentRequest, ) -> Result, AnalyzerError> { let limit = validate_limit(request.limit, MAX_RECENT_LIMIT)?; @@ -90,7 +82,6 @@ impl JournalReadService for DefaultJournalReadService { async fn search_text( &self, - _ctx: &UserContext, request: SearchTextRequest, ) -> Result, AnalyzerError> { let limit = validate_limit(request.limit, MAX_TEXT_SEARCH_LIMIT)?; @@ -113,7 +104,6 @@ impl JournalReadService for DefaultJournalReadService { async fn search_semantic( &self, - ctx: &UserContext, request: SearchSemanticRequest, ) -> Result, AnalyzerError> { let limit = validate_limit(request.limit, MAX_SEMANTIC_LIMIT)?; @@ -128,7 +118,6 @@ impl JournalReadService for DefaultJournalReadService { let hits = self .semantic .search( - &ctx.user_id, trimmed, request.from_date, request.to_date_exclusive, @@ -139,11 +128,7 @@ impl JournalReadService for DefaultJournalReadService { Ok(hits) } - async fn get_by_id( - &self, - _ctx: &UserContext, - id: &str, - ) -> Result, AnalyzerError> { + async fn get_by_id(&self, id: &str) -> Result, AnalyzerError> { let mut rows = self .repository .fetch_by_ids(&[id.to_string()]) @@ -212,7 +197,6 @@ mod tests { impl SemanticJournalSearcher for StubSemanticSearcher { async fn search( &self, - _user_id: &str, _query: &str, from_date: Option, to_date_exclusive: Option, @@ -249,10 +233,6 @@ mod tests { } } - fn ctx() -> UserContext { - UserContext::new("user-1") - } - fn at(y: i32, m: u32, d: u32, h: u32, mi: u32) -> DateTime { Utc.with_ymd_and_hms(y, m, d, h, mi, 0).unwrap() } @@ -262,7 +242,6 @@ mod tests { source: MessageSource::Telegram, source_conversation_id: "42".to_string(), source_message_id: source_message_id.to_string(), - user_id: "user-1".to_string(), text: text.to_string(), received_at, } @@ -295,7 +274,7 @@ mod tests { .await .unwrap(); - let result = service.get_recent(&ctx(), recent(10)).await.unwrap(); + let result = service.get_recent(recent(10)).await.unwrap(); assert_eq!(result.len(), 2); assert_eq!(result[0].text, "second"); @@ -321,7 +300,7 @@ mod tests { from_date: Some(NaiveDate::from_ymd_opt(2026, 4, 28).unwrap()), to_date_exclusive: Some(NaiveDate::from_ymd_opt(2026, 4, 29).unwrap()), }; - let result = service.get_recent(&ctx(), req).await.unwrap(); + let result = service.get_recent(req).await.unwrap(); assert_eq!(result.len(), 1); assert_eq!(result[0].text, "in"); @@ -342,7 +321,7 @@ mod tests { from_date: Some(NaiveDate::from_ymd_opt(2026, 4, 28).unwrap()), to_date_exclusive: None, }; - let result = service.get_recent(&ctx(), req).await.unwrap(); + let result = service.get_recent(req).await.unwrap(); assert_eq!(result.len(), 1); assert_eq!(result[0].text, "in"); @@ -351,7 +330,7 @@ mod tests { #[tokio::test] async fn get_recent_rejects_zero_limit() { let (service, _) = setup().await; - let err = service.get_recent(&ctx(), recent(0)).await.unwrap_err(); + let err = service.get_recent(recent(0)).await.unwrap_err(); assert!(matches!(err, AnalyzerError::InvalidArgument(_))); } @@ -359,7 +338,7 @@ mod tests { async fn get_recent_rejects_limit_above_max() { let (service, _) = setup().await; let err = service - .get_recent(&ctx(), recent(MAX_RECENT_LIMIT + 1)) + .get_recent(recent(MAX_RECENT_LIMIT + 1)) .await .unwrap_err(); assert!(matches!(err, AnalyzerError::LimitTooLarge { .. })); @@ -373,7 +352,7 @@ mod tests { from_date: Some(NaiveDate::from_ymd_opt(2026, 4, 29).unwrap()), to_date_exclusive: Some(NaiveDate::from_ymd_opt(2026, 4, 28).unwrap()), }; - let err = service.get_recent(&ctx(), req).await.unwrap_err(); + let err = service.get_recent(req).await.unwrap_err(); assert!(matches!(err, AnalyzerError::InvalidArgument(_))); } @@ -387,14 +366,13 @@ mod tests { source: MessageSource::Telegram, source_conversation_id: "99".to_string(), source_message_id: "2".to_string(), - user_id: "user-2".to_string(), text: "theirs".to_string(), received_at: at(2026, 4, 28, 11, 0), }) .await .unwrap(); - let result = service.get_recent(&ctx(), recent(10)).await.unwrap(); + let result = service.get_recent(recent(10)).await.unwrap(); assert_eq!(result.len(), 2); assert_eq!(result[0].text, "theirs"); @@ -411,7 +389,7 @@ mod tests { .expect("entry stored"); let result = service - .get_by_id(&ctx(), &id) + .get_by_id(&id) .await .unwrap() .expect("entry present"); @@ -423,7 +401,7 @@ mod tests { #[tokio::test] async fn get_by_id_returns_none_when_missing() { let (service, _) = setup().await; - let result = service.get_by_id(&ctx(), "missing").await.unwrap(); + let result = service.get_by_id("missing").await.unwrap(); assert!(result.is_none()); } @@ -441,10 +419,7 @@ mod tests { .await .unwrap(); - let result = service - .search_text(&ctx(), search("ANXIOUS", 10)) - .await - .unwrap(); + let result = service.search_text(search("ANXIOUS", 10)).await.unwrap(); assert_eq!(result.len(), 1); assert_eq!(result[0].text, "felt anxious before the call"); @@ -453,10 +428,7 @@ mod tests { #[tokio::test] async fn search_text_trims_whitespace_and_rejects_empty_query() { let (service, _) = setup().await; - let err = service - .search_text(&ctx(), search(" ", 10)) - .await - .unwrap_err(); + let err = service.search_text(search(" ", 10)).await.unwrap_err(); assert!(matches!(err, AnalyzerError::InvalidArgument(_))); } @@ -464,7 +436,7 @@ mod tests { async fn search_text_rejects_limit_above_max() { let (service, _) = setup().await; let err = service - .search_text(&ctx(), search("x", MAX_TEXT_SEARCH_LIMIT + 1)) + .search_text(search("x", MAX_TEXT_SEARCH_LIMIT + 1)) .await .unwrap_err(); assert!(matches!(err, AnalyzerError::LimitTooLarge { .. })); @@ -480,17 +452,13 @@ mod tests { source: MessageSource::Telegram, source_conversation_id: "99".to_string(), source_message_id: "2".to_string(), - user_id: "user-2".to_string(), text: "theirs matches".to_string(), received_at: at(2026, 4, 28, 11, 0), }) .await .unwrap(); - let result = service - .search_text(&ctx(), search("matches", 10)) - .await - .unwrap(); + let result = service.search_text(search("matches", 10)).await.unwrap(); assert_eq!(result.len(), 2); assert_eq!(result[0].text, "theirs matches"); @@ -513,7 +481,7 @@ mod tests { from_date: Some(NaiveDate::from_ymd_opt(2026, 4, 28).unwrap()), to_date_exclusive: Some(NaiveDate::from_ymd_opt(2026, 4, 29).unwrap()), }; - let result = service.search_text(&ctx(), req).await.unwrap(); + let result = service.search_text(req).await.unwrap(); assert_eq!(result.len(), 1); assert_eq!(result[0].text, "match within"); @@ -530,7 +498,7 @@ mod tests { let (service, _) = setup_with_semantic(Arc::new(stub.clone())).await; let result = service - .search_semantic(&ctx(), semantic_req("anxiety", 5)) + .search_semantic(semantic_req("anxiety", 5)) .await .unwrap(); @@ -543,7 +511,7 @@ mod tests { async fn search_semantic_rejects_zero_limit() { let (service, _) = setup().await; let err = service - .search_semantic(&ctx(), semantic_req("x", 0)) + .search_semantic(semantic_req("x", 0)) .await .unwrap_err(); assert!(matches!(err, AnalyzerError::InvalidArgument(_))); @@ -553,7 +521,7 @@ mod tests { async fn search_semantic_rejects_limit_above_max() { let (service, _) = setup().await; let err = service - .search_semantic(&ctx(), semantic_req("x", MAX_SEMANTIC_LIMIT + 1)) + .search_semantic(semantic_req("x", MAX_SEMANTIC_LIMIT + 1)) .await .unwrap_err(); assert!(matches!(err, AnalyzerError::LimitTooLarge { .. })); @@ -563,7 +531,7 @@ mod tests { async fn search_semantic_rejects_blank_query() { let (service, _) = setup().await; let err = service - .search_semantic(&ctx(), semantic_req(" ", 5)) + .search_semantic(semantic_req(" ", 5)) .await .unwrap_err(); assert!(matches!(err, AnalyzerError::InvalidArgument(_))); @@ -578,7 +546,7 @@ mod tests { from_date: Some(NaiveDate::from_ymd_opt(2026, 4, 29).unwrap()), to_date_exclusive: Some(NaiveDate::from_ymd_opt(2026, 4, 28).unwrap()), }; - let err = service.search_semantic(&ctx(), req).await.unwrap_err(); + let err = service.search_semantic(req).await.unwrap_err(); assert!(matches!(err, AnalyzerError::InvalidArgument(_))); } @@ -598,7 +566,7 @@ mod tests { from_date: Some(NaiveDate::from_ymd_opt(2026, 4, 28).unwrap()), to_date_exclusive: Some(NaiveDate::from_ymd_opt(2026, 4, 29).unwrap()), }; - let result = service.search_semantic(&ctx(), req).await.unwrap(); + let result = service.search_semantic(req).await.unwrap(); assert_eq!(result.len(), 1); assert_eq!(result[0].text, "in-1"); @@ -618,10 +586,7 @@ mod tests { let stub = StubSemanticSearcher::with_hits(vec![]); let (service, _) = setup_with_semantic(Arc::new(stub.clone())).await; - let _ = service - .search_semantic(&ctx(), semantic_req("x", 3)) - .await - .unwrap(); + let _ = service.search_semantic(semantic_req("x", 3)).await.unwrap(); assert_eq!(stub.last_limit(), Some(3)); } diff --git a/src/journal/analyzer/mod.rs b/src/journal/analyzer/mod.rs index d9b669ba..3759e43b 100644 --- a/src/journal/analyzer/mod.rs +++ b/src/journal/analyzer/mod.rs @@ -1,9 +1,8 @@ //! Read-only services exposed to MCP clients. //! -//! Froid is a single-user journal. Every method takes a fixed [`UserContext`] -//! from server startup, never an LLM-supplied user id. The MCP adapter is -//! restricted to loopback binds by configuration validation, and each service -//! caps requested limits to a maximum. +//! Each service is bound to one tenant's isolated database — user identity +//! comes from the bearer token at the HTTP layer, never from LLM-supplied +//! input — and each service caps requested limits to a maximum. pub mod journal; pub mod review; @@ -15,5 +14,4 @@ mod validation; pub mod wiring; pub use semantic::{DefaultSemanticJournalSearcher, SemanticJournalSearcher}; -pub use types::UserContext; pub use wiring::{AnalyzerMcpComponents, build_analyzer_mcp_components}; diff --git a/src/journal/analyzer/review.rs b/src/journal/analyzer/review.rs index 350b4e4e..379d28d4 100644 --- a/src/journal/analyzer/review.rs +++ b/src/journal/analyzer/review.rs @@ -8,22 +8,18 @@ use crate::journal::week_review::repository::{ }; use crate::journal::week_review::{WeeklyReview, WeeklyReviewStatus}; -use super::types::{ - AnalyzerError, DailyReviewView, GetReviewsRequest, UserContext, WeeklyReviewView, -}; +use super::types::{AnalyzerError, DailyReviewView, GetReviewsRequest, WeeklyReviewView}; use super::validation::validate_range; #[async_trait] pub trait ReviewReadService: Send + Sync { async fn get_daily_reviews( &self, - ctx: &UserContext, request: GetReviewsRequest, ) -> Result, AnalyzerError>; async fn get_weekly_reviews( &self, - ctx: &UserContext, request: GetReviewsRequest, ) -> Result, AnalyzerError>; @@ -31,7 +27,6 @@ pub trait ReviewReadService: Send + Sync { /// completed review exists for that date. async fn get_daily_review( &self, - ctx: &UserContext, review_date: NaiveDate, ) -> Result, AnalyzerError>; @@ -39,7 +34,6 @@ pub trait ReviewReadService: Send + Sync { /// or `None` if no completed review exists for that week. async fn get_weekly_review( &self, - ctx: &UserContext, week_start: NaiveDate, ) -> Result, AnalyzerError>; } @@ -68,7 +62,6 @@ fn map_weekly_error(err: WeeklyReviewRepositoryError) -> AnalyzerError { impl ReviewReadService for DefaultReviewReadService { async fn get_daily_reviews( &self, - _ctx: &UserContext, request: GetReviewsRequest, ) -> Result, AnalyzerError> { validate_range(request.from_date, request.to_date_exclusive)?; @@ -91,7 +84,6 @@ impl ReviewReadService for DefaultReviewReadService { async fn get_weekly_reviews( &self, - _ctx: &UserContext, request: GetReviewsRequest, ) -> Result, AnalyzerError> { validate_range(request.from_date, request.to_date_exclusive)?; @@ -115,7 +107,6 @@ impl ReviewReadService for DefaultReviewReadService { async fn get_daily_review( &self, - _ctx: &UserContext, review_date: NaiveDate, ) -> Result, AnalyzerError> { let row = self @@ -129,7 +120,6 @@ impl ReviewReadService for DefaultReviewReadService { async fn get_weekly_review( &self, - _ctx: &UserContext, week_start: NaiveDate, ) -> Result, AnalyzerError> { let row = self @@ -195,10 +185,6 @@ mod tests { (service, daily, weekly) } - fn ctx() -> UserContext { - UserContext::new("user-1") - } - fn ymd(y: i32, m: u32, d: u32) -> NaiveDate { NaiveDate::from_ymd_opt(y, m, d).unwrap() } @@ -227,7 +213,7 @@ mod tests { .unwrap(); let result = service - .get_daily_reviews(&ctx(), req(ymd(2026, 4, 27), ymd(2026, 4, 29))) + .get_daily_reviews(req(ymd(2026, 4, 27), ymd(2026, 4, 29))) .await .unwrap(); @@ -251,7 +237,7 @@ mod tests { .unwrap(); let result = service - .get_daily_reviews(&ctx(), req(ymd(2026, 4, 27), ymd(2026, 4, 29))) + .get_daily_reviews(req(ymd(2026, 4, 27), ymd(2026, 4, 29))) .await .unwrap(); @@ -272,7 +258,7 @@ mod tests { .unwrap(); let result = service - .get_daily_reviews(&ctx(), req(ymd(2026, 4, 27), ymd(2026, 4, 28))) + .get_daily_reviews(req(ymd(2026, 4, 27), ymd(2026, 4, 28))) .await .unwrap(); @@ -284,7 +270,7 @@ mod tests { async fn get_daily_reviews_rejects_inverted_range() { let (service, _, _) = setup().await; let err = service - .get_daily_reviews(&ctx(), req(ymd(2026, 4, 28), ymd(2026, 4, 27))) + .get_daily_reviews(req(ymd(2026, 4, 28), ymd(2026, 4, 27))) .await .unwrap_err(); assert!(matches!(err, AnalyzerError::InvalidArgument(_))); @@ -294,7 +280,7 @@ mod tests { async fn get_daily_reviews_rejects_equal_bounds() { let (service, _, _) = setup().await; let err = service - .get_daily_reviews(&ctx(), req(ymd(2026, 4, 28), ymd(2026, 4, 28))) + .get_daily_reviews(req(ymd(2026, 4, 28), ymd(2026, 4, 28))) .await .unwrap_err(); assert!(matches!(err, AnalyzerError::InvalidArgument(_))); @@ -316,7 +302,7 @@ mod tests { .unwrap(); let result = service - .get_weekly_reviews(&ctx(), req(w1, ymd(2026, 5, 4))) + .get_weekly_reviews(req(w1, ymd(2026, 5, 4))) .await .unwrap(); @@ -341,7 +327,7 @@ mod tests { weekly.upsert_failed(w2, "m", "v1", "boom").await.unwrap(); let result = service - .get_weekly_reviews(&ctx(), req(w1, ymd(2026, 5, 4))) + .get_weekly_reviews(req(w1, ymd(2026, 5, 4))) .await .unwrap(); @@ -363,7 +349,7 @@ mod tests { .unwrap(); let result = service - .get_weekly_reviews(&ctx(), req(w, ymd(2026, 4, 27))) + .get_weekly_reviews(req(w, ymd(2026, 4, 27))) .await .unwrap(); @@ -380,7 +366,7 @@ mod tests { .unwrap(); let result = service - .get_daily_review(&ctx(), ymd(2026, 4, 28)) + .get_daily_review(ymd(2026, 4, 28)) .await .unwrap() .expect("completed review present"); @@ -392,10 +378,7 @@ mod tests { #[tokio::test] async fn get_daily_review_returns_none_when_missing() { let (service, _, _) = setup().await; - let result = service - .get_daily_review(&ctx(), ymd(2026, 4, 28)) - .await - .unwrap(); + let result = service.get_daily_review(ymd(2026, 4, 28)).await.unwrap(); assert!(result.is_none()); } @@ -407,10 +390,7 @@ mod tests { .await .unwrap(); - let result = service - .get_daily_review(&ctx(), ymd(2026, 4, 28)) - .await - .unwrap(); + let result = service.get_daily_review(ymd(2026, 4, 28)).await.unwrap(); assert!(result.is_none()); } @@ -424,7 +404,7 @@ mod tests { .unwrap(); let result = service - .get_weekly_review(&ctx(), w) + .get_weekly_review(w) .await .unwrap() .expect("completed review present"); @@ -437,10 +417,7 @@ mod tests { #[tokio::test] async fn get_weekly_review_returns_none_when_missing() { let (service, _, _) = setup().await; - let result = service - .get_weekly_review(&ctx(), ymd(2026, 4, 20)) - .await - .unwrap(); + let result = service.get_weekly_review(ymd(2026, 4, 20)).await.unwrap(); assert!(result.is_none()); } @@ -450,7 +427,7 @@ mod tests { let w = ymd(2026, 4, 20); weekly.upsert_failed(w, "m", "v1", "boom").await.unwrap(); - let result = service.get_weekly_review(&ctx(), w).await.unwrap(); + let result = service.get_weekly_review(w).await.unwrap(); assert!(result.is_none()); } @@ -458,7 +435,7 @@ mod tests { async fn get_weekly_reviews_rejects_inverted_range() { let (service, _, _) = setup().await; let err = service - .get_weekly_reviews(&ctx(), req(ymd(2026, 4, 28), ymd(2026, 4, 27))) + .get_weekly_reviews(req(ymd(2026, 4, 28), ymd(2026, 4, 27))) .await .unwrap_err(); assert!(matches!(err, AnalyzerError::InvalidArgument(_))); diff --git a/src/journal/analyzer/semantic.rs b/src/journal/analyzer/semantic.rs index 288f82ab..ab01334a 100644 --- a/src/journal/analyzer/semantic.rs +++ b/src/journal/analyzer/semantic.rs @@ -12,10 +12,9 @@ use super::types::{AnalyzerError, SemanticHit}; #[async_trait] pub trait SemanticJournalSearcher: Send + Sync { /// Returns up to `limit` journal entries semantically similar to `query`, - /// scoped to `user_id` and the optional date bounds. + /// scoped to the optional date bounds. async fn search( &self, - user_id: &str, query: &str, from_date: Option, to_date_exclusive: Option, @@ -52,7 +51,6 @@ where { async fn search( &self, - user_id: &str, query: &str, from_date: Option, to_date_exclusive: Option, @@ -68,14 +66,7 @@ where let index_results = self .index - .search_for_user( - user_id, - &embedding, - model, - from_date, - to_date_exclusive, - limit, - ) + .search(&embedding, model, from_date, to_date_exclusive, limit) .await .map_err(|e| AnalyzerError::Internal(Box::new(e)))?; @@ -154,7 +145,6 @@ mod tests { source: MessageSource::Telegram, source_conversation_id: "42".to_string(), source_message_id: msg_id.to_string(), - user_id: "user-1".to_string(), text: text.to_string(), received_at, }; @@ -259,9 +249,8 @@ mod tests { ) -> Result { unreachable!() } - async fn search_for_user( + async fn search( &self, - _user_id: &str, _embedding: &Embedding, _embedding_model: &str, _from_date: Option, @@ -282,10 +271,7 @@ mod tests { let searcher = DefaultSemanticJournalSearcher::new(index, FakeEmbedder::succeeds(TEST_MODEL, 1), repo); - let hits = searcher - .search("user-1", "query", None, None, 10) - .await - .unwrap(); + let hits = searcher.search("query", None, None, 10).await.unwrap(); assert!(!hits.is_empty()); assert_eq!(hits[0].text, "closest"); @@ -311,10 +297,7 @@ mod tests { let searcher = DefaultSemanticJournalSearcher::new(index, FakeEmbedder::succeeds(TEST_MODEL, 0), repo); - let hits = searcher - .search("user-1", "query", None, None, 3) - .await - .unwrap(); + let hits = searcher.search("query", None, None, 3).await.unwrap(); assert_eq!(hits.len(), 3); } @@ -328,7 +311,6 @@ mod tests { source: MessageSource::Telegram, source_conversation_id: "99".to_string(), source_message_id: "2".to_string(), - user_id: "user-2".to_string(), text: "theirs".to_string(), received_at: at(11, 0), }; @@ -352,10 +334,7 @@ mod tests { let searcher = DefaultSemanticJournalSearcher::new(index, FakeEmbedder::succeeds(TEST_MODEL, 0), repo); - let hits = searcher - .search("user-1", "query", None, None, 10) - .await - .unwrap(); + let hits = searcher.search("query", None, None, 10).await.unwrap(); assert_eq!(hits.len(), 2); assert_eq!(hits[0].text, "mine"); @@ -368,10 +347,7 @@ mod tests { let searcher = DefaultSemanticJournalSearcher::new(index, FakeEmbedder::fails(TEST_MODEL), repo); - let err = searcher - .search("user-1", "query", None, None, 5) - .await - .unwrap_err(); + let err = searcher.search("query", None, None, 5).await.unwrap_err(); assert!(matches!(err, AnalyzerError::Internal(_))); } @@ -384,7 +360,6 @@ mod tests { source: MessageSource::Telegram, source_conversation_id: "42".to_string(), source_message_id: "1".to_string(), - user_id: "user-1".to_string(), text: "kept".to_string(), received_at: at(10, 0), }; @@ -415,10 +390,7 @@ mod tests { repo, ); - let hits = searcher - .search("user-1", "query", None, None, 5) - .await - .unwrap(); + let hits = searcher.search("query", None, None, 5).await.unwrap(); assert_eq!(hits.len(), 1); assert_eq!(hits[0].text, "kept"); diff --git a/src/journal/analyzer/signal.rs b/src/journal/analyzer/signal.rs index 27bc8d24..dfef818a 100644 --- a/src/journal/analyzer/signal.rs +++ b/src/journal/analyzer/signal.rs @@ -4,18 +4,13 @@ use crate::journal::review::signals::repository::{ DailyReviewSignalRepository, DailyReviewSignalRepositoryError, SignalSearchFilters, }; -use super::types::{ - AnalyzerError, MAX_SIGNAL_LIMIT, SearchSignalsRequest, SignalView, UserContext, -}; +use super::types::{AnalyzerError, MAX_SIGNAL_LIMIT, SearchSignalsRequest, SignalView}; use super::validation::{validate_limit, validate_optional_range}; #[async_trait] pub trait SignalReadService: Send + Sync { - async fn search( - &self, - ctx: &UserContext, - request: SearchSignalsRequest, - ) -> Result, AnalyzerError>; + async fn search(&self, request: SearchSignalsRequest) + -> Result, AnalyzerError>; } #[derive(Debug, Clone)] @@ -64,7 +59,6 @@ fn normalize_label_contains(value: Option) -> Result, Ana impl SignalReadService for DefaultSignalReadService { async fn search( &self, - _ctx: &UserContext, request: SearchSignalsRequest, ) -> Result, AnalyzerError> { let limit = validate_limit(request.limit, MAX_SIGNAL_LIMIT)?; @@ -131,10 +125,6 @@ mod tests { (service, signals, reviews) } - fn ctx() -> UserContext { - UserContext::new("user-1") - } - fn ymd(y: i32, m: u32, d: u32) -> NaiveDate { NaiveDate::from_ymd_opt(y, m, d).unwrap() } @@ -178,7 +168,6 @@ mod tests { async fn seed( signals: &DailyReviewSignalRepository, reviews: &DailyReviewRepository, - _user_id: &str, date: NaiveDate, candidates: &[DailyReviewSignalCandidate], ) { @@ -202,18 +191,11 @@ mod tests { #[tokio::test] async fn search_returns_signals_in_date_then_id_order() { let (service, signals, reviews) = setup().await; - seed(&signals, &reviews, "user-1", ymd(2026, 4, 27), &[theme()]).await; - seed(&signals, &reviews, "user-1", ymd(2026, 4, 28), &[need()]).await; - seed( - &signals, - &reviews, - "user-1", - ymd(2026, 4, 29), - &[behavior()], - ) - .await; - - let result = service.search(&ctx(), req(10)).await.unwrap(); + seed(&signals, &reviews, ymd(2026, 4, 27), &[theme()]).await; + seed(&signals, &reviews, ymd(2026, 4, 28), &[need()]).await; + seed(&signals, &reviews, ymd(2026, 4, 29), &[behavior()]).await; + + let result = service.search(req(10)).await.unwrap(); assert_eq!(result.len(), 3); let dates: Vec<_> = result.iter().map(|s| s.review_date).collect(); @@ -226,26 +208,16 @@ mod tests { #[tokio::test] async fn search_applies_filters() { let (service, signals, reviews) = setup().await; - seed(&signals, &reviews, "user-1", ymd(2026, 4, 27), &[theme()]).await; - seed(&signals, &reviews, "user-1", ymd(2026, 4, 28), &[need()]).await; - seed( - &signals, - &reviews, - "user-1", - ymd(2026, 4, 29), - &[behavior()], - ) - .await; + seed(&signals, &reviews, ymd(2026, 4, 27), &[theme()]).await; + seed(&signals, &reviews, ymd(2026, 4, 28), &[need()]).await; + seed(&signals, &reviews, ymd(2026, 4, 29), &[behavior()]).await; let result = service - .search( - &ctx(), - SearchSignalsRequest { - signal_type: Some(SignalType::Behavior), - valence: Some(BehaviorValence::Negative), - ..req(10) - }, - ) + .search(SearchSignalsRequest { + signal_type: Some(SignalType::Behavior), + valence: Some(BehaviorValence::Negative), + ..req(10) + }) .await .unwrap(); @@ -257,28 +229,22 @@ mod tests { #[tokio::test] async fn search_trims_label_contains_and_rejects_blank() { let (service, signals, reviews) = setup().await; - seed(&signals, &reviews, "user-1", ymd(2026, 4, 27), &[theme()]).await; + seed(&signals, &reviews, ymd(2026, 4, 27), &[theme()]).await; let result = service - .search( - &ctx(), - SearchSignalsRequest { - label_contains: Some(" PHYSICAL ".to_string()), - ..req(10) - }, - ) + .search(SearchSignalsRequest { + label_contains: Some(" PHYSICAL ".to_string()), + ..req(10) + }) .await .unwrap(); assert_eq!(result.len(), 1); let err = service - .search( - &ctx(), - SearchSignalsRequest { - label_contains: Some(" ".to_string()), - ..req(10) - }, - ) + .search(SearchSignalsRequest { + label_contains: Some(" ".to_string()), + ..req(10) + }) .await .unwrap_err(); assert!(matches!(err, AnalyzerError::InvalidArgument(_))); @@ -289,13 +255,10 @@ mod tests { let (service, _, _) = setup().await; for invalid in [-0.1, 1.1] { let err = service - .search( - &ctx(), - SearchSignalsRequest { - min_strength: Some(invalid), - ..req(10) - }, - ) + .search(SearchSignalsRequest { + min_strength: Some(invalid), + ..req(10) + }) .await .unwrap_err(); assert!(matches!(err, AnalyzerError::InvalidArgument(_))); @@ -305,17 +268,14 @@ mod tests { #[tokio::test] async fn search_rejects_zero_limit() { let (service, _, _) = setup().await; - let err = service.search(&ctx(), req(0)).await.unwrap_err(); + let err = service.search(req(0)).await.unwrap_err(); assert!(matches!(err, AnalyzerError::InvalidArgument(_))); } #[tokio::test] async fn search_rejects_limit_above_max() { let (service, _, _) = setup().await; - let err = service - .search(&ctx(), req(MAX_SIGNAL_LIMIT + 1)) - .await - .unwrap_err(); + let err = service.search(req(MAX_SIGNAL_LIMIT + 1)).await.unwrap_err(); assert!(matches!(err, AnalyzerError::LimitTooLarge { .. })); } @@ -323,14 +283,11 @@ mod tests { async fn search_rejects_inverted_range() { let (service, _, _) = setup().await; let err = service - .search( - &ctx(), - SearchSignalsRequest { - from_date: Some(ymd(2026, 4, 29)), - to_date_exclusive: Some(ymd(2026, 4, 28)), - ..req(10) - }, - ) + .search(SearchSignalsRequest { + from_date: Some(ymd(2026, 4, 29)), + to_date_exclusive: Some(ymd(2026, 4, 28)), + ..req(10) + }) .await .unwrap_err(); assert!(matches!(err, AnalyzerError::InvalidArgument(_))); @@ -339,10 +296,10 @@ mod tests { #[tokio::test] async fn search_uses_single_user_scope() { let (service, signals, reviews) = setup().await; - seed(&signals, &reviews, "user-1", ymd(2026, 4, 27), &[theme()]).await; - seed(&signals, &reviews, "user-2", ymd(2026, 4, 27), &[theme()]).await; + seed(&signals, &reviews, ymd(2026, 4, 27), &[theme()]).await; + seed(&signals, &reviews, ymd(2026, 4, 27), &[theme()]).await; - let result = service.search(&ctx(), req(10)).await.unwrap(); + let result = service.search(req(10)).await.unwrap(); assert_eq!(result.len(), 1); } @@ -350,9 +307,9 @@ mod tests { #[tokio::test] async fn search_view_preserves_signal_fields() { let (service, signals, reviews) = setup().await; - seed(&signals, &reviews, "user-1", ymd(2026, 4, 28), &[need()]).await; + seed(&signals, &reviews, ymd(2026, 4, 28), &[need()]).await; - let result = service.search(&ctx(), req(10)).await.unwrap(); + let result = service.search(req(10)).await.unwrap(); assert_eq!(result.len(), 1); let s = &result[0]; diff --git a/src/journal/analyzer/tools/journal.rs b/src/journal/analyzer/tools/journal.rs index 4d388544..daa61939 100644 --- a/src/journal/analyzer/tools/journal.rs +++ b/src/journal/analyzer/tools/journal.rs @@ -10,7 +10,6 @@ use super::{Tool, ToolError, deserialize_input, schema_value, serialize_output}; use crate::journal::analyzer::journal::JournalReadService; use crate::journal::analyzer::types::{ GetRecentRequest, JournalEntryView, SearchSemanticRequest, SearchTextRequest, SemanticHit, - UserContext, }; #[derive(Debug, Deserialize, JsonSchema)] @@ -68,18 +67,15 @@ impl Tool for JournalGetRecentTool { fn input_schema(&self) -> Value { schema_value::() } - async fn dispatch(&self, ctx: &UserContext, args: Value) -> Result { + async fn dispatch(&self, args: Value) -> Result { let input: GetRecentInput = deserialize_input(args)?; let entries = self .service - .get_recent( - ctx, - GetRecentRequest { - limit: input.limit, - from_date: input.from_date, - to_date_exclusive: input.to_date_exclusive, - }, - ) + .get_recent(GetRecentRequest { + limit: input.limit, + from_date: input.from_date, + to_date_exclusive: input.to_date_exclusive, + }) .await?; serialize_output(JournalEntriesOutput { entries: entries.into_iter().map(JournalEntryItem::from).collect(), @@ -122,19 +118,16 @@ impl Tool for JournalSearchTextTool { fn input_schema(&self) -> Value { schema_value::() } - async fn dispatch(&self, ctx: &UserContext, args: Value) -> Result { + async fn dispatch(&self, args: Value) -> Result { let input: SearchTextInput = deserialize_input(args)?; let entries = self .service - .search_text( - ctx, - SearchTextRequest { - query: input.query, - limit: input.limit, - from_date: input.from_date, - to_date_exclusive: input.to_date_exclusive, - }, - ) + .search_text(SearchTextRequest { + query: input.query, + limit: input.limit, + from_date: input.from_date, + to_date_exclusive: input.to_date_exclusive, + }) .await?; serialize_output(JournalEntriesOutput { entries: entries.into_iter().map(JournalEntryItem::from).collect(), @@ -202,19 +195,16 @@ impl Tool for JournalSearchSemanticTool { fn input_schema(&self) -> Value { schema_value::() } - async fn dispatch(&self, ctx: &UserContext, args: Value) -> Result { + async fn dispatch(&self, args: Value) -> Result { let input: SearchSemanticInput = deserialize_input(args)?; let hits = self .service - .search_semantic( - ctx, - SearchSemanticRequest { - query: input.query, - limit: input.limit, - from_date: input.from_date, - to_date_exclusive: input.to_date_exclusive, - }, - ) + .search_semantic(SearchSemanticRequest { + query: input.query, + limit: input.limit, + from_date: input.from_date, + to_date_exclusive: input.to_date_exclusive, + }) .await?; serialize_output(SemanticHitsOutput { hits: hits.into_iter().map(SemanticHitItem::from).collect(), @@ -244,7 +234,6 @@ mod tests { impl JournalReadService for StubJournalService { async fn get_recent( &self, - _ctx: &UserContext, request: GetRecentRequest, ) -> Result, AnalyzerError> { *self.last_recent.lock().unwrap() = Some(request); @@ -252,7 +241,6 @@ mod tests { } async fn search_text( &self, - _ctx: &UserContext, request: SearchTextRequest, ) -> Result, AnalyzerError> { *self.last_text.lock().unwrap() = Some(request); @@ -260,25 +248,16 @@ mod tests { } async fn search_semantic( &self, - _ctx: &UserContext, request: SearchSemanticRequest, ) -> Result, AnalyzerError> { *self.last_semantic.lock().unwrap() = Some(request); Ok(self.semantic_response.lock().unwrap().clone()) } - async fn get_by_id( - &self, - _ctx: &UserContext, - _id: &str, - ) -> Result, AnalyzerError> { + async fn get_by_id(&self, _id: &str) -> Result, AnalyzerError> { Ok(None) } } - fn ctx() -> UserContext { - UserContext::new("user-1") - } - fn at(h: u32) -> DateTime { Utc.with_ymd_and_hms(2026, 4, 28, h, 0, 0).unwrap() } @@ -295,14 +274,11 @@ mod tests { let tool = JournalGetRecentTool::new(stub.clone()); let out = tool - .dispatch( - &ctx(), - json!({ - "limit": 5, - "from_date": "2026-04-28", - "to_date_exclusive": "2026-04-29" - }), - ) + .dispatch(json!({ + "limit": 5, + "from_date": "2026-04-28", + "to_date_exclusive": "2026-04-29" + })) .await .unwrap(); @@ -322,7 +298,7 @@ mod tests { let stub = Arc::new(StubJournalService::default()); let tool = JournalGetRecentTool::new(stub.clone()); - let _ = tool.dispatch(&ctx(), json!({"limit": 3})).await.unwrap(); + let _ = tool.dispatch(json!({"limit": 3})).await.unwrap(); let captured = stub.last_recent.lock().unwrap().clone().unwrap(); assert_eq!(captured.limit, 3); @@ -335,7 +311,7 @@ mod tests { let tool = JournalGetRecentTool::new(Arc::new(StubJournalService::default())); let err = tool - .dispatch(&ctx(), json!({"limit": "not a number"})) + .dispatch(json!({"limit": "not a number"})) .await .unwrap_err(); assert!(matches!(err, ToolError::InvalidInput(_))); @@ -348,39 +324,29 @@ mod tests { impl JournalReadService for FailingService { async fn get_recent( &self, - _: &UserContext, _: GetRecentRequest, ) -> Result, AnalyzerError> { Err(AnalyzerError::LimitTooLarge { max: 50 }) } async fn search_text( &self, - _: &UserContext, _: SearchTextRequest, ) -> Result, AnalyzerError> { unreachable!() } async fn search_semantic( &self, - _: &UserContext, _: SearchSemanticRequest, ) -> Result, AnalyzerError> { unreachable!() } - async fn get_by_id( - &self, - _: &UserContext, - _: &str, - ) -> Result, AnalyzerError> { + async fn get_by_id(&self, _: &str) -> Result, AnalyzerError> { unreachable!() } } let tool = JournalGetRecentTool::new(Arc::new(FailingService)); - let err = tool - .dispatch(&ctx(), json!({"limit": 999})) - .await - .unwrap_err(); + let err = tool.dispatch(json!({"limit": 999})).await.unwrap_err(); assert!(matches!( err, @@ -399,7 +365,7 @@ mod tests { let tool = JournalSearchTextTool::new(stub.clone()); let out = tool - .dispatch(&ctx(), json!({"query": "anxious", "limit": 5})) + .dispatch(json!({"query": "anxious", "limit": 5})) .await .unwrap(); @@ -421,7 +387,7 @@ mod tests { let tool = JournalSearchSemanticTool::new(stub.clone()); let out = tool - .dispatch(&ctx(), json!({"query": "avoidance", "limit": 3})) + .dispatch(json!({"query": "avoidance", "limit": 3})) .await .unwrap(); diff --git a/src/journal/analyzer/tools/mod.rs b/src/journal/analyzer/tools/mod.rs index 908473e8..192319ed 100644 --- a/src/journal/analyzer/tools/mod.rs +++ b/src/journal/analyzer/tools/mod.rs @@ -3,7 +3,7 @@ //! Each tool exposes a stable name, a JSON schema for inputs, and a dispatch //! method that takes JSON in and returns JSON out. The analyzer agent loop //! looks tools up by name in [`ToolRegistry`] and invokes them with the -//! authenticated [`UserContext`] — `user_id` is never part of the tool input. +//! per-tenant database — user identity is never part of the tool input. pub mod journal; pub mod review; @@ -16,7 +16,7 @@ use thiserror::Error; use async_trait::async_trait; use serde_json::Value; -use super::types::{AnalyzerError, UserContext}; +use super::types::AnalyzerError; #[derive(Debug, Error)] pub enum ToolError { @@ -39,7 +39,7 @@ pub trait Tool: Send + Sync { fn name(&self) -> &'static str; fn description(&self) -> &'static str; fn input_schema(&self) -> Value; - async fn dispatch(&self, ctx: &UserContext, args: Value) -> Result; + async fn dispatch(&self, args: Value) -> Result; } /// Holds an ordered set of tools indexed by name. @@ -74,16 +74,11 @@ impl ToolRegistry { &self.tools } - pub async fn dispatch( - &self, - name: &str, - ctx: &UserContext, - args: Value, - ) -> Result { + pub async fn dispatch(&self, name: &str, args: Value) -> Result { let tool = self .get(name) .ok_or_else(|| ToolError::UnknownTool(name.to_string()))?; - tool.dispatch(ctx, args).await + tool.dispatch(args).await } } @@ -126,24 +121,17 @@ mod tests { fn input_schema(&self) -> Value { json!({"type": "object"}) } - async fn dispatch(&self, _ctx: &UserContext, args: Value) -> Result { + async fn dispatch(&self, args: Value) -> Result { Ok(args) } } - fn ctx() -> UserContext { - UserContext::new("u") - } - #[tokio::test] async fn registry_dispatches_by_name() { let mut registry = ToolRegistry::new(); registry.register(Arc::new(EchoTool)); - let out = registry - .dispatch("echo", &ctx(), json!({"x": 1})) - .await - .unwrap(); + let out = registry.dispatch("echo", json!({"x": 1})).await.unwrap(); assert_eq!(out, json!({"x": 1})); } @@ -152,10 +140,7 @@ mod tests { async fn registry_returns_unknown_tool_error() { let registry = ToolRegistry::new(); - let err = registry - .dispatch("missing", &ctx(), json!({})) - .await - .unwrap_err(); + let err = registry.dispatch("missing", json!({})).await.unwrap_err(); assert!(matches!(err, ToolError::UnknownTool(name) if name == "missing")); } diff --git a/src/journal/analyzer/tools/review.rs b/src/journal/analyzer/tools/review.rs index ef5756e0..086e6eb3 100644 --- a/src/journal/analyzer/tools/review.rs +++ b/src/journal/analyzer/tools/review.rs @@ -8,9 +8,7 @@ use serde_json::Value; use super::{Tool, ToolError, deserialize_input, schema_value, serialize_output}; use crate::journal::analyzer::review::ReviewReadService; -use crate::journal::analyzer::types::{ - DailyReviewView, GetReviewsRequest, UserContext, WeeklyReviewView, -}; +use crate::journal::analyzer::types::{DailyReviewView, GetReviewsRequest, WeeklyReviewView}; #[derive(Debug, Deserialize, JsonSchema)] struct GetReviewsInput { @@ -64,17 +62,14 @@ impl Tool for DailyReviewGetRangeTool { fn input_schema(&self) -> Value { schema_value::() } - async fn dispatch(&self, ctx: &UserContext, args: Value) -> Result { + async fn dispatch(&self, args: Value) -> Result { let input: GetReviewsInput = deserialize_input(args)?; let reviews = self .service - .get_daily_reviews( - ctx, - GetReviewsRequest { - from_date: input.from_date, - to_date_exclusive: input.to_date_exclusive, - }, - ) + .get_daily_reviews(GetReviewsRequest { + from_date: input.from_date, + to_date_exclusive: input.to_date_exclusive, + }) .await?; serialize_output(DailyReviewsOutput { reviews: reviews.into_iter().map(DailyReviewItem::from).collect(), @@ -127,17 +122,14 @@ impl Tool for WeeklyReviewGetRangeTool { fn input_schema(&self) -> Value { schema_value::() } - async fn dispatch(&self, ctx: &UserContext, args: Value) -> Result { + async fn dispatch(&self, args: Value) -> Result { let input: GetReviewsInput = deserialize_input(args)?; let reviews = self .service - .get_weekly_reviews( - ctx, - GetReviewsRequest { - from_date: input.from_date, - to_date_exclusive: input.to_date_exclusive, - }, - ) + .get_weekly_reviews(GetReviewsRequest { + from_date: input.from_date, + to_date_exclusive: input.to_date_exclusive, + }) .await?; serialize_output(WeeklyReviewsOutput { reviews: reviews.into_iter().map(WeeklyReviewItem::from).collect(), @@ -170,7 +162,6 @@ mod tests { impl ReviewReadService for StubReviewService { async fn get_daily_reviews( &self, - _ctx: &UserContext, request: GetReviewsRequest, ) -> Result, AnalyzerError> { *self.last_daily.lock().unwrap() = Some(request); @@ -178,7 +169,6 @@ mod tests { } async fn get_weekly_reviews( &self, - _ctx: &UserContext, request: GetReviewsRequest, ) -> Result, AnalyzerError> { *self.last_weekly.lock().unwrap() = Some(request); @@ -186,7 +176,6 @@ mod tests { } async fn get_daily_review( &self, - _ctx: &UserContext, review_date: NaiveDate, ) -> Result, AnalyzerError> { *self.last_daily_single.lock().unwrap() = Some(review_date); @@ -194,7 +183,6 @@ mod tests { } async fn get_weekly_review( &self, - _ctx: &UserContext, week_start: NaiveDate, ) -> Result, AnalyzerError> { *self.last_weekly_single.lock().unwrap() = Some(week_start); @@ -202,10 +190,6 @@ mod tests { } } - fn ctx() -> UserContext { - UserContext::new("user-1") - } - fn ymd(y: i32, m: u32, d: u32) -> NaiveDate { NaiveDate::from_ymd_opt(y, m, d).unwrap() } @@ -225,13 +209,10 @@ mod tests { let tool = DailyReviewGetRangeTool::new(stub.clone()); let out = tool - .dispatch( - &ctx(), - json!({ - "from_date": "2026-04-27", - "to_date_exclusive": "2026-04-29" - }), - ) + .dispatch(json!({ + "from_date": "2026-04-27", + "to_date_exclusive": "2026-04-29" + })) .await .unwrap(); @@ -247,7 +228,7 @@ mod tests { let tool = DailyReviewGetRangeTool::new(Arc::new(StubReviewService::default())); let err = tool - .dispatch(&ctx(), json!({"from_date": "2026-04-27"})) + .dispatch(json!({"from_date": "2026-04-27"})) .await .unwrap_err(); assert!(matches!(err, ToolError::InvalidInput(_))); @@ -265,13 +246,10 @@ mod tests { let tool = WeeklyReviewGetRangeTool::new(stub.clone()); let out = tool - .dispatch( - &ctx(), - json!({ - "from_date": "2026-04-20", - "to_date_exclusive": "2026-04-27" - }), - ) + .dispatch(json!({ + "from_date": "2026-04-20", + "to_date_exclusive": "2026-04-27" + })) .await .unwrap(); diff --git a/src/journal/analyzer/tools/signal.rs b/src/journal/analyzer/tools/signal.rs index 776efc13..b5305d30 100644 --- a/src/journal/analyzer/tools/signal.rs +++ b/src/journal/analyzer/tools/signal.rs @@ -8,7 +8,7 @@ use serde_json::Value; use super::{Tool, ToolError, deserialize_input, schema_value, serialize_output}; use crate::journal::analyzer::signal::SignalReadService; -use crate::journal::analyzer::types::{SearchSignalsRequest, SignalView, UserContext}; +use crate::journal::analyzer::types::{SearchSignalsRequest, SignalView}; use crate::journal::extraction::{BehaviorValence, NeedStatus}; use crate::journal::review::signals::types::SignalType; @@ -95,23 +95,20 @@ impl Tool for SignalsSearchTool { fn input_schema(&self) -> Value { schema_value::() } - async fn dispatch(&self, ctx: &UserContext, args: Value) -> Result { + async fn dispatch(&self, args: Value) -> Result { let input: SearchSignalsInput = deserialize_input(args)?; let signals = self .service - .search( - ctx, - SearchSignalsRequest { - signal_type: input.signal_type, - label_contains: input.label_contains, - status: input.status, - valence: input.valence, - from_date: input.from_date, - to_date_exclusive: input.to_date_exclusive, - min_strength: input.min_strength, - limit: input.limit, - }, - ) + .search(SearchSignalsRequest { + signal_type: input.signal_type, + label_contains: input.label_contains, + status: input.status, + valence: input.valence, + from_date: input.from_date, + to_date_exclusive: input.to_date_exclusive, + min_strength: input.min_strength, + limit: input.limit, + }) .await?; serialize_output(SignalsOutput { signals: signals.into_iter().map(SignalItem::from).collect(), @@ -137,7 +134,6 @@ mod tests { impl SignalReadService for StubSignalService { async fn search( &self, - _ctx: &UserContext, request: SearchSignalsRequest, ) -> Result, AnalyzerError> { *self.last.lock().unwrap() = Some(request); @@ -145,10 +141,6 @@ mod tests { } } - fn ctx() -> UserContext { - UserContext::new("user-1") - } - #[tokio::test] async fn signals_search_dispatches_filters_and_serializes_output() { let stub = Arc::new(StubSignalService::default()); @@ -166,17 +158,14 @@ mod tests { let tool = SignalsSearchTool::new(stub.clone()); let out = tool - .dispatch( - &ctx(), - json!({ - "limit": 10, - "signal_type": "need", - "status": "unmet", - "min_strength": 0.5, - "from_date": "2026-04-01", - "to_date_exclusive": "2026-05-01" - }), - ) + .dispatch(json!({ + "limit": 10, + "signal_type": "need", + "status": "unmet", + "min_strength": 0.5, + "from_date": "2026-04-01", + "to_date_exclusive": "2026-05-01" + })) .await .unwrap(); @@ -202,7 +191,7 @@ mod tests { let stub = Arc::new(StubSignalService::default()); let tool = SignalsSearchTool::new(stub.clone()); - let _ = tool.dispatch(&ctx(), json!({"limit": 5})).await.unwrap(); + let _ = tool.dispatch(json!({"limit": 5})).await.unwrap(); let captured = stub.last.lock().unwrap().clone().unwrap(); assert_eq!(captured.limit, 5); @@ -220,7 +209,7 @@ mod tests { let tool = SignalsSearchTool::new(Arc::new(StubSignalService::default())); let err = tool - .dispatch(&ctx(), json!({"limit": 5, "signal_type": "diagnosis"})) + .dispatch(json!({"limit": 5, "signal_type": "diagnosis"})) .await .unwrap_err(); assert!(matches!(err, ToolError::InvalidInput(_))); diff --git a/src/journal/analyzer/types.rs b/src/journal/analyzer/types.rs index f0a5b9c8..b918860a 100644 --- a/src/journal/analyzer/types.rs +++ b/src/journal/analyzer/types.rs @@ -11,19 +11,6 @@ pub const MAX_TEXT_SEARCH_LIMIT: u32 = 50; pub const MAX_SIGNAL_LIMIT: u32 = 50; pub const MAX_SEMANTIC_LIMIT: u32 = 20; -#[derive(Debug, Clone, PartialEq, Eq)] -pub struct UserContext { - pub user_id: String, -} - -impl UserContext { - pub fn new(user_id: impl Into) -> Self { - Self { - user_id: user_id.into(), - } - } -} - #[derive(Debug, thiserror::Error)] pub enum AnalyzerError { #[error("invalid argument: {0}")] diff --git a/src/journal/analyzer/wiring.rs b/src/journal/analyzer/wiring.rs index a13abd3b..15b0456b 100644 --- a/src/journal/analyzer/wiring.rs +++ b/src/journal/analyzer/wiring.rs @@ -80,7 +80,7 @@ mod tests { use super::*; use crate::database; - use crate::journal::analyzer::types::{AnalyzerError, SemanticHit, UserContext}; + use crate::journal::analyzer::types::{AnalyzerError, SemanticHit}; struct StubSemanticSearcher; @@ -88,7 +88,6 @@ mod tests { impl SemanticJournalSearcher for StubSemanticSearcher { async fn search( &self, - _user_id: &str, _query: &str, _from_date: Option, _to_date_exclusive: Option, @@ -133,11 +132,9 @@ mod tests { async fn registered_tools_are_dispatchable_by_name() { let pool = pool().await; let components = build_analyzer_mcp_components(pool, Arc::new(StubSemanticSearcher)); - let ctx = UserContext::new("u"); - let result = components .registry - .dispatch("journal_get_recent", &ctx, serde_json::json!({"limit": 5})) + .dispatch("journal_get_recent", serde_json::json!({"limit": 5})) .await .unwrap(); assert!(result["entries"].is_array()); diff --git a/src/journal/command.rs b/src/journal/command.rs index 4453beb2..1da83a75 100644 --- a/src/journal/command.rs +++ b/src/journal/command.rs @@ -9,7 +9,6 @@ pub const MAX_RECENT_LIMIT: u32 = 50; pub struct JournalCommandRequest { pub source: MessageSource, pub source_conversation_id: String, - pub user_id: String, pub received_at: DateTime, pub command: JournalCommand, } diff --git a/src/journal/embedding/mod.rs b/src/journal/embedding/mod.rs index 51d162c4..e891c528 100644 --- a/src/journal/embedding/mod.rs +++ b/src/journal/embedding/mod.rs @@ -48,7 +48,6 @@ mod tests { source: MessageSource::Telegram, source_conversation_id: "42".to_string(), source_message_id: source_message_id.to_string(), - user_id: "7".to_string(), text: text.to_string(), received_at, } @@ -177,9 +176,8 @@ mod tests { .map_err(Into::into) } - async fn search_for_user( + async fn search( &self, - user_id: &str, embedding: &Embedding, embedding_model: &str, from_date: Option, @@ -187,8 +185,7 @@ mod tests { limit: usize, ) -> Result>, EmbeddingRepositoryError> { self.inner - .search_for_user( - user_id, + .search( embedding, embedding_model, from_date, diff --git a/src/journal/embedding/repository.rs b/src/journal/embedding/repository.rs index 0685c094..31779c6e 100644 --- a/src/journal/embedding/repository.rs +++ b/src/journal/embedding/repository.rs @@ -56,9 +56,8 @@ pub trait EmbeddingIndex: Send + Sync { embedding_model: &str, ) -> Result; - async fn search_for_user( + async fn search( &self, - user_id: &str, embedding: &Embedding, embedding_model: &str, from_date: Option, @@ -69,9 +68,8 @@ pub trait EmbeddingIndex: Send + Sync { #[async_trait] pub trait PendingEmbeddingCounter: Send + Sync { - async fn count_entries_missing_embedding_for_user( + async fn count_entries_missing_embedding( &self, - user_id: &str, embedding_model: &str, ) -> Result; } @@ -248,9 +246,8 @@ impl SqliteEmbeddingRepository { Ok(count as u32) } - pub async fn count_entries_missing_embedding_for_user( + pub async fn count_entries_missing_embedding( &self, - _user_id: &str, embedding_model: &str, ) -> Result { sqlx::query_scalar( @@ -269,39 +266,10 @@ impl SqliteEmbeddingRepository { .await } - #[cfg(test)] - pub(crate) async fn search( + pub async fn search( &self, embedding: &Embedding, embedding_model: &str, - limit: usize, - ) -> Result>, sqlx::Error> { - let rows = sqlx::query( - r#" - SELECT - m.journal_entry_id, - vec_distance_cosine(v.embedding, ?) AS distance - FROM journal_entry_embedding_metadata m - JOIN journal_entry_embedding_vec v ON v.rowid = m.id - WHERE m.embedding_model = ? - ORDER BY distance ASC - LIMIT ? - "#, - ) - .bind(embedding.to_blob()) - .bind(embedding_model) - .bind(limit as i64) - .fetch_all(&self.pool) - .await?; - - Ok(rows.into_iter().map(map_search_result).collect()) - } - - pub async fn search_for_user( - &self, - _user_id: &str, - embedding: &Embedding, - embedding_model: &str, from_date: Option, to_date_exclusive: Option, limit: usize, @@ -447,18 +415,16 @@ impl EmbeddingIndex for SqliteEmbeddingRepository { .map_err(Into::into) } - async fn search_for_user( + async fn search( &self, - user_id: &str, embedding: &Embedding, embedding_model: &str, from_date: Option, to_date_exclusive: Option, limit: usize, ) -> Result>, EmbeddingRepositoryError> { - SqliteEmbeddingRepository::search_for_user( + SqliteEmbeddingRepository::search( self, - user_id, embedding, embedding_model, from_date, @@ -472,18 +438,13 @@ impl EmbeddingIndex for SqliteEmbeddingRepository { #[async_trait] impl PendingEmbeddingCounter for SqliteEmbeddingRepository { - async fn count_entries_missing_embedding_for_user( + async fn count_entries_missing_embedding( &self, - user_id: &str, embedding_model: &str, ) -> Result { - SqliteEmbeddingRepository::count_entries_missing_embedding_for_user( - self, - user_id, - embedding_model, - ) - .await - .map_err(Into::into) + SqliteEmbeddingRepository::count_entries_missing_embedding(self, embedding_model) + .await + .map_err(Into::into) } } @@ -529,7 +490,6 @@ mod tests { source: MessageSource::Telegram, source_conversation_id: "42".to_string(), source_message_id: source_message_id.to_string(), - user_id: "7".to_string(), text: text.to_string(), received_at, } @@ -961,7 +921,7 @@ mod tests { let query = directional_embedding(1, 1.0); let results = embedding_repository - .search(&query, TEST_EMBEDDING_MODEL, 3) + .search(&query, TEST_EMBEDDING_MODEL, None, None, 3) .await .unwrap(); @@ -991,7 +951,13 @@ mod tests { } let results = embedding_repository - .search(&directional_embedding(0, 1.0), TEST_EMBEDDING_MODEL, 2) + .search( + &directional_embedding(0, 1.0), + TEST_EMBEDDING_MODEL, + None, + None, + 2, + ) .await .unwrap(); @@ -1014,7 +980,7 @@ mod tests { .unwrap(); let results = embedding_repository - .search(&embedding(1.0), "model-b", 10) + .search(&embedding(1.0), "model-b", None, None, 10) .await .unwrap(); @@ -1022,7 +988,7 @@ mod tests { } #[tokio::test] - async fn search_for_user_applies_date_filter_before_limit() { + async fn search_applies_date_filter_before_limit() { let (journal_repository, embedding_repository) = setup().await; let outside = store_entry( &journal_repository, @@ -1059,8 +1025,7 @@ mod tests { .unwrap(); let results = embedding_repository - .search_for_user( - "user-1", + .search( &directional_embedding(0, 1.0), TEST_EMBEDDING_MODEL, NaiveDate::from_ymd_opt(2026, 4, 28), @@ -1079,7 +1044,7 @@ mod tests { let (_, embedding_repository) = setup().await; let results = embedding_repository - .search(&embedding(1.0), TEST_EMBEDDING_MODEL, 10) + .search(&embedding(1.0), TEST_EMBEDDING_MODEL, None, None, 10) .await .unwrap(); diff --git a/src/journal/extraction/backfill.rs b/src/journal/extraction/backfill.rs index c0f96932..9a12569a 100644 --- a/src/journal/extraction/backfill.rs +++ b/src/journal/extraction/backfill.rs @@ -144,7 +144,6 @@ mod tests { source: MessageSource::Telegram, source_conversation_id: "42".to_string(), source_message_id: source_message_id.to_string(), - user_id: "7".to_string(), text: text.to_string(), received_at, }; diff --git a/src/journal/extraction/repository.rs b/src/journal/extraction/repository.rs index d403a419..62138465 100644 --- a/src/journal/extraction/repository.rs +++ b/src/journal/extraction/repository.rs @@ -281,7 +281,6 @@ mod tests { source: MessageSource::Telegram, source_conversation_id: "42".to_string(), source_message_id: "100".to_string(), - user_id: "7".to_string(), text: "hello froid".to_string(), received_at: chrono::Utc::now(), }; @@ -398,7 +397,6 @@ mod tests { source: MessageSource::Telegram, source_conversation_id: "42".to_string(), source_message_id: source_message_id.to_string(), - user_id: "7".to_string(), text: text.to_string(), received_at, }; diff --git a/src/journal/extraction/service.rs b/src/journal/extraction/service.rs index aa448264..5c80baef 100644 --- a/src/journal/extraction/service.rs +++ b/src/journal/extraction/service.rs @@ -237,7 +237,6 @@ mod tests { source: MessageSource::Telegram, source_conversation_id: "42".to_string(), source_message_id: "100".to_string(), - user_id: "7".to_string(), text: text.to_string(), received_at: chrono::Utc::now(), }; diff --git a/src/journal/repository.rs b/src/journal/repository.rs index 71a89717..9352a06f 100644 --- a/src/journal/repository.rs +++ b/src/journal/repository.rs @@ -3,13 +3,12 @@ use std::sync::{Arc, Mutex}; use chrono::{Duration, NaiveDate, TimeZone, Utc}; use sqlx::{Row, SqlitePool, sqlite::SqliteRow}; -use crate::messages::{IncomingMessage, MessageSource, SINGLE_USER_ID}; +use crate::messages::{IncomingMessage, MessageSource}; use super::entry::{JournalEntry, JournalStats, StoredJournalEntry}; #[derive(Debug, Clone, PartialEq, Eq)] pub struct JournalConversation { - pub user_id: String, pub source_conversation_id: String, } @@ -438,7 +437,6 @@ impl JournalRepository { Ok(rows .into_iter() .map(|row| JournalConversation { - user_id: SINGLE_USER_ID.to_string(), source_conversation_id: row.get("source_conversation_id"), }) .collect()) @@ -472,7 +470,6 @@ impl JournalRepository { Ok(rows .into_iter() .map(|row| JournalConversation { - user_id: SINGLE_USER_ID.to_string(), source_conversation_id: row.get("source_conversation_id"), }) .collect()) diff --git a/src/journal/repository_tests.rs b/src/journal/repository_tests.rs index 868319d7..8aaa2e27 100644 --- a/src/journal/repository_tests.rs +++ b/src/journal/repository_tests.rs @@ -30,7 +30,6 @@ fn incoming_for_conversation( source: MessageSource::Telegram, source_conversation_id: source_conversation_id.to_string(), source_message_id: source_message_id.to_string(), - user_id: "7".to_string(), text: text.to_string(), received_at, } @@ -498,7 +497,6 @@ async fn fetch_in_range_returns_all_entries_in_single_user_journal() { source: MessageSource::Telegram, source_conversation_id: "99".to_string(), source_message_id: "2".to_string(), - user_id: "other_user".to_string(), text: "theirs".to_string(), received_at: at_on(2026, 4, 28, 11, 0), }) @@ -637,7 +635,6 @@ async fn search_text_returns_matches_from_the_single_user_journal() { source: MessageSource::Telegram, source_conversation_id: "99".to_string(), source_message_id: "2".to_string(), - user_id: "other_user".to_string(), text: "theirs matches too".to_string(), received_at: at(11, 0), }) @@ -688,7 +685,7 @@ async fn search_text_returns_empty_when_no_match() { } #[tokio::test] -async fn search_text_ignores_caller_user_id() { +async fn search_text_matches_any_stored_entry() { let repo = setup().await; repo.store(&incoming("1", "match", at(10, 0))) @@ -715,7 +712,6 @@ async fn conversations_with_entries_for_date_returns_distinct_source_conversatio source: MessageSource::Telegram, source_conversation_id: "99".to_string(), source_message_id: "3".to_string(), - user_id: "8".to_string(), text: "other user".to_string(), received_at: at(12, 0), }) @@ -725,7 +721,6 @@ async fn conversations_with_entries_for_date_returns_distinct_source_conversatio source: MessageSource::Telegram, source_conversation_id: "100".to_string(), source_message_id: "4".to_string(), - user_id: "9".to_string(), text: "tomorrow".to_string(), received_at: Utc.with_ymd_and_hms(2026, 4, 29, 9, 0, 0).unwrap(), }) @@ -741,11 +736,9 @@ async fn conversations_with_entries_for_date_returns_distinct_source_conversatio conversations, vec![ JournalConversation { - user_id: crate::messages::SINGLE_USER_ID.to_string(), source_conversation_id: "42".to_string(), }, JournalConversation { - user_id: crate::messages::SINGLE_USER_ID.to_string(), source_conversation_id: "99".to_string(), }, ] @@ -801,7 +794,6 @@ async fn fetch_by_ids_returns_all_matching_entries_in_single_user_journal() { source: MessageSource::Telegram, source_conversation_id: "99".to_string(), source_message_id: "2".to_string(), - user_id: "other_user".to_string(), text: "theirs".to_string(), received_at: at(11, 0), }; diff --git a/src/journal/review/embedding_repository.rs b/src/journal/review/embedding_repository.rs index 6111d9f2..62e9f424 100644 --- a/src/journal/review/embedding_repository.rs +++ b/src/journal/review/embedding_repository.rs @@ -156,9 +156,8 @@ impl SqliteDailyReviewEmbeddingRepository { Ok(count as u32) } - async fn search_for_user( + async fn search( &self, - _user_id: &str, embedding: &Embedding, embedding_model: &str, from_date: Option, @@ -259,17 +258,15 @@ impl EmbeddingIndex for SqliteDailyReviewEmbeddingRepository { .map_err(Into::into) } - async fn search_for_user( + async fn search( &self, - user_id: &str, embedding: &Embedding, embedding_model: &str, from_date: Option, to_date_exclusive: Option, limit: usize, ) -> Result>, EmbeddingRepositoryError> { - self.search_for_user( - user_id, + self.search( embedding, embedding_model, from_date, @@ -467,7 +464,7 @@ mod tests { let query = directional_embedding(1); let results = embedding_repo - .search_for_user("user-1", &query, TEST_EMBEDDING_MODEL, None, None, 10) + .search(&query, TEST_EMBEDDING_MODEL, None, None, 10) .await .unwrap(); diff --git a/src/journal/review/mod.rs b/src/journal/review/mod.rs index 33b258a3..588a6c8c 100644 --- a/src/journal/review/mod.rs +++ b/src/journal/review/mod.rs @@ -28,7 +28,6 @@ pub struct JournalEntryWithExtraction { #[derive(Debug, Clone, PartialEq, Eq)] pub struct DailyReview { pub id: i64, - pub user_id: String, pub review_date: NaiveDate, pub review_text: Option, pub model: String, @@ -75,7 +74,6 @@ pub enum DailyReviewResult { #[derive(Debug, Clone, PartialEq, Eq)] pub struct DailyReviewFailure { - pub user_id: String, pub review_date: NaiveDate, pub model: String, pub prompt_version: String, diff --git a/src/journal/review/repository.rs b/src/journal/review/repository.rs index 00bcf75e..9e4ee953 100644 --- a/src/journal/review/repository.rs +++ b/src/journal/review/repository.rs @@ -4,8 +4,6 @@ use thiserror::Error; use crate::errors::from_error_string; -use crate::messages::SINGLE_USER_ID; - use super::{DailyReview, DailyReviewStatus, SignalGenerationStatus}; #[derive(Debug, Clone, PartialEq, Eq, Error)] @@ -361,7 +359,6 @@ fn row_to_daily_review(row: SqliteRow) -> Result Result, DailyReviewSearchError>; } @@ -58,7 +57,6 @@ where { async fn search( &self, - user_id: &str, query: &str, ) -> Result, DailyReviewSearchError> { let embedding = self @@ -71,7 +69,7 @@ where let index_results: Vec> = self .index - .search_for_user(user_id, &embedding, model, None, None, 5) + .search(&embedding, model, None, None, 5) .await .map_err(DailyReviewSearchError::Index)?; @@ -225,9 +223,8 @@ mod tests { unreachable!() } - async fn search_for_user( + async fn search( &self, - _user_id: &str, _embedding: &Embedding, _embedding_model: &str, _from_date: Option, @@ -258,7 +255,7 @@ mod tests { repo, ); - let results = service.search("user-1", "query").await.unwrap(); + let results = service.search("query").await.unwrap(); assert_eq!(results.len(), 1); assert_eq!( @@ -289,7 +286,7 @@ mod tests { repo, ); - let results = service.search("user-1", "query").await.unwrap(); + let results = service.search("query").await.unwrap(); assert_eq!(results.len(), 1); assert_eq!( @@ -307,7 +304,7 @@ mod tests { repo, ); - let results = service.search("user-1", "query").await.unwrap(); + let results = service.search("query").await.unwrap(); assert!(results.is_empty()); } @@ -318,7 +315,7 @@ mod tests { let service = SemanticDailyReviewSearchService::new(index, FakeEmbedder::fails(TEST_MODEL), repo); - let err = service.search("user-1", "query").await.unwrap_err(); + let err = service.search("query").await.unwrap_err(); assert!(matches!(err, DailyReviewSearchError::Embedder(_))); } diff --git a/src/journal/review/service.rs b/src/journal/review/service.rs index 2ec518dd..5e234309 100644 --- a/src/journal/review/service.rs +++ b/src/journal/review/service.rs @@ -43,13 +43,11 @@ pub struct DailyReviewService { pub trait DailyReviewRunner: Send + Sync { async fn review_day( &self, - user_id: &str, utc_date: NaiveDate, ) -> Result; async fn fetch_review( &self, - user_id: &str, utc_date: NaiveDate, ) -> Result, DailyReviewServiceError>; } @@ -74,7 +72,6 @@ impl DailyReviewService { pub async fn review_day( &self, - user_id: &str, utc_date: NaiveDate, ) -> Result { let existing = self.daily_reviews.find_by_user_and_date(utc_date).await?; @@ -89,9 +86,8 @@ impl DailyReviewService { return Ok(DailyReviewResult::Existing(review.clone())); } - let entries_with_extractions: Vec = self - .fetch_entries_with_extractions(user_id, utc_date) - .await?; + let entries_with_extractions: Vec = + self.fetch_entries_with_extractions(utc_date).await?; if entries_with_extractions.is_empty() { return Ok(DailyReviewResult::EmptyDay); } @@ -108,13 +104,7 @@ impl DailyReviewService { let review_text = review_text.trim(); if review_text.is_empty() { return self - .store_failed_review( - user_id, - utc_date, - model, - &prompt_version, - EMPTY_REVIEW_ERROR, - ) + .store_failed_review(utc_date, model, &prompt_version, EMPTY_REVIEW_ERROR) .await; } @@ -127,7 +117,7 @@ impl DailyReviewService { Err(error) => { let prompt_version = self.generator.prompt_version(); let error_message = error.to_string(); - self.store_failed_review(user_id, utc_date, model, &prompt_version, &error_message) + self.store_failed_review(utc_date, model, &prompt_version, &error_message) .await } } @@ -135,7 +125,6 @@ impl DailyReviewService { pub async fn fetch_review( &self, - _user_id: &str, utc_date: NaiveDate, ) -> Result, DailyReviewServiceError> { let review = self.daily_reviews.find_by_user_and_date(utc_date).await?; @@ -149,7 +138,6 @@ impl DailyReviewService { async fn fetch_entries_with_extractions( &self, - _user_id: &str, date: NaiveDate, ) -> Result, DailyReviewServiceError> { let entries = self.journal_entries.fetch_today(date).await?; @@ -178,7 +166,6 @@ impl DailyReviewService { async fn store_failed_review( &self, - user_id: &str, utc_date: NaiveDate, model: &str, prompt_version: &str, @@ -188,7 +175,6 @@ impl DailyReviewService { .upsert_failed(utc_date, model, prompt_version, error_message) .await?; Ok(DailyReviewResult::GenerationFailed(DailyReviewFailure { - user_id: user_id.to_string(), review_date: utc_date, model: model.to_string(), prompt_version: prompt_version.to_string(), @@ -201,18 +187,16 @@ impl DailyReviewService { impl DailyReviewRunner for DailyReviewService { async fn review_day( &self, - user_id: &str, utc_date: NaiveDate, ) -> Result { - DailyReviewService::review_day(self, user_id, utc_date).await + DailyReviewService::review_day(self, utc_date).await } async fn fetch_review( &self, - user_id: &str, utc_date: NaiveDate, ) -> Result, DailyReviewServiceError> { - DailyReviewService::fetch_review(self, user_id, utc_date).await + DailyReviewService::fetch_review(self, utc_date).await } } @@ -272,7 +256,6 @@ mod tests { } fn incoming( - user_id: &str, source_message_id: &str, text: &str, received_at: chrono::DateTime, @@ -281,7 +264,6 @@ mod tests { source: MessageSource::Telegram, source_conversation_id: "42".to_string(), source_message_id: source_message_id.to_string(), - user_id: user_id.to_string(), text: text.to_string(), received_at, } @@ -289,7 +271,6 @@ mod tests { fn at_date(day: u32, source_message_id: &str, text: &str) -> IncomingMessage { incoming( - "user-1", source_message_id, text, Utc.with_ymd_and_hms(2026, 4, day, 10, 0, 0).unwrap(), @@ -312,7 +293,7 @@ mod tests { .await .unwrap(); - let result = service.review_day("user-1", date()).await.unwrap(); + let result = service.review_day(date()).await.unwrap(); assert_eq!(result, DailyReviewResult::Existing(existing)); assert_eq!(generator.calls(), 0); @@ -331,7 +312,7 @@ mod tests { .await .unwrap(); - let review = generated_review(service.review_day("user-1", date()).await.unwrap()); + let review = generated_review(service.review_day(date()).await.unwrap()); assert_eq!(review.review_text, Some("generated review".to_string())); assert_eq!(review.status, DailyReviewStatus::Completed); @@ -357,12 +338,11 @@ mod tests { .await .unwrap(); - let result = service.review_day("user-1", date()).await.unwrap(); + let result = service.review_day(date()).await.unwrap(); assert_eq!( result, DailyReviewResult::GenerationFailed(DailyReviewFailure { - user_id: "user-1".to_string(), review_date: date(), model: "fake-review-model".to_string(), prompt_version: "fake-prompt-v1".to_string(), @@ -397,7 +377,7 @@ mod tests { .await .unwrap(); - let review = generated_review(service.review_day("user-1", date()).await.unwrap()); + let review = generated_review(service.review_day(date()).await.unwrap()); assert_eq!(generator.calls(), 1); assert_eq!(review.id, existing.id); @@ -410,7 +390,7 @@ mod tests { let (service, _daily_reviews, _journal_entries, _extractions, generator) = setup(FakeReviewGenerator::succeeding("generated review")).await; - let result = service.review_day("user-1", date()).await.unwrap(); + let result = service.review_day(date()).await.unwrap(); assert_eq!(result, DailyReviewResult::EmptyDay); assert_eq!(generator.calls(), 0); @@ -425,12 +405,11 @@ mod tests { .await .unwrap(); - let result = service.review_day("user-1", date()).await.unwrap(); + let result = service.review_day(date()).await.unwrap(); assert_eq!( result, DailyReviewResult::GenerationFailed(DailyReviewFailure { - user_id: "user-1".to_string(), review_date: date(), model: "fake-review-model".to_string(), prompt_version: "fake-prompt-v1".to_string(), @@ -465,7 +444,7 @@ mod tests { .unwrap(); assert!(matches!( - service.review_day("user-1", date()).await.unwrap(), + service.review_day(date()).await.unwrap(), DailyReviewResult::GenerationFailed(_) )); let failed = daily_reviews @@ -475,7 +454,7 @@ mod tests { .unwrap(); std::thread::sleep(std::time::Duration::from_millis(2)); - let completed = generated_review(service.review_day("user-1", date()).await.unwrap()); + let completed = generated_review(service.review_day(date()).await.unwrap()); assert_eq!(generator.calls(), 2); assert_eq!(completed.id, failed.id); @@ -487,7 +466,7 @@ mod tests { } #[tokio::test] - async fn same_date_reviews_share_single_user_scope() { + async fn repeated_review_day_returns_existing_without_regenerating() { let (service, _daily_reviews, journal_entries, _extractions, generator) = setup(FakeReviewGenerator::new(vec![ Ok("user one review".to_string()), @@ -496,7 +475,6 @@ mod tests { .await; journal_entries .store(&incoming( - "user-1", "1", "user one entry", Utc.with_ymd_and_hms(2026, 4, 28, 10, 0, 0).unwrap(), @@ -505,7 +483,6 @@ mod tests { .unwrap(); journal_entries .store(&incoming( - "user-2", "2", "user two entry", Utc.with_ymd_and_hms(2026, 4, 28, 11, 0, 0).unwrap(), @@ -513,15 +490,14 @@ mod tests { .await .unwrap(); - let user_one = generated_review(service.review_day("user-1", date()).await.unwrap()); - let user_two = match service.review_day("user-2", date()).await.unwrap() { + let first = generated_review(service.review_day(date()).await.unwrap()); + let second = match service.review_day(date()).await.unwrap() { DailyReviewResult::Existing(review) => review, other => panic!("expected existing review, got {other:?}"), }; - assert_eq!(user_one.review_text, Some("user one review".to_string())); - assert_eq!(user_two.review_text, Some("user one review".to_string())); - assert_eq!(user_one.user_id, user_two.user_id); + assert_eq!(first.review_text, Some("user one review".to_string())); + assert_eq!(second.review_text, Some("user one review".to_string())); assert_eq!(generator.calls(), 1); } @@ -542,7 +518,7 @@ mod tests { .await .unwrap(); - let review = generated_review(service.review_day("user-1", date()).await.unwrap()); + let review = generated_review(service.review_day(date()).await.unwrap()); assert_eq!(review.review_date, date()); assert_eq!(generator.calls(), 1); @@ -570,8 +546,8 @@ mod tests { .await .unwrap(); - let first = generated_review(service.review_day("user-1", date()).await.unwrap()); - let second = generated_review(service.review_day("user-1", next_date).await.unwrap()); + let first = generated_review(service.review_day(date()).await.unwrap()); + let second = generated_review(service.review_day(next_date).await.unwrap()); assert_eq!(first.review_date, date()); assert_eq!(second.review_date, next_date); @@ -588,7 +564,7 @@ mod tests { .await .unwrap(); - let result = service.fetch_review("user-1", date()).await.unwrap(); + let result = service.fetch_review(date()).await.unwrap(); assert_eq!(result.unwrap().review_text, Some("review text".to_string())); } @@ -598,7 +574,7 @@ mod tests { let (service, _daily_reviews, _journal_entries, _extractions, _generator) = setup(FakeReviewGenerator::succeeding("any")).await; - let result = service.fetch_review("user-1", date()).await.unwrap(); + let result = service.fetch_review(date()).await.unwrap(); assert!(result.is_none()); } @@ -612,7 +588,7 @@ mod tests { .await .unwrap(); - let result = service.fetch_review("user-1", date()).await.unwrap(); + let result = service.fetch_review(date()).await.unwrap(); assert!(result.is_none()); } @@ -637,7 +613,7 @@ mod tests { PoolClosingGenerator { pool }, ); - let error = service.review_day("user-1", date()).await.unwrap_err(); + let error = service.review_day(date()).await.unwrap_err(); assert!(matches!(error, DailyReviewServiceError::Storage(_))); } @@ -672,7 +648,6 @@ mod tests { let entry_id = journal_entries .store(&incoming( - "user-1", "1", "entry with extraction", Utc.with_ymd_and_hms(2026, 4, 28, 10, 0, 0).unwrap(), @@ -682,7 +657,6 @@ mod tests { .unwrap(); journal_entries .store(&incoming( - "user-1", "2", "entry without extraction", Utc.with_ymd_and_hms(2026, 4, 28, 10, 1, 0).unwrap(), @@ -712,7 +686,7 @@ mod tests { .await .unwrap(); - service.review_day("user-1", date()).await.unwrap(); + service.review_day(date()).await.unwrap(); let seen = generator.entries_seen(); assert_eq!(seen.len(), 1); diff --git a/src/journal/review/signals/backfill.rs b/src/journal/review/signals/backfill.rs index 1033536f..8c89c221 100644 --- a/src/journal/review/signals/backfill.rs +++ b/src/journal/review/signals/backfill.rs @@ -57,15 +57,11 @@ impl DailyReviewSignalBackfillService { remaining: 0, }; - for (daily_review_id, user_id, review_date) in candidates { - if let Err(error) = self - .process_candidate(&user_id, review_date, daily_review_id) - .await - { + for (daily_review_id, review_date) in candidates { + if let Err(error) = self.process_candidate(review_date, daily_review_id).await { result.errored += 1; warn!( daily_review_id, - user_id = %user_id, review_date = %review_date, error = %error, "signal backfill candidate failed" @@ -83,13 +79,12 @@ impl DailyReviewSignalBackfillService { async fn process_candidate( &self, - user_id: &str, review_date: NaiveDate, daily_review_id: i64, ) -> Result<(), ProcessCandidateError> { let result = self .service - .generate_signals_for_review(user_id, review_date) + .generate_signals_for_review(review_date) .await .map_err(ProcessCandidateError::Service)?; @@ -189,7 +184,6 @@ mod tests { source: MessageSource::Telegram, source_conversation_id: "42".to_string(), source_message_id: msg_id.to_string(), - user_id: "user-1".to_string(), text: text.to_string(), received_at: chrono::Utc::now(), }; diff --git a/src/journal/review/signals/repository.rs b/src/journal/review/signals/repository.rs index 73e45101..ec87161c 100644 --- a/src/journal/review/signals/repository.rs +++ b/src/journal/review/signals/repository.rs @@ -4,10 +4,7 @@ use thiserror::Error; use crate::errors::from_error_string; -use crate::{ - journal::extraction::{BehaviorValence, NeedStatus}, - messages::SINGLE_USER_ID, -}; +use crate::journal::extraction::{BehaviorValence, NeedStatus}; use super::types::{DailyReviewSignal, DailyReviewSignalCandidate, SignalType}; @@ -115,7 +112,7 @@ impl DailyReviewSignalRepository { pub async fn find_completed_reviews_missing_signals( &self, limit: u32, - ) -> Result, DailyReviewSignalRepositoryError> { + ) -> Result, DailyReviewSignalRepositoryError> { let rows = sqlx::query( r#" SELECT dr.id, dr.review_date @@ -139,11 +136,7 @@ impl DailyReviewSignalRepository { NaiveDate::parse_from_str(&review_date_str, "%Y-%m-%d").map_err(|_| { DailyReviewSignalRepositoryError::InvalidReviewDate(review_date_str) })?; - Ok(( - row.get::("id"), - SINGLE_USER_ID.to_string(), - review_date, - )) + Ok((row.get::("id"), review_date)) }) .collect() } @@ -298,7 +291,6 @@ fn row_to_signal(row: SqliteRow) -> Result i64 { - let user_id = "user-1"; let msg = IncomingMessage { source: MessageSource::Telegram, source_conversation_id: "42".to_string(), source_message_id: "1".to_string(), - user_id: user_id.to_string(), text: "entry text".to_string(), received_at: chrono::Utc::now(), }; @@ -524,7 +514,7 @@ mod tests { assert_eq!(candidates.len(), 1); } - async fn insert_daily_review_for(pool: &SqlitePool, _user_id: &str, date: NaiveDate) -> i64 { + async fn insert_daily_review_for(pool: &SqlitePool, date: NaiveDate) -> i64 { let review_repo = DailyReviewRepository::new(pool.clone()); review_repo .upsert_completed(date, "review text", "model", "v1") @@ -540,9 +530,9 @@ mod tests { let wednesday = NaiveDate::from_ymd_opt(2026, 4, 29).unwrap(); let outside = NaiveDate::from_ymd_opt(2026, 5, 4).unwrap(); - let monday_review = insert_daily_review_for(&pool, "user-1", monday).await; - let wednesday_review = insert_daily_review_for(&pool, "user-1", wednesday).await; - let outside_review = insert_daily_review_for(&pool, "user-1", outside).await; + let monday_review = insert_daily_review_for(&pool, monday).await; + let wednesday_review = insert_daily_review_for(&pool, wednesday).await; + let outside_review = insert_daily_review_for(&pool, outside).await; repo.replace_in_transaction(monday_review, monday, &[theme_candidate()], "m", "v1") .await @@ -586,9 +576,9 @@ mod tests { let tuesday = NaiveDate::from_ymd_opt(2026, 4, 28).unwrap(); let wednesday = NaiveDate::from_ymd_opt(2026, 4, 29).unwrap(); - let monday_review = insert_daily_review_for(pool, "user-1", monday).await; - let tuesday_review = insert_daily_review_for(pool, "user-1", tuesday).await; - let wednesday_review = insert_daily_review_for(pool, "user-1", wednesday).await; + let monday_review = insert_daily_review_for(pool, monday).await; + let tuesday_review = insert_daily_review_for(pool, tuesday).await; + let wednesday_review = insert_daily_review_for(pool, wednesday).await; let repo = DailyReviewSignalRepository::new(pool.clone()); repo.replace_in_transaction(monday_review, monday, &[theme_candidate()], "m", "v1") @@ -747,8 +737,8 @@ mod tests { async fn search_uses_single_user_scope() { let (repo, _reviews, pool) = setup().await; let target = NaiveDate::from_ymd_opt(2026, 4, 27).unwrap(); - let user_one_review = insert_daily_review_for(&pool, "user-1", target).await; - let user_two_review = insert_daily_review_for(&pool, "user-2", target).await; + let user_one_review = insert_daily_review_for(&pool, target).await; + let user_two_review = insert_daily_review_for(&pool, target).await; repo.replace_in_transaction(user_one_review, target, &[theme_candidate()], "m", "v1") .await @@ -760,7 +750,6 @@ mod tests { let rows = repo.search(&filters(10)).await.unwrap(); assert_eq!(rows.len(), 1); - assert!(rows.iter().all(|row| row.user_id == SINGLE_USER_ID)); } #[tokio::test] @@ -802,8 +791,8 @@ mod tests { async fn find_by_user_in_range_uses_single_user_scope() { let (repo, _reviews, pool) = setup().await; let target = NaiveDate::from_ymd_opt(2026, 4, 27).unwrap(); - let user_one_review = insert_daily_review_for(&pool, "user-1", target).await; - let user_two_review = insert_daily_review_for(&pool, "user-2", target).await; + let user_one_review = insert_daily_review_for(&pool, target).await; + let user_two_review = insert_daily_review_for(&pool, target).await; repo.replace_in_transaction(user_one_review, target, &[theme_candidate()], "m", "v1") .await @@ -818,6 +807,5 @@ mod tests { .unwrap(); assert_eq!(rows.len(), 1); - assert!(rows.iter().all(|row| row.user_id == SINGLE_USER_ID)); } } diff --git a/src/journal/review/signals/service.rs b/src/journal/review/signals/service.rs index 6158ae78..d9a24203 100644 --- a/src/journal/review/signals/service.rs +++ b/src/journal/review/signals/service.rs @@ -23,11 +23,8 @@ use crate::journal::{ #[derive(Debug, Clone, PartialEq, Eq, Error)] pub enum DailyReviewSignalServiceError { - #[error("no completed daily review found for user '{user_id}' on {review_date}")] - NoDailyReview { - user_id: String, - review_date: NaiveDate, - }, + #[error("no completed daily review found on {review_date}")] + NoDailyReview { review_date: NaiveDate }, #[error("{0}")] Storage(String), } @@ -86,7 +83,6 @@ impl DailyReviewSignalService { pub async fn generate_signals_for_review( &self, - user_id: &str, review_date: NaiveDate, ) -> Result { let review = self @@ -112,9 +108,7 @@ impl DailyReviewSignalService { .mark_signals_pending(review.id, self.generator.model(), &pending_prompt_version) .await?; - let entries = self - .fetch_entries_with_extractions(user_id, review_date) - .await?; + let entries = self.fetch_entries_with_extractions(review_date).await?; let generation_result = self .generator @@ -186,7 +180,6 @@ impl DailyReviewSignalService { pub async fn fetch_signals( &self, - _user_id: &str, review_date: NaiveDate, ) -> Result, DailyReviewSignalServiceError> { Ok(self.signals.find_by_user_and_date(review_date).await?) @@ -194,7 +187,6 @@ impl DailyReviewSignalService { async fn fetch_entries_with_extractions( &self, - _user_id: &str, date: NaiveDate, ) -> Result< Vec, @@ -280,12 +272,11 @@ mod tests { NaiveDate::from_ymd_opt(2026, 4, 28).unwrap() } - fn incoming(user_id: &str, text: &str) -> IncomingMessage { + fn incoming(text: &str) -> IncomingMessage { IncomingMessage { source: MessageSource::Telegram, source_conversation_id: "42".to_string(), source_message_id: "1".to_string(), - user_id: user_id.to_string(), text: text.to_string(), received_at: Utc.with_ymd_and_hms(2026, 4, 28, 10, 0, 0).unwrap(), } @@ -323,10 +314,7 @@ mod tests { async fn returns_no_daily_review_when_review_does_not_exist() { let (service, _, _) = setup(FakeSignalGenerator::succeeding(output_with(vec![]))).await; - let result = service - .generate_signals_for_review("user-1", date()) - .await - .unwrap(); + let result = service.generate_signals_for_review(date()).await.unwrap(); assert_eq!(result, DailyReviewSignalResult::NoDailyReview); } @@ -340,10 +328,7 @@ mod tests { .await .unwrap(); - let result = service - .generate_signals_for_review("user-1", date()) - .await - .unwrap(); + let result = service.generate_signals_for_review(date()).await.unwrap(); assert_eq!(result, DailyReviewSignalResult::NoDailyReview); } @@ -352,23 +337,17 @@ mod tests { async fn generates_and_stores_signals_for_completed_review() { let generator = FakeSignalGenerator::succeeding(output_with(vec![theme_signal()])); let (service, reviews, entries) = setup(generator).await; - entries - .store(&incoming("user-1", "entry text")) - .await - .unwrap(); + entries.store(&incoming("entry text")).await.unwrap(); reviews .upsert_completed(date(), "review text", "model", "v1") .await .unwrap(); - let result = service - .generate_signals_for_review("user-1", date()) - .await - .unwrap(); + let result = service.generate_signals_for_review(date()).await.unwrap(); assert_eq!(result, DailyReviewSignalResult::Generated { count: 1 }); - let signals = service.fetch_signals("user-1", date()).await.unwrap(); + let signals = service.fetch_signals(date()).await.unwrap(); assert_eq!(signals.len(), 1); assert_eq!(signals[0].signal_type, SignalType::Theme); assert_eq!(signals[0].label, "physical appearance"); @@ -383,16 +362,10 @@ mod tests { .await .unwrap(); - service - .generate_signals_for_review("user-1", date()) - .await - .unwrap(); - service - .generate_signals_for_review("user-1", date()) - .await - .unwrap(); + service.generate_signals_for_review(date()).await.unwrap(); + service.generate_signals_for_review(date()).await.unwrap(); - let signals = service.fetch_signals("user-1", date()).await.unwrap(); + let signals = service.fetch_signals(date()).await.unwrap(); assert_eq!(signals.len(), 1, "second run must not duplicate signals"); } @@ -405,10 +378,7 @@ mod tests { .await .unwrap(); - let result = service - .generate_signals_for_review("user-1", date()) - .await - .unwrap(); + let result = service.generate_signals_for_review(date()).await.unwrap(); assert_eq!( result, @@ -417,7 +387,7 @@ mod tests { } ); - let signals = service.fetch_signals("user-1", date()).await.unwrap(); + let signals = service.fetch_signals(date()).await.unwrap(); assert!(signals.is_empty()); } @@ -434,14 +404,11 @@ mod tests { .await .unwrap(); - let result = service - .generate_signals_for_review("user-1", date()) - .await - .unwrap(); + let result = service.generate_signals_for_review(date()).await.unwrap(); assert_eq!(result, DailyReviewSignalResult::Generated { count: 1 }); - let signals = service.fetch_signals("user-1", date()).await.unwrap(); + let signals = service.fetch_signals(date()).await.unwrap(); assert_eq!(signals.len(), 1); assert_eq!(signals[0].signal_type, SignalType::Need); } @@ -459,19 +426,13 @@ mod tests { .await .unwrap(); - service - .generate_signals_for_review("user-1", date()) - .await - .unwrap(); - service - .generate_signals_for_review("user-2", date()) - .await - .unwrap(); + service.generate_signals_for_review(date()).await.unwrap(); + service.generate_signals_for_review(date()).await.unwrap(); - let user_one = service.fetch_signals("user-1", date()).await.unwrap(); - let user_two = service.fetch_signals("user-2", date()).await.unwrap(); + let user_one = service.fetch_signals(date()).await.unwrap(); + let user_two = service.fetch_signals(date()).await.unwrap(); let user_one_other_date = service - .fetch_signals("user-1", NaiveDate::from_ymd_opt(2026, 4, 29).unwrap()) + .fetch_signals(NaiveDate::from_ymd_opt(2026, 4, 29).unwrap()) .await .unwrap(); @@ -489,20 +450,14 @@ mod tests { .await .unwrap(); - service - .generate_signals_for_review("user-1", date()) - .await - .unwrap(); + service.generate_signals_for_review(date()).await.unwrap(); assert_eq!(generator.calls(), 1); // Second run replaces with same signals — still only 1 total - service - .generate_signals_for_review("user-1", date()) - .await - .unwrap(); + service.generate_signals_for_review(date()).await.unwrap(); assert_eq!(generator.calls(), 2); - let signals = service.fetch_signals("user-1", date()).await.unwrap(); + let signals = service.fetch_signals(date()).await.unwrap(); assert_eq!(signals.len(), 1); } } diff --git a/src/journal/review/signals/types.rs b/src/journal/review/signals/types.rs index eedb59db..4b1aa533 100644 --- a/src/journal/review/signals/types.rs +++ b/src/journal/review/signals/types.rs @@ -66,7 +66,6 @@ pub struct DailyReviewSignalsOutput { pub struct DailyReviewSignal { pub id: i64, pub daily_review_id: i64, - pub user_id: String, pub review_date: NaiveDate, pub signal_type: SignalType, pub label: String, diff --git a/src/journal/review/wiring.rs b/src/journal/review/wiring.rs index da70cd70..d35b9120 100644 --- a/src/journal/review/wiring.rs +++ b/src/journal/review/wiring.rs @@ -134,7 +134,6 @@ mod tests { .command(&JournalCommandRequest { source: MessageSource::Telegram, source_conversation_id: "42".to_string(), - user_id: "7".to_string(), received_at: Utc::now(), command: JournalCommand::DayReviewLast, }) @@ -217,7 +216,6 @@ mod tests { .command(&JournalCommandRequest { source: MessageSource::Telegram, source_conversation_id: "42".to_string(), - user_id: "7".to_string(), received_at: Utc::now(), command: JournalCommand::DayReviewLast, }) @@ -266,7 +264,6 @@ mod tests { .command(&JournalCommandRequest { source: MessageSource::Telegram, source_conversation_id: "42".to_string(), - user_id: "7".to_string(), received_at: Utc::now(), command: JournalCommand::Status, }) diff --git a/src/journal/search.rs b/src/journal/search.rs index ba6ef6a4..201e8d21 100644 --- a/src/journal/search.rs +++ b/src/journal/search.rs @@ -32,11 +32,7 @@ pub enum SemanticSearchError { #[async_trait] pub(crate) trait SearchService: Send + Sync { - async fn search( - &self, - user_id: &str, - query: &str, - ) -> Result, SemanticSearchError>; + async fn search(&self, query: &str) -> Result, SemanticSearchError>; } #[derive(Clone)] @@ -66,11 +62,7 @@ where I: EmbeddingIndex + Send + Sync, E: Embedder + Send + Sync, { - async fn search( - &self, - user_id: &str, - query: &str, - ) -> Result, SemanticSearchError> { + async fn search(&self, query: &str) -> Result, SemanticSearchError> { let embedding = self .embedder .embed(query) @@ -81,7 +73,7 @@ where let index_results: Vec> = self .index - .search_for_user(user_id, &embedding, model, None, None, MAX_SEARCH_LIMIT) + .search(&embedding, model, None, None, MAX_SEARCH_LIMIT) .await .map_err(SemanticSearchError::Index)?; @@ -185,7 +177,6 @@ mod tests { source: MessageSource::Telegram, source_conversation_id: "42".to_string(), source_message_id: msg_id.to_string(), - user_id: "7".to_string(), text: text.to_string(), received_at, } @@ -304,9 +295,8 @@ mod tests { unreachable!("search tests do not count missing through FakeIndex") } - async fn search_for_user( + async fn search( &self, - _user_id: &str, _embedding: &Embedding, _embedding_model: &str, _from_date: Option, @@ -360,7 +350,7 @@ mod tests { let service = make_service(index, FakeEmbedder::succeeds(TEST_MODEL, 1), repo); - let results = service.search("7", "query").await.unwrap(); + let results = service.search("query").await.unwrap(); assert!(!results.is_empty()); assert_eq!(results[0].journal_entry.text, "most relevant entry"); @@ -376,7 +366,7 @@ mod tests { let service = make_service(index, FakeEmbedder::succeeds(TEST_MODEL, 0), repo); - let results = service.search("7", "query").await.unwrap(); + let results = service.search("query").await.unwrap(); assert!(results.is_empty()); } @@ -392,7 +382,7 @@ mod tests { let service = make_service(index, FakeEmbedder::succeeds(TEST_MODEL, 0), repo); - let results = service.search("7", "query").await.unwrap(); + let results = service.search("query").await.unwrap(); assert_eq!(results.len(), 1); assert_eq!(results[0].journal_entry.text, "embedded entry"); @@ -404,7 +394,7 @@ mod tests { let service = make_service(index, FakeEmbedder::fails(TEST_MODEL), repo); - let error = service.search("7", "query").await.unwrap_err(); + let error = service.search("query").await.unwrap_err(); assert!(matches!(error, SemanticSearchError::Embedder(_))); } @@ -419,7 +409,7 @@ mod tests { // Query with a model that has no stored embeddings. let service = make_service(index, FakeEmbedder::succeeds("other-model", 0), repo); - let results = service.search("7", "query").await.unwrap(); + let results = service.search("query").await.unwrap(); assert!(results.is_empty()); } @@ -435,7 +425,7 @@ mod tests { let service = make_service(index, FakeEmbedder::succeeds(TEST_MODEL, 1), repo); - let results = service.search("7", "query").await.unwrap(); + let results = service.search("query").await.unwrap(); assert_eq!(results[0].journal_entry.text, "entry B"); for window in results.windows(2) { @@ -455,7 +445,6 @@ mod tests { source: MessageSource::Telegram, source_conversation_id: "99".to_string(), source_message_id: "2".to_string(), - user_id: "other_user".to_string(), text: "other user entry".to_string(), received_at: at(11, 0), }; @@ -479,7 +468,7 @@ mod tests { // Query at dim 1: in single-user mode, every stored entry is part of the same journal. let service = make_service(index, FakeEmbedder::succeeds(TEST_MODEL, 1), repo); - let results = service.search("7", "query").await.unwrap(); + let results = service.search("query").await.unwrap(); assert_eq!(results.len(), 2); assert_eq!(results[0].journal_entry.text, "other user entry"); @@ -497,7 +486,6 @@ mod tests { source: MessageSource::Telegram, source_conversation_id: format!("other-{i}"), source_message_id: format!("other-{i}"), - user_id: "other_user".to_string(), text: format!("other user entry {i}"), received_at: at(11, i as u32), }; @@ -522,7 +510,7 @@ mod tests { let service = make_service(index, FakeEmbedder::succeeds(TEST_MODEL, 1), repo); - let results = service.search("7", "query").await.unwrap(); + let results = service.search("query").await.unwrap(); assert_eq!(results.len(), DEFAULT_SEARCH_LIMIT); assert!( @@ -558,7 +546,7 @@ mod tests { }; let service = make_fake_index_service(index, FakeEmbedder::succeeds(TEST_MODEL, 0), repo); - let results = service.search("7", "query").await.unwrap(); + let results = service.search("query").await.unwrap(); assert_eq!(results.len(), 1); assert_eq!(results[0].journal_entry.text, "kept entry"); @@ -583,7 +571,7 @@ mod tests { let service = make_service(index, FakeEmbedder::succeeds(TEST_MODEL, 0), repo); - let results = service.search("7", "query").await.unwrap(); + let results = service.search("query").await.unwrap(); assert_eq!(results.len(), DEFAULT_SEARCH_LIMIT); } diff --git a/src/journal/service/commands.rs b/src/journal/service/commands.rs index 1fc4d121..0115a320 100644 --- a/src/journal/service/commands.rs +++ b/src/journal/service/commands.rs @@ -47,43 +47,29 @@ impl JournalService { }), JournalCommand::Last => self.last(request).await, JournalCommand::Undo => self.undo(request).await, - JournalCommand::Recent { requested_limit } => { - self.recent(&request.user_id, *requested_limit).await - } + JournalCommand::Recent { requested_limit } => self.recent(*requested_limit).await, JournalCommand::RecentUsage => Ok(OutgoingMessage { text: recent_usage_response(), }), - JournalCommand::Today => { - self.today(&request.user_id, request.received_at.date_naive()) - .await - } - JournalCommand::Stats => { - self.stats(&request.user_id, request.received_at.date_naive()) - .await - } - JournalCommand::Status => { - self.status(&request.user_id, request.received_at.date_naive()) - .await + JournalCommand::Today => self.today(request.received_at.date_naive()).await, + JournalCommand::Stats => self.stats(request.received_at.date_naive()).await, + JournalCommand::Status => self.status(request.received_at.date_naive()).await, + JournalCommand::DayReviewLast => { + Ok(self.day_review_last(request.received_at.date_naive()).await) } - JournalCommand::DayReviewLast => Ok(self - .day_review_last(&request.user_id, request.received_at.date_naive()) - .await), JournalCommand::WeekReviewLast => Ok(self - .week_review_last(&request.user_id, request.received_at.date_naive()) + .week_review_last(request.received_at.date_naive()) .await), - JournalCommand::Search { query } => { - Ok(self.search_command(&request.user_id, query).await) - } + JournalCommand::Search { query } => Ok(self.search_command(query).await), JournalCommand::SearchUsage => Ok(OutgoingMessage { text: search_usage_response(), }), } } - async fn day_review_last(&self, user_id: &str, today: chrono::NaiveDate) -> OutgoingMessage { + async fn day_review_last(&self, today: chrono::NaiveDate) -> OutgoingMessage { let yesterday = today - Duration::days(1); self.run_review( - user_id, yesterday, |r| format_daily_review_for_date(r, yesterday), daily_review_not_available_for_date_response(yesterday), @@ -93,7 +79,6 @@ impl JournalService { async fn run_review( &self, - user_id: &str, date: chrono::NaiveDate, format_review: impl Fn(&DailyReview) -> String, not_found_text: String, @@ -104,7 +89,7 @@ impl JournalService { }; }; - match daily_review.fetch_review(user_id, date).await { + match daily_review.fetch_review(date).await { Ok(Some(review)) => OutgoingMessage { text: format_review(&review), }, @@ -120,7 +105,7 @@ impl JournalService { } } - async fn week_review_last(&self, user_id: &str, today: NaiveDate) -> OutgoingMessage { + async fn week_review_last(&self, today: NaiveDate) -> OutgoingMessage { let Some(weekly_review) = &self.weekly_review else { return OutgoingMessage { text: weekly_review_unavailable_response(), @@ -129,7 +114,7 @@ impl JournalService { let week_start = previous_iso_week_monday(today); - match weekly_review.fetch_review(user_id, week_start).await { + match weekly_review.fetch_review(week_start).await { Ok(Some(review)) => OutgoingMessage { text: format_weekly_review_for_week(&review, week_start), }, @@ -145,14 +130,14 @@ impl JournalService { } } - async fn search_command(&self, user_id: &str, query: &str) -> OutgoingMessage { + async fn search_command(&self, query: &str) -> OutgoingMessage { let Some(search) = &self.search else { return OutgoingMessage { text: search_unavailable_response(), }; }; - match search.search(user_id, query).await { + match search.search(query).await { Ok(results) if results.is_empty() => OutgoingMessage { text: search_empty_response(), }, @@ -197,7 +182,7 @@ impl JournalService { }) } - async fn recent(&self, _user_id: &str, limit: u32) -> Result { + async fn recent(&self, limit: u32) -> Result { let limit = limit.min(MAX_RECENT_LIMIT); let entries = self.repository.fetch_recent(limit).await?; @@ -212,11 +197,7 @@ impl JournalService { }) } - async fn today( - &self, - _user_id: &str, - date: chrono::NaiveDate, - ) -> Result { + async fn today(&self, date: chrono::NaiveDate) -> Result { let entries = self.repository.fetch_today(date).await?; if entries.is_empty() { @@ -230,11 +211,7 @@ impl JournalService { }) } - async fn stats( - &self, - _user_id: &str, - today: chrono::NaiveDate, - ) -> Result { + async fn stats(&self, today: chrono::NaiveDate) -> Result { let stats = self.repository.stats(today).await?; Ok(OutgoingMessage { @@ -242,13 +219,9 @@ impl JournalService { }) } - async fn status( - &self, - user_id: &str, - today: chrono::NaiveDate, - ) -> Result { + async fn status(&self, today: chrono::NaiveDate) -> Result { let journal = self.repository.stats(today).await?; - let embeddings = self.embedding_status(user_id).await; + let embeddings = self.embedding_status().await; let daily_review = self.daily_review_status(); Ok(OutgoingMessage { @@ -260,7 +233,7 @@ impl JournalService { }) } - async fn embedding_status(&self, user_id: &str) -> EmbeddingStatus { + async fn embedding_status(&self) -> EmbeddingStatus { let semantic_search = if self.search.is_some() && self.embedding_status_config.is_some() { SemanticSearchStatus::Enabled } else { @@ -272,10 +245,7 @@ impl JournalService { self.pending_embedding_counter.as_ref(), ) { (Some(config), Some(counter)) => { - match counter - .count_entries_missing_embedding_for_user(user_id, &config.model) - .await - { + match counter.count_entries_missing_embedding(&config.model).await { Ok(count) => Some(count), Err(error) => { warn!(%error, "failed to count pending embeddings for status"); diff --git a/src/journal/service/tests.rs b/src/journal/service/tests.rs index e99b06ba..6712926e 100644 --- a/src/journal/service/tests.rs +++ b/src/journal/service/tests.rs @@ -158,7 +158,7 @@ const TEST_MODEL: &str = "test-model"; #[derive(Debug, Clone)] struct FakeDailyReviewRunner { fetch_result: Result, DailyReviewServiceError>, - calls: Arc>>, + calls: Arc>>, } impl FakeDailyReviewRunner { @@ -175,7 +175,7 @@ impl FakeDailyReviewRunner { } } - fn calls(&self) -> Vec<(String, NaiveDate)> { + fn calls(&self) -> Vec { self.calls.lock().unwrap().clone() } } @@ -184,7 +184,6 @@ impl FakeDailyReviewRunner { impl DailyReviewRunner for FakeDailyReviewRunner { async fn review_day( &self, - _user_id: &str, _utc_date: NaiveDate, ) -> Result { Ok(DailyReviewResult::EmptyDay) @@ -192,13 +191,9 @@ impl DailyReviewRunner for FakeDailyReviewRunner { async fn fetch_review( &self, - user_id: &str, utc_date: NaiveDate, ) -> Result, DailyReviewServiceError> { - self.calls - .lock() - .unwrap() - .push((user_id.to_string(), utc_date)); + self.calls.lock().unwrap().push(utc_date); self.fetch_result.clone() } } @@ -286,9 +281,8 @@ struct FailingPendingEmbeddingCounter; #[async_trait::async_trait] impl PendingEmbeddingCounter for FailingPendingEmbeddingCounter { - async fn count_entries_missing_embedding_for_user( + async fn count_entries_missing_embedding( &self, - _user_id: &str, _embedding_model: &str, ) -> Result { Err(EmbeddingRepositoryError::Database( @@ -315,7 +309,6 @@ fn incoming_for_conversation( source: MessageSource::Telegram, source_conversation_id: source_conversation_id.to_string(), source_message_id: source_message_id.to_string(), - user_id: "7".to_string(), text: text.to_string(), received_at, } @@ -329,7 +322,6 @@ fn command(command: JournalCommand) -> JournalCommandRequest { JournalCommandRequest { source: MessageSource::Telegram, source_conversation_id: "42".to_string(), - user_id: "7".to_string(), received_at: at(12, 0), command, } @@ -338,7 +330,6 @@ fn command(command: JournalCommand) -> JournalCommandRequest { fn daily_review(review_text: &str) -> DailyReview { DailyReview { id: 1, - user_id: "7".to_string(), review_date: date(), review_text: Some(review_text.to_string()), model: "test-model".to_string(), @@ -548,7 +539,6 @@ async fn status_uses_single_user_journal_stats_and_command_received_at_date() { source: MessageSource::Telegram, source_conversation_id: "42".to_string(), source_message_id: "3".to_string(), - user_id: "8".to_string(), text: "other user".to_string(), received_at: Utc.with_ymd_and_hms(2026, 4, 29, 9, 0, 0).unwrap(), }) @@ -559,7 +549,6 @@ async fn status_uses_single_user_journal_stats_and_command_received_at_date() { .command(&JournalCommandRequest { source: MessageSource::Telegram, source_conversation_id: "42".to_string(), - user_id: "7".to_string(), received_at: Utc.with_ymd_and_hms(2026, 4, 29, 12, 0, 0).unwrap(), command: JournalCommand::Status, }) @@ -622,7 +611,6 @@ async fn status_reports_configured_embedding_status_and_single_user_pending_coun source: MessageSource::Telegram, source_conversation_id: "42".to_string(), source_message_id: "3".to_string(), - user_id: "8".to_string(), text: "other user pending entry".to_string(), received_at: at(12, 0), }) @@ -758,7 +746,7 @@ async fn day_review_last_returns_existing_review() { yesterday.format("%Y-%m-%d") ) ); - assert_eq!(runner.calls(), vec![("7".to_string(), yesterday)]); + assert_eq!(runner.calls(), vec![yesterday]); } #[tokio::test] @@ -992,7 +980,6 @@ async fn undo_deletes_latest_entry_for_current_conversation() { .command(&JournalCommandRequest { source: MessageSource::Telegram, source_conversation_id: "99".to_string(), - user_id: "7".to_string(), received_at: at(12, 0), command: JournalCommand::Last, }) @@ -1311,7 +1298,7 @@ async fn process_survives_panic_in_capture_time_extraction() { #[derive(Debug, Clone)] struct FakeWeeklyReviewRunner { fetch_result: Result, WeeklyReviewServiceError>, - calls: Arc>>, + calls: Arc>>, } impl FakeWeeklyReviewRunner { @@ -1324,7 +1311,7 @@ impl FakeWeeklyReviewRunner { } } - fn calls(&self) -> Vec<(String, NaiveDate)> { + fn calls(&self) -> Vec { self.calls.lock().unwrap().clone() } } @@ -1333,7 +1320,6 @@ impl FakeWeeklyReviewRunner { impl WeeklyReviewRunner for FakeWeeklyReviewRunner { async fn review_week( &self, - _user_id: &str, _week_start: NaiveDate, ) -> Result { Ok(WeeklyReviewResult::SparseWeek) @@ -1341,13 +1327,9 @@ impl WeeklyReviewRunner for FakeWeeklyReviewRunner { async fn fetch_review( &self, - user_id: &str, week_start: NaiveDate, ) -> Result, WeeklyReviewServiceError> { - self.calls - .lock() - .unwrap() - .push((user_id.to_string(), week_start)); + self.calls.lock().unwrap().push(week_start); self.fetch_result.clone() } } @@ -1355,7 +1337,6 @@ impl WeeklyReviewRunner for FakeWeeklyReviewRunner { fn weekly_review(text: &str, week_start: NaiveDate) -> WeeklyReview { WeeklyReview { id: 1, - user_id: "7".to_string(), week_start_date: week_start, review_text: Some(text.to_string()), model: "test-model".to_string(), @@ -1417,10 +1398,7 @@ async fn week_review_fetches_previous_iso_week_monday() { outgoing.text, "Weekly review for week of 2026-04-20\n\nstored weekly review" ); - assert_eq!( - runner.calls(), - vec![("7".to_string(), previous_week_monday())] - ); + assert_eq!(runner.calls(), vec![previous_week_monday()]); } #[tokio::test] diff --git a/src/journal/store.rs b/src/journal/store.rs index 3459b41c..b8a3e341 100644 --- a/src/journal/store.rs +++ b/src/journal/store.rs @@ -201,7 +201,6 @@ mod tests { source: MessageSource::Telegram, source_conversation_id: "42".to_string(), source_message_id: source_message_id.to_string(), - user_id: "7".to_string(), text: text.to_string(), received_at, } diff --git a/src/journal/transfer.rs b/src/journal/transfer.rs index 55077a33..d7ec82c2 100644 --- a/src/journal/transfer.rs +++ b/src/journal/transfer.rs @@ -228,7 +228,6 @@ mod tests { source: MessageSource::Telegram, source_conversation_id: "42".to_string(), source_message_id: message_id.to_string(), - user_id: "42".to_string(), text: text.to_string(), received_at: Utc::now(), } diff --git a/src/journal/week_review/generator.rs b/src/journal/week_review/generator.rs index e7b271a8..7b8d6982 100644 --- a/src/journal/week_review/generator.rs +++ b/src/journal/week_review/generator.rs @@ -487,7 +487,6 @@ mod tests { DailyReviewSignal { id: 0, daily_review_id: 0, - user_id: "user-1".to_string(), review_date: day(0), signal_type, label: label.to_string(), diff --git a/src/journal/week_review/mod.rs b/src/journal/week_review/mod.rs index c451816a..2c8961c8 100644 --- a/src/journal/week_review/mod.rs +++ b/src/journal/week_review/mod.rs @@ -16,7 +16,6 @@ use crate::journal::review::signals::types::DailyReviewSignal; #[derive(Debug, Clone, PartialEq, Eq)] pub struct WeeklyReview { pub id: i64, - pub user_id: String, pub week_start_date: NaiveDate, pub review_text: Option, pub model: String, diff --git a/src/journal/week_review/repository.rs b/src/journal/week_review/repository.rs index 152d03b2..8e789fbd 100644 --- a/src/journal/week_review/repository.rs +++ b/src/journal/week_review/repository.rs @@ -4,8 +4,6 @@ use thiserror::Error; use crate::errors::from_error_string; -use crate::messages::SINGLE_USER_ID; - use super::{WeeklyReview, WeeklyReviewStatus}; #[derive(Debug, Clone, PartialEq, Eq, Error)] @@ -210,7 +208,6 @@ fn row_to_weekly_review(row: SqliteRow) -> Result Result; async fn fetch_review( &self, - user_id: &str, week_start: NaiveDate, ) -> Result, WeeklyReviewServiceError>; } @@ -99,7 +96,6 @@ impl WeeklyReviewService { pub async fn review_week( &self, - user_id: &str, week_start: NaiveDate, ) -> Result { let existing = self @@ -163,13 +159,7 @@ impl WeeklyReviewService { let trimmed = review_text.trim(); if trimmed.is_empty() { return self - .store_failed( - user_id, - week_start, - model, - &prompt_version, - EMPTY_REVIEW_ERROR, - ) + .store_failed(week_start, model, &prompt_version, EMPTY_REVIEW_ERROR) .await; } @@ -194,7 +184,7 @@ impl WeeklyReviewService { Err(error) => { let prompt_version = self.generator.prompt_version(); let message = error.to_string(); - self.store_failed(user_id, week_start, model, &prompt_version, &message) + self.store_failed(week_start, model, &prompt_version, &message) .await } } @@ -202,7 +192,6 @@ impl WeeklyReviewService { pub async fn fetch_review( &self, - _user_id: &str, week_start: NaiveDate, ) -> Result, WeeklyReviewServiceError> { let review = self @@ -219,7 +208,6 @@ impl WeeklyReviewService { async fn store_failed( &self, - user_id: &str, week_start: NaiveDate, model: &str, prompt_version: &str, @@ -229,7 +217,6 @@ impl WeeklyReviewService { .upsert_failed(week_start, model, prompt_version, error_message) .await?; Ok(WeeklyReviewResult::GenerationFailed(WeeklyReviewFailure { - user_id: user_id.to_string(), week_start_date: week_start, model: model.to_string(), prompt_version: prompt_version.to_string(), @@ -242,18 +229,16 @@ impl WeeklyReviewService { impl WeeklyReviewRunner for WeeklyReviewService { async fn review_week( &self, - user_id: &str, week_start: NaiveDate, ) -> Result { - WeeklyReviewService::review_week(self, user_id, week_start).await + WeeklyReviewService::review_week(self, week_start).await } async fn fetch_review( &self, - user_id: &str, week_start: NaiveDate, ) -> Result, WeeklyReviewServiceError> { - WeeklyReviewService::fetch_review(self, user_id, week_start).await + WeeklyReviewService::fetch_review(self, week_start).await } } @@ -343,7 +328,6 @@ mod tests { async fn seed_completed_daily( daily_reviews: &DailyReviewRepository, - _user_id: &str, date: NaiveDate, text: &str, ) -> i64 { @@ -370,7 +354,7 @@ mod tests { .await .unwrap(); - let result = service.review_week("user-1", week_start()).await.unwrap(); + let result = service.review_week(week_start()).await.unwrap(); assert_eq!(result, WeeklyReviewResult::Existing(existing)); assert_eq!(generator.calls(), 0); @@ -380,10 +364,10 @@ mod tests { async fn sparse_week_skips_generation_when_below_min_daily_reviews() { let (service, _weekly, daily, _signals, generator) = setup(FakeWeeklyReviewGenerator::succeeding("ignored")).await; - seed_completed_daily(&daily, "user-1", day(0), "monday").await; - seed_completed_daily(&daily, "user-1", day(2), "wednesday").await; + seed_completed_daily(&daily, day(0), "monday").await; + seed_completed_daily(&daily, day(2), "wednesday").await; - let result = service.review_week("user-1", week_start()).await.unwrap(); + let result = service.review_week(week_start()).await.unwrap(); assert_eq!(result, WeeklyReviewResult::SparseWeek); assert_eq!(generator.calls(), 0); @@ -393,11 +377,11 @@ mod tests { async fn generates_review_when_threshold_met() { let (service, weekly_reviews, daily, _signals, generator) = setup(FakeWeeklyReviewGenerator::succeeding("week review")).await; - seed_completed_daily(&daily, "user-1", day(0), "monday").await; - seed_completed_daily(&daily, "user-1", day(1), "tuesday").await; - seed_completed_daily(&daily, "user-1", day(2), "wednesday").await; + seed_completed_daily(&daily, day(0), "monday").await; + seed_completed_daily(&daily, day(1), "tuesday").await; + seed_completed_daily(&daily, day(2), "wednesday").await; - let review = generated(service.review_week("user-1", week_start()).await.unwrap()); + let review = generated(service.review_week(week_start()).await.unwrap()); assert_eq!(review.review_text, Some("week review".to_string())); assert_eq!(review.status, WeeklyReviewStatus::Completed); @@ -422,9 +406,9 @@ mod tests { async fn passes_signals_grouped_by_day_to_generator() { let (service, _weekly, daily, signals_repo, generator) = setup(FakeWeeklyReviewGenerator::succeeding("text")).await; - let monday_review = seed_completed_daily(&daily, "user-1", day(0), "monday").await; - let tuesday_review = seed_completed_daily(&daily, "user-1", day(1), "tuesday").await; - seed_completed_daily(&daily, "user-1", day(2), "wednesday").await; + let monday_review = seed_completed_daily(&daily, day(0), "monday").await; + let tuesday_review = seed_completed_daily(&daily, day(1), "tuesday").await; + seed_completed_daily(&daily, day(2), "wednesday").await; signals_repo .replace_in_transaction(monday_review, day(0), &[theme_candidate()], "model", "v1") @@ -441,7 +425,7 @@ mod tests { .await .unwrap(); - service.review_week("user-1", week_start()).await.unwrap(); + service.review_week(week_start()).await.unwrap(); let inputs = generator.inputs_seen(); let monday_slice = inputs[0].days.iter().find(|d| d.date == day(0)).unwrap(); @@ -456,18 +440,18 @@ mod tests { async fn excludes_dailies_and_signals_outside_the_week() { let (service, _weekly, daily, signals_repo, generator) = setup(FakeWeeklyReviewGenerator::succeeding("text")).await; - let prior_review = seed_completed_daily(&daily, "user-1", day(-1), "previous sunday").await; - seed_completed_daily(&daily, "user-1", day(0), "monday").await; - seed_completed_daily(&daily, "user-1", day(1), "tuesday").await; - seed_completed_daily(&daily, "user-1", day(2), "wednesday").await; - seed_completed_daily(&daily, "user-1", day(7), "next monday").await; + let prior_review = seed_completed_daily(&daily, day(-1), "previous sunday").await; + seed_completed_daily(&daily, day(0), "monday").await; + seed_completed_daily(&daily, day(1), "tuesday").await; + seed_completed_daily(&daily, day(2), "wednesday").await; + seed_completed_daily(&daily, day(7), "next monday").await; signals_repo .replace_in_transaction(prior_review, day(-1), &[theme_candidate()], "m", "v1") .await .unwrap(); - service.review_week("user-1", week_start()).await.unwrap(); + service.review_week(week_start()).await.unwrap(); let inputs = generator.inputs_seen(); let dates: Vec<_> = inputs[0].days.iter().map(|d| d.date).collect(); @@ -482,15 +466,14 @@ mod tests { let (service, weekly_reviews, daily, _signals, generator) = setup(FakeWeeklyReviewGenerator::succeeding(" \n\t")).await; for offset in 0..3 { - seed_completed_daily(&daily, "user-1", day(offset), "text").await; + seed_completed_daily(&daily, day(offset), "text").await; } - let result = service.review_week("user-1", week_start()).await.unwrap(); + let result = service.review_week(week_start()).await.unwrap(); assert_eq!( result, WeeklyReviewResult::GenerationFailed(WeeklyReviewFailure { - user_id: "user-1".to_string(), week_start_date: week_start(), model: "fake-weekly-review-model".to_string(), prompt_version: "fake-weekly-prompt-v1".to_string(), @@ -514,15 +497,14 @@ mod tests { let (service, weekly_reviews, daily, _signals, _generator) = setup(FakeWeeklyReviewGenerator::failing("provider down")).await; for offset in 0..3 { - seed_completed_daily(&daily, "user-1", day(offset), "text").await; + seed_completed_daily(&daily, day(offset), "text").await; } - let result = service.review_week("user-1", week_start()).await.unwrap(); + let result = service.review_week(week_start()).await.unwrap(); assert_eq!( result, WeeklyReviewResult::GenerationFailed(WeeklyReviewFailure { - user_id: "user-1".to_string(), week_start_date: week_start(), model: "fake-weekly-review-model".to_string(), prompt_version: "fake-weekly-prompt-v1".to_string(), @@ -547,11 +529,11 @@ mod tests { ]); let (service, weekly_reviews, daily, _signals, generator) = setup(generator).await; for offset in 0..3 { - seed_completed_daily(&daily, "user-1", day(offset), "text").await; + seed_completed_daily(&daily, day(offset), "text").await; } assert!(matches!( - service.review_week("user-1", week_start()).await.unwrap(), + service.review_week(week_start()).await.unwrap(), WeeklyReviewResult::GenerationFailed(_) )); let failed = weekly_reviews @@ -561,7 +543,7 @@ mod tests { .unwrap(); std::thread::sleep(std::time::Duration::from_millis(2)); - let completed = generated(service.review_week("user-1", week_start()).await.unwrap()); + let completed = generated(service.review_week(week_start()).await.unwrap()); assert_eq!(generator.calls(), 2); assert_eq!(completed.id, failed.id); @@ -577,14 +559,14 @@ mod tests { let (service, weekly_reviews, daily, _signals, generator) = setup(FakeWeeklyReviewGenerator::succeeding("regenerated")).await; for offset in 0..3 { - seed_completed_daily(&daily, "user-1", day(offset), "text").await; + seed_completed_daily(&daily, day(offset), "text").await; } let existing = weekly_reviews .upsert_completed(week_start(), "", "old-model", "v0", "{}") .await .unwrap(); - let review = generated(service.review_week("user-1", week_start()).await.unwrap()); + let review = generated(service.review_week(week_start()).await.unwrap()); assert_eq!(generator.calls(), 1); assert_eq!(review.id, existing.id); @@ -593,7 +575,7 @@ mod tests { } #[tokio::test] - async fn same_week_reviews_share_single_user_scope() { + async fn repeated_review_week_returns_existing_without_regenerating() { let (service, _weekly, daily, _signals, generator) = setup(FakeWeeklyReviewGenerator::new(vec![ Ok("user one review".to_string()), @@ -601,19 +583,18 @@ mod tests { ])) .await; for offset in 0..3 { - seed_completed_daily(&daily, "user-1", day(offset), "user one").await; - seed_completed_daily(&daily, "user-2", day(offset), "user two").await; + seed_completed_daily(&daily, day(offset), "user one").await; + seed_completed_daily(&daily, day(offset), "user two").await; } - let one = generated(service.review_week("user-1", week_start()).await.unwrap()); - let two = match service.review_week("user-2", week_start()).await.unwrap() { + let one = generated(service.review_week(week_start()).await.unwrap()); + let two = match service.review_week(week_start()).await.unwrap() { WeeklyReviewResult::Existing(review) => review, other => panic!("expected existing review, got {other:?}"), }; assert_eq!(one.review_text, Some("user one review".to_string())); assert_eq!(two.review_text, Some("user one review".to_string())); - assert_eq!(one.user_id, two.user_id); assert_eq!(generator.calls(), 1); } @@ -626,7 +607,7 @@ mod tests { .await .unwrap(); - let result = service.fetch_review("user-1", week_start()).await.unwrap(); + let result = service.fetch_review(week_start()).await.unwrap(); assert_eq!(result.unwrap().review_text, Some("review text".to_string())); } @@ -636,7 +617,7 @@ mod tests { let (service, _weekly, _daily, _signals, _generator) = setup(FakeWeeklyReviewGenerator::succeeding("any")).await; - let result = service.fetch_review("user-1", week_start()).await.unwrap(); + let result = service.fetch_review(week_start()).await.unwrap(); assert!(result.is_none()); } @@ -650,7 +631,7 @@ mod tests { .await .unwrap(); - let result = service.fetch_review("user-1", week_start()).await.unwrap(); + let result = service.fetch_review(week_start()).await.unwrap(); assert!(result.is_none()); } @@ -659,15 +640,15 @@ mod tests { async fn generated_review_persists_serialized_inputs_snapshot() { let (service, weekly_reviews, daily, signals_repo, _generator) = setup(FakeWeeklyReviewGenerator::succeeding("week review")).await; - let monday_review = seed_completed_daily(&daily, "user-1", day(0), "monday text").await; - seed_completed_daily(&daily, "user-1", day(1), "tuesday text").await; - seed_completed_daily(&daily, "user-1", day(2), "wednesday text").await; + let monday_review = seed_completed_daily(&daily, day(0), "monday text").await; + seed_completed_daily(&daily, day(1), "tuesday text").await; + seed_completed_daily(&daily, day(2), "wednesday text").await; signals_repo .replace_in_transaction(monday_review, day(0), &[theme_candidate()], "model", "v1") .await .unwrap(); - service.review_week("user-1", week_start()).await.unwrap(); + service.review_week(week_start()).await.unwrap(); let stored = weekly_reviews .find_by_user_and_week(week_start()) diff --git a/src/messages.rs b/src/messages.rs index 1789eb85..e95a592b 100644 --- a/src/messages.rs +++ b/src/messages.rs @@ -1,7 +1,5 @@ use chrono::{DateTime, Utc}; -pub const SINGLE_USER_ID: &str = "default"; - #[derive(Debug, Clone, PartialEq, Eq)] pub enum MessageSource { Telegram, @@ -20,7 +18,6 @@ pub struct IncomingMessage { pub source: MessageSource, pub source_conversation_id: String, pub source_message_id: String, - pub user_id: String, pub text: String, pub received_at: DateTime, } diff --git a/src/workers/daily_review.rs b/src/workers/daily_review.rs index cd00d8b9..d849a2d4 100644 --- a/src/workers/daily_review.rs +++ b/src/workers/daily_review.rs @@ -107,11 +107,7 @@ where }; for target in targets { - let review = match self - .review_runner - .review_day(&target.user_id, review_date) - .await - { + let review = match self.review_runner.review_day(review_date).await { Ok(DailyReviewResult::Existing(review) | DailyReviewResult::Generated(review)) => { review } @@ -121,7 +117,6 @@ where } Ok(DailyReviewResult::GenerationFailed(failure)) => { warn!( - user_id = %failure.user_id, review_date = %failure.review_date, error = %failure.error_message, "daily review generation failed during delivery" @@ -130,8 +125,7 @@ where continue; } Err(error) => { - self.record_review_runner_error(&target.user_id, review_date, error) - .await?; + self.record_review_runner_error(review_date, error).await?; result.failed += 1; continue; } @@ -160,7 +154,6 @@ where .mark_delivery_failed(review_date, &error) .await?; warn!( - user_id = %target.user_id, source_conversation_id = %target.source_conversation_id, review_date = %review_date, error = %error, @@ -176,12 +169,10 @@ where async fn record_review_runner_error( &self, - user_id: &str, review_date: NaiveDate, error: DailyReviewServiceError, ) -> Result<(), DailyReviewDeliveryWorkerError> { warn!( - user_id = %user_id, review_date = %review_date, error = %error, "daily review runner failed during delivery" @@ -372,7 +363,6 @@ mod tests { impl DailyReviewRunner for FakeRunner { async fn review_day( &self, - _user_id: &str, _utc_date: NaiveDate, ) -> Result { self.result.clone() @@ -380,7 +370,6 @@ mod tests { async fn fetch_review( &self, - _user_id: &str, _utc_date: NaiveDate, ) -> Result, DailyReviewServiceError> { Ok(None) @@ -421,17 +410,11 @@ mod tests { NaiveDate::from_ymd_opt(2026, 4, 28).unwrap() } - fn entry_for( - user_id: &str, - conversation_id: &str, - message_id: &str, - text: &str, - ) -> IncomingMessage { + fn entry_for(conversation_id: &str, message_id: &str, text: &str) -> IncomingMessage { IncomingMessage { source: MessageSource::Telegram, source_conversation_id: conversation_id.to_string(), source_message_id: message_id.to_string(), - user_id: user_id.to_string(), text: text.to_string(), received_at: Utc.with_ymd_and_hms(2026, 4, 28, 12, 0, 0).unwrap(), } @@ -442,7 +425,6 @@ mod tests { source: MessageSource::Telegram, source_conversation_id: "42".to_string(), source_message_id: source_message_id.to_string(), - user_id: "7".to_string(), text: text.to_string(), received_at: Utc.with_ymd_and_hms(2026, 4, 28, 12, 0, 0).unwrap(), } @@ -695,11 +677,11 @@ mod tests { let (worker, daily_reviews, journal_entries, sender) = setup(FakeReviewGenerator::succeeding("review text"), sender).await; journal_entries - .store(&entry_for("7", "42", "1", "first")) + .store(&entry_for("42", "1", "first")) .await .unwrap(); journal_entries - .store(&entry_for("8", "99", "2", "second")) + .store(&entry_for("99", "2", "second")) .await .unwrap(); diff --git a/src/workers/embedding.rs b/src/workers/embedding.rs index d356be50..175d6cfe 100644 --- a/src/workers/embedding.rs +++ b/src/workers/embedding.rs @@ -105,7 +105,6 @@ mod tests { source: MessageSource::Telegram, source_conversation_id: "42".to_string(), source_message_id: source_message_id.to_string(), - user_id: "7".to_string(), text: text.to_string(), received_at, } diff --git a/src/workers/extraction.rs b/src/workers/extraction.rs index 6fd2b4e7..1853db0c 100644 --- a/src/workers/extraction.rs +++ b/src/workers/extraction.rs @@ -106,7 +106,6 @@ mod tests { source: MessageSource::Telegram, source_conversation_id: "42".to_string(), source_message_id: source_message_id.to_string(), - user_id: "7".to_string(), text: text.to_string(), received_at, }; diff --git a/src/workers/signals.rs b/src/workers/signals.rs index ef6a080c..81a6e861 100644 --- a/src/workers/signals.rs +++ b/src/workers/signals.rs @@ -145,7 +145,6 @@ mod tests { source: MessageSource::Telegram, source_conversation_id: "42".to_string(), source_message_id: msg_id.to_string(), - user_id: "user-1".to_string(), text: text.to_string(), received_at: chrono::Utc::now(), }; diff --git a/src/workers/weekly_review.rs b/src/workers/weekly_review.rs index 0dea6738..3d5cbce0 100644 --- a/src/workers/weekly_review.rs +++ b/src/workers/weekly_review.rs @@ -119,11 +119,7 @@ where }; for target in targets { - let review = match self - .review_runner - .review_week(&target.user_id, week_start) - .await - { + let review = match self.review_runner.review_week(week_start).await { Ok( WeeklyReviewResult::Existing(review) | WeeklyReviewResult::Generated(review), ) => review, @@ -133,7 +129,6 @@ where } Ok(WeeklyReviewResult::GenerationFailed(failure)) => { warn!( - user_id = %failure.user_id, week_start = %failure.week_start_date, error = %failure.error_message, "weekly review generation failed during delivery" @@ -142,8 +137,7 @@ where continue; } Err(error) => { - self.record_review_runner_error(&target.user_id, week_start, error) - .await?; + self.record_review_runner_error(week_start, error).await?; result.failed += 1; continue; } @@ -172,7 +166,6 @@ where .mark_delivery_failed(week_start, &error) .await?; warn!( - user_id = %target.user_id, source_conversation_id = %target.source_conversation_id, week_start = %week_start, error = %error, @@ -188,12 +181,10 @@ where async fn record_review_runner_error( &self, - _user_id: &str, week_start: NaiveDate, error: WeeklyReviewServiceError, ) -> Result<(), WeeklyReviewDeliveryWorkerError> { warn!( - user_id = %_user_id, week_start = %week_start, error = %error, "weekly review runner failed during delivery" @@ -410,23 +401,17 @@ mod tests { Utc.with_ymd_and_hms(2026, 4, 27, 12, 0, 0).unwrap() + Duration::days(offset) } - fn entry( - user_id: &str, - conversation_id: &str, - message_id: &str, - day_offset: i64, - ) -> IncomingMessage { + fn entry(conversation_id: &str, message_id: &str, day_offset: i64) -> IncomingMessage { IncomingMessage { source: MessageSource::Telegram, source_conversation_id: conversation_id.to_string(), source_message_id: message_id.to_string(), - user_id: user_id.to_string(), text: format!("entry on day {day_offset}"), received_at: day_within_week(day_offset), } } - async fn seed_three_daily_reviews(daily_reviews: &DailyReviewRepository, _user_id: &str) { + async fn seed_three_daily_reviews(daily_reviews: &DailyReviewRepository) { for offset in 0..3 { let date = week_start() + Duration::days(offset); daily_reviews @@ -467,8 +452,8 @@ mod tests { FakeSender::succeeding(), ) .await; - seed_three_daily_reviews(&daily, "7").await; - journal.store(&entry("7", "42", "1", 0)).await.unwrap(); + seed_three_daily_reviews(&daily).await; + journal.store(&entry("42", "1", 0)).await.unwrap(); // Tuesday — not the configured kickoff day (Monday). let tuesday = Utc.with_ymd_and_hms(2026, 5, 5, 9, 0, 0).unwrap(); @@ -495,8 +480,8 @@ mod tests { sender, ) .await; - seed_three_daily_reviews(&daily, "7").await; - journal.store(&entry("7", "42", "1", 0)).await.unwrap(); + seed_three_daily_reviews(&daily).await; + journal.store(&entry("42", "1", 0)).await.unwrap(); // Monday after the target week. let monday = Utc.with_ymd_and_hms(2026, 5, 4, 6, 0, 0).unwrap(); @@ -540,7 +525,7 @@ mod tests { .await .unwrap(); } - journal.store(&entry("7", "42", "1", 0)).await.unwrap(); + journal.store(&entry("42", "1", 0)).await.unwrap(); let result = worker.run_once_for_week(week_start()).await.unwrap(); @@ -561,8 +546,8 @@ mod tests { let sender = FakeSender::failing("telegram unavailable"); let (worker, weekly_reviews, daily, journal, sender) = setup(FakeWeeklyReviewGenerator::succeeding("week review"), sender).await; - seed_three_daily_reviews(&daily, "7").await; - journal.store(&entry("7", "42", "1", 0)).await.unwrap(); + seed_three_daily_reviews(&daily).await; + journal.store(&entry("42", "1", 0)).await.unwrap(); let result = worker.run_once_for_week(week_start()).await.unwrap(); @@ -594,8 +579,8 @@ mod tests { let sender = FakeSender::skipped(); let (worker, weekly_reviews, daily, journal, sender) = setup(FakeWeeklyReviewGenerator::succeeding("week review"), sender).await; - seed_three_daily_reviews(&daily, "7").await; - journal.store(&entry("7", "42", "1", 0)).await.unwrap(); + seed_three_daily_reviews(&daily).await; + journal.store(&entry("42", "1", 0)).await.unwrap(); let result = worker.run_once_for_week(week_start()).await.unwrap(); @@ -624,8 +609,8 @@ mod tests { let sender = FakeSender::succeeding(); let (worker, weekly_reviews, daily, journal, sender) = setup(FakeWeeklyReviewGenerator::succeeding("ignored"), sender).await; - seed_three_daily_reviews(&daily, "7").await; - journal.store(&entry("7", "42", "1", 0)).await.unwrap(); + seed_three_daily_reviews(&daily).await; + journal.store(&entry("42", "1", 0)).await.unwrap(); weekly_reviews .upsert_completed(week_start(), "existing", "m", "v1", "{}") .await @@ -675,8 +660,8 @@ mod tests { sender, ) .await; - seed_three_daily_reviews(&daily, "7").await; - journal.store(&entry("7", "42", "1", 0)).await.unwrap(); + seed_three_daily_reviews(&daily).await; + journal.store(&entry("42", "1", 0)).await.unwrap(); let result = worker.run_once_for_week(week_start()).await.unwrap(); diff --git a/tests/mcp_server.rs b/tests/mcp_server.rs index 9d9d35c7..7ae46db3 100644 --- a/tests/mcp_server.rs +++ b/tests/mcp_server.rs @@ -23,7 +23,7 @@ use froid::{ adapters::mcp::AnalyzerMcpServer, database, journal::analyzer::{ - SemanticJournalSearcher, UserContext, build_analyzer_mcp_components, + SemanticJournalSearcher, build_analyzer_mcp_components, types::{AnalyzerError, SemanticHit}, }, }; @@ -34,7 +34,6 @@ struct StubSemanticSearcher; impl SemanticJournalSearcher for StubSemanticSearcher { async fn search( &self, - _user_id: &str, _query: &str, _from_date: Option, _to_date_exclusive: Option, @@ -55,7 +54,7 @@ async fn fresh_pool() -> SqlitePool { async fn lists_and_calls_analyzer_tools_over_streamable_http() { let pool = fresh_pool().await; let components = build_analyzer_mcp_components(pool, Arc::new(StubSemanticSearcher)); - let server = AnalyzerMcpServer::new(components, UserContext::new("test-user")); + let server = AnalyzerMcpServer::new(components); let cancel = CancellationToken::new(); let service = StreamableHttpService::new( diff --git a/tests/multiuser_tests.rs b/tests/multiuser_tests.rs index 6ad4944d..96a279d8 100644 --- a/tests/multiuser_tests.rs +++ b/tests/multiuser_tests.rs @@ -9,7 +9,7 @@ use froid::{ review::signals::wiring::DailyReviewSignalRuntimeConfig, week_review::WeeklyReviewRuntimeConfig, }, - messages::{IncomingMessage, MessageSource, SINGLE_USER_ID}, + messages::{IncomingMessage, MessageSource}, }; use tokio_util::sync::CancellationToken; @@ -60,7 +60,6 @@ async fn test_multiuser_database_isolation_and_routing() { source: MessageSource::Telegram, source_conversation_id: "user_a".to_string(), source_message_id: "msg_1".to_string(), - user_id: SINGLE_USER_ID.to_string(), text: "Today was a productive day writing Rust integration tests.".to_string(), received_at: chrono::Utc::now(), }; @@ -86,7 +85,6 @@ async fn test_multiuser_database_isolation_and_routing() { source: MessageSource::Telegram, source_conversation_id: "user_b".to_string(), source_message_id: "msg_2".to_string(), - user_id: SINGLE_USER_ID.to_string(), text: "I spent the afternoon gardening in the backyard.".to_string(), received_at: chrono::Utc::now(), }; @@ -112,7 +110,6 @@ async fn test_multiuser_database_isolation_and_routing() { let cmd_a = JournalCommandRequest { source: MessageSource::Telegram, source_conversation_id: "user_a".to_string(), - user_id: SINGLE_USER_ID.to_string(), received_at: chrono::Utc::now(), command: JournalCommand::Recent { requested_limit: 10, @@ -135,7 +132,6 @@ async fn test_multiuser_database_isolation_and_routing() { let cmd_b = JournalCommandRequest { source: MessageSource::Telegram, source_conversation_id: "user_b".to_string(), - user_id: SINGLE_USER_ID.to_string(), received_at: chrono::Utc::now(), command: JournalCommand::Recent { requested_limit: 10,