Skip to content

Commit 497500b

Browse files
Merge remote-tracking branch 'origin/fix/user-message-search-storage-scope' into fix/user-message-search-storage-scope
2 parents 8e123fa + 6ca1816 commit 497500b

4 files changed

Lines changed: 95 additions & 40 deletions

File tree

src/mcp/tools/handlers/mod.rs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -59,10 +59,10 @@ pub async fn handle_user_lcm_tool(
5959
.to_string(),
6060
});
6161
}
62-
let sessions_db_path = crate::sessions::user_sessions_db_path(profile_root);
6362
if tool_name == "tracedecay_message_search" {
64-
return session::handle_user_message_search(&sessions_db_path, args).await;
63+
return session::handle_user_message_search(profile_root, args).await;
6564
}
65+
let sessions_db_path = crate::sessions::user_sessions_db_path(profile_root);
6666
let context = session::LcmHandlerContext::user(&sessions_db_path);
6767
match tool_name {
6868
"tracedecay_lcm_status" => session::handle_lcm_status(context, args).await,

src/mcp/tools/handlers/session.rs

Lines changed: 28 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -2527,11 +2527,36 @@ pub(super) async fn handle_message_search(
25272527
}
25282528

25292529
pub(super) async fn handle_user_message_search(
2530-
sessions_db_path: &Path,
2530+
profile_root: &Path,
25312531
args: Value,
25322532
) -> Result<ToolResult> {
25332533
let request = parse_message_search_request(&args)?;
2534-
let Some(db) = GlobalDb::open_read_only_at(sessions_db_path).await else {
2534+
let sessions_db_path = crate::sessions::user_sessions_db_path(profile_root);
2535+
let catch_up_leader = if request.catch_up {
2536+
let key = format!(
2537+
"user:{}:{}",
2538+
sessions_db_path.display(),
2539+
request.provider_scope.response_label()
2540+
);
2541+
match claim_message_catch_up(key) {
2542+
MessageCatchUpClaim::Wait(done) => {
2543+
wait_for_message_catch_up(done).await;
2544+
None
2545+
}
2546+
MessageCatchUpClaim::Leader(leader) => Some(leader),
2547+
}
2548+
} else {
2549+
None
2550+
};
2551+
if catch_up_leader.is_some() {
2552+
let _ = crate::sessions::ingest_user_global_sources_for_provider_at(
2553+
profile_root,
2554+
request.provider_scope.provider(),
2555+
)
2556+
.await;
2557+
}
2558+
drop(catch_up_leader);
2559+
let Some(db) = GlobalDb::open_read_only_at(&sessions_db_path).await else {
25352560
return Ok(tool_json(
25362561
None,
25372562
&args,
@@ -2549,7 +2574,7 @@ pub(super) async fn handle_user_message_search(
25492574
} else {
25502575
search_session_messages_in_db(&db, &request).await
25512576
};
2552-
let payload = message_search_payload(&request, &results, false);
2577+
let payload = message_search_payload(&request, &results, request.catch_up);
25532578
Ok(tool_json_with_md(None, &args, &payload, || {
25542579
render_message_search_md(&payload)
25552580
}))

src/sessions/mod.rs

Lines changed: 63 additions & 34 deletions
Original file line numberDiff line numberDiff line change
@@ -54,6 +54,15 @@ pub async fn registered_project_roots() -> Vec<PathBuf> {
5454
/// is projectless.
5555
pub async fn try_registered_project_roots() -> Option<Vec<PathBuf>> {
5656
let global = GlobalDb::open().await?;
57+
registered_project_roots_from(&global).await
58+
}
59+
60+
async fn try_registered_project_roots_at(profile_root: &Path) -> Option<Vec<PathBuf>> {
61+
let global = GlobalDb::open_at(&profile_root.join("global.db")).await?;
62+
registered_project_roots_from(&global).await
63+
}
64+
65+
async fn registered_project_roots_from(global: &GlobalDb) -> Option<Vec<PathBuf>> {
5766
let mut roots = global
5867
.list_project_paths()
5968
.await
@@ -72,27 +81,42 @@ pub async fn ingest_user_codex_sessions(session_id: Option<String>) -> Transcrip
7281
let Ok(profile_root) = crate::storage::default_profile_root() else {
7382
return TranscriptIngestStats::default();
7483
};
75-
let Some(db) = open_user_session_db(&profile_root).await else {
84+
let Some(registered_roots) = try_registered_project_roots().await else {
7685
return TranscriptIngestStats::default();
7786
};
78-
let Some(source) = codex::CodexSource::new() else {
87+
ingest_user_codex_sessions_at(&profile_root, session_id, registered_roots).await
88+
}
89+
90+
async fn ingest_user_codex_sessions_at(
91+
profile_root: &Path,
92+
session_id: Option<String>,
93+
registered_roots: Vec<PathBuf>,
94+
) -> TranscriptIngestStats {
95+
let Some(db) = open_user_session_db(profile_root).await else {
7996
return TranscriptIngestStats::default();
8097
};
81-
let Some(registered_roots) = try_registered_project_roots().await else {
98+
let Some(source) = codex::CodexSource::new() else {
8299
return TranscriptIngestStats::default();
83100
};
84101
let source = source.for_user_scope(session_id, registered_roots);
85-
ingest_source(&db, &source, &profile_root, None).await
102+
ingest_source(&db, &source, profile_root, None).await
86103
}
87104

88105
pub async fn ingest_user_cursor_sessions() -> TranscriptIngestStats {
89106
let Ok(profile_root) = crate::storage::default_profile_root() else {
90107
return TranscriptIngestStats::default();
91108
};
92-
let Some(db) = open_user_session_db(&profile_root).await else {
109+
let Some(registered_roots) = try_registered_project_roots().await else {
93110
return TranscriptIngestStats::default();
94111
};
95-
let Some(registered_roots) = try_registered_project_roots().await else {
112+
ingest_user_cursor_sessions_at(&profile_root, registered_roots).await
113+
}
114+
115+
async fn ingest_user_cursor_sessions_at(
116+
profile_root: &Path,
117+
registered_roots: Vec<PathBuf>,
118+
) -> TranscriptIngestStats {
119+
let Some(db) = open_user_session_db(profile_root).await else {
96120
return TranscriptIngestStats::default();
97121
};
98122
let (composer_stats, owned) = if let Some(source) = cursor_composer::CursorComposerSource::new()
@@ -123,7 +147,7 @@ pub async fn ingest_user_cursor_sessions() -> TranscriptIngestStats {
123147
let source = source
124148
.with_skip_session_ids(owned)
125149
.for_user_scope(&registered_roots);
126-
composer_stats.merge(ingest_source(&db, &source, &profile_root, None).await)
150+
composer_stats.merge(ingest_source(&db, &source, profile_root, None).await)
127151
}
128152

129153
pub async fn ingest_user_global_sources() -> TranscriptIngestStats {
@@ -138,29 +162,47 @@ fn provider_selected(scope: Option<SessionProvider>, candidate: SessionProvider)
138162
/// outside an explicitly requested message-search scope.
139163
pub async fn ingest_user_global_sources_for_provider(
140164
provider: Option<SessionProvider>,
165+
) -> TranscriptIngestStats {
166+
let Ok(profile_root) = crate::storage::default_profile_root() else {
167+
return TranscriptIngestStats::default();
168+
};
169+
let Some(roots) = try_registered_project_roots().await else {
170+
return TranscriptIngestStats::default();
171+
};
172+
ingest_user_global_sources_for_provider_with_roots(&profile_root, provider, roots).await
173+
}
174+
175+
pub(crate) async fn ingest_user_global_sources_for_provider_at(
176+
profile_root: &Path,
177+
provider: Option<SessionProvider>,
178+
) -> TranscriptIngestStats {
179+
let Some(roots) = try_registered_project_roots_at(profile_root).await else {
180+
return TranscriptIngestStats::default();
181+
};
182+
ingest_user_global_sources_for_provider_with_roots(profile_root, provider, roots).await
183+
}
184+
185+
async fn ingest_user_global_sources_for_provider_with_roots(
186+
profile_root: &Path,
187+
provider: Option<SessionProvider>,
188+
roots: Vec<PathBuf>,
141189
) -> TranscriptIngestStats {
142190
let mut stats = TranscriptIngestStats::default();
143191
if provider_selected(provider, SessionProvider::Codex) {
144-
stats = stats.merge(ingest_user_codex_sessions(None).await);
192+
stats = stats.merge(ingest_user_codex_sessions_at(profile_root, None, roots.clone()).await);
145193
}
146194
if provider_selected(provider, SessionProvider::Cursor) {
147-
stats = stats.merge(ingest_user_cursor_sessions().await);
195+
stats = stats.merge(ingest_user_cursor_sessions_at(profile_root, roots.clone()).await);
148196
}
149-
let Ok(profile_root) = crate::storage::default_profile_root() else {
150-
return stats;
151-
};
152-
let Some(db) = open_user_session_db(&profile_root).await else {
153-
return stats;
154-
};
155-
let Some(roots) = try_registered_project_roots().await else {
197+
let Some(db) = open_user_session_db(profile_root).await else {
156198
return stats;
157199
};
158200
if provider_selected(provider, SessionProvider::Hermes) {
159201
stats = stats.merge(hermes::ingest_user_sessions(&db, &roots).await);
160202
}
161203
if provider_selected(provider, SessionProvider::Claude) {
162-
stats = stats
163-
.merge(claude::ingest_user_sessions(&db, &profile_root, None, roots.clone()).await);
204+
stats =
205+
stats.merge(claude::ingest_user_sessions(&db, profile_root, None, roots.clone()).await);
164206
}
165207
let mut sources: Vec<Box<dyn TranscriptSource>> = Vec::new();
166208
if provider_selected(provider, SessionProvider::Vibe)
@@ -189,7 +231,7 @@ pub async fn ingest_user_global_sources_for_provider(
189231
sources.push(Box::new(source.for_user_scope(roots)));
190232
}
191233
for source in sources {
192-
let source_stats = ingest_source(&db, source.as_ref(), &profile_root, None).await;
234+
let source_stats = ingest_source(&db, source.as_ref(), profile_root, None).await;
193235
stats = stats.merge(source_stats);
194236
}
195237
if stats.messages_upserted > 0 {
@@ -307,20 +349,7 @@ pub async fn ingest_global_sources_for_provider(
307349
project_root: &Path,
308350
provider: Option<SessionProvider>,
309351
) -> TranscriptIngestStats {
310-
ingest_global_sources_for_provider_inner(db, project_root, provider, true).await
311-
}
312-
313-
async fn ingest_global_sources_for_provider_inner(
314-
db: &GlobalDb,
315-
project_root: &Path,
316-
provider: Option<SessionProvider>,
317-
include_user_scope: bool,
318-
) -> TranscriptIngestStats {
319-
// Keep the profile-level projectless history current alongside every
320-
// project sweep while respecting an explicit provider scope.
321-
if include_user_scope {
322-
let _ = ingest_user_global_sources_for_provider(provider).await;
323-
}
352+
let _ = ingest_user_global_sources_for_provider(provider).await;
324353
ingest_project_sources_for_provider(db, project_root, provider, true).await
325354
}
326355

@@ -415,7 +444,7 @@ pub(crate) async fn ingest_global_sources_for_startup(
415444
project_root: &Path,
416445
) -> TranscriptIngestStats {
417446
let user = ingest_user_global_sources_for_startup().await;
418-
user.merge(ingest_global_sources_for_provider_inner(db, project_root, None, false).await)
447+
user.merge(ingest_project_sources_for_provider(db, project_root, None, true).await)
419448
}
420449

421450
/// Runs the bounded commit-attribution sweep against the correlation store.

tests/mcp_suite/mcp_handler_test.rs

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -12470,7 +12470,7 @@ async fn user_message_search_cli_bridge_accepts_storage_scope() {
1247012470
"tracedecay_message_search",
1247112471
"--json",
1247212472
"--args",
12473-
r#"{"storage_scope":"user","provider":"codex","query":"apricot","catch_up":false,"format":"json"}"#,
12473+
r#"{"storage_scope":"user","provider":"codex","query":"apricot","format":"json"}"#,
1247412474
])
1247512475
.output()
1247612476
.unwrap();
@@ -12483,6 +12483,7 @@ async fn user_message_search_cli_bridge_accepts_storage_scope() {
1248312483
let envelope: Value = serde_json::from_slice(&search_output.stdout).unwrap();
1248412484
let payload = extract_first_json_content(&envelope);
1248512485
assert_eq!(payload["status"], "ok");
12486+
assert_eq!(payload["catch_up_performed"], true);
1248612487
assert_eq!(payload["count"], 1);
1248712488
assert_eq!(payload["results"][0]["session"]["project_key"], "user");
1248812489
}

0 commit comments

Comments
 (0)