use std::collections::HashMap;
use std::sync::Arc;
use uuid::Uuid;
use crate::app::events::AppEvent;
use crate::entities::chat::Chat;
use crate::entities::message::{Message, MessageRole};
use crate::features::chat_search::{self, SearchGroup, SearchHit};
use crate::features::chat_search_sort::SortMode;
use crate::shared::storage::Storage;
use crate::shared::storage::cache::{IndexScope, IndexedMessage, MessageHit};
use super::{ChatView, Orchestrator};
impl Orchestrator {
pub(super) fn handle_search_chats(&self, query: String) {
let chat_ids = match crate::features::chat_search::to_fts_query(&query) {
None => None,
Some(fts) => match self.storage.cache().search_chats(&fts) {
Ok(ids) => Some(ids),
Err(err) => {
tracing::warn!(query_chars = query.chars().count(), error = %format!("{err:#}"),
"chat content search failed");
None
}
},
};
let _ = self
.evt_tx
.send(AppEvent::ChatSearchResults { query, chat_ids });
}
pub(super) fn handle_search_messages(&self, query: String, sort: SortMode) {
let (groups, total) = self.message_search(&query, sort);
let _ = self.evt_tx.send(AppEvent::MessageSearchResults {
query,
groups,
total,
});
}
fn message_search(&self, query: &str, sort: SortMode) -> (Vec<SearchGroup>, usize) {
let Some(fts) = chat_search::to_fts_query(query) else {
return (Vec::new(), 0);
};
let cache = self.storage.cache();
let hits = match cache.search_messages(&fts, chat_search::HIT_CAP) {
Ok(hits) => hits,
Err(err) => {
tracing::warn!(query_chars = query.chars().count(), error = %format!("{err:#}"),
"message content search failed");
return (Vec::new(), 0);
}
};
let total = if hits.len() < chat_search::HIT_CAP {
hits.len()
} else {
cache.count_matching_messages(&fts).unwrap_or(hits.len())
};
(self.group_hits(hits, query, sort), total)
}
fn group_hits(&self, hits: Vec<MessageHit>, query: &str, sort: SortMode) -> Vec<SearchGroup> {
let mut by_scope: HashMap<(Uuid, Option<Uuid>), Vec<MessageHit>> = HashMap::new();
for hit in hits {
by_scope
.entry((hit.chat_id, hit.sub_id))
.or_default()
.push(hit);
}
let mut chats: Vec<&Chat> = self
.chats
.iter()
.filter(|c| !c.is_hidden)
.filter(|c| {
by_scope.contains_key(&(c.id, None))
|| c.children()
.any(|r| by_scope.contains_key(&(c.id, Some(r.id))))
})
.collect();
chats.sort_by_key(|c| {
std::cmp::Reverse(match sort {
SortMode::Created => c.created_at,
SortMode::Modified => c.modified_at,
})
});
let mut groups = Vec::new();
for chat in chats {
let own = by_scope.remove(&(chat.id, None)).unwrap_or_default();
groups.push(make_group(
chat.id,
None,
&chat.title,
&chat.messages,
own,
query,
));
for run in chat.children() {
if let Some(hits) = by_scope.remove(&(chat.id, Some(run.id))) {
groups.push(make_group(
run.id,
Some(chat.id),
&run.title,
&run.messages,
hits,
query,
));
}
}
}
groups
}
pub(super) fn first_match_in_chat(&self, id: Uuid, query: &str) -> Option<Uuid> {
let fts = chat_search::to_fts_query(query)?;
let (scope, messages) = match self.view(id)? {
ChatView::Top(chat) => (IndexScope::Chat(id), &chat.messages),
ChatView::Child { run, .. } => (IndexScope::Transcript(id), &run.messages),
};
let ids = self
.storage
.cache()
.matching_messages_in_chat(&fts, scope)
.inspect_err(|err| {
tracing::warn!(chat = %id, error = %format!("{err:#}"),
"resolving the first match in a chat failed");
})
.ok()?;
let order = message_order(messages);
ids.into_iter()
.filter_map(|id| order.get(&id).map(|pos| (*pos, id)))
.min()
.map(|(_, id)| id)
}
pub(super) fn index_saved_chat(&self, chat: &Chat) {
if chat.is_hidden {
self.forget_chat_index(chat.id);
return;
}
let Some(info) = self.storage.json().chat_file_info(chat.id) else {
tracing::warn!(chat = %chat.id, "could not stat a just-saved chat file for the search index");
return;
};
let messages = indexed_messages(chat);
if let Err(err) =
self.storage
.cache()
.index_chat(chat.id, info.mtime_ms, info.size, &messages)
{
tracing::warn!(chat = %chat.id, error = %format!("{err:#}"),
"failed to update the chat search index");
}
}
pub(super) fn forget_chat_index(&self, chat_id: Uuid) {
if let Err(err) = self.storage.cache().forget_chat(chat_id) {
tracing::warn!(chat = %chat_id, error = %format!("{err:#}"),
"failed to drop a chat from the search index");
}
}
}
fn make_group(
id: Uuid,
parent: Option<Uuid>,
title: &str,
messages: &[Message],
mut hits: Vec<MessageHit>,
query: &str,
) -> SearchGroup {
let order = message_order(messages);
hits.sort_by_key(|h| order.get(&h.message_id).copied().unwrap_or(usize::MAX));
SearchGroup {
chat_id: id,
parent,
title: title.to_string(),
hits: hits
.into_iter()
.map(|h| SearchHit {
message_id: h.message_id,
role: h.role,
ts: h.ts,
snippet: chat_search::build_snippet(
&h.text,
query,
chat_search::SNIPPET_BUDGET_CHARS,
),
})
.collect(),
}
}
fn indexed_messages(chat: &Chat) -> Vec<IndexedMessage> {
let own = chat.messages.iter().map(|m| (None, m));
let runs = chat
.children()
.flat_map(|run| run.messages.iter().map(move |m| (Some(run.id), m)));
own.chain(runs)
.filter(|(_, m)| !m.text.trim().is_empty())
.map(|(sub_id, m)| IndexedMessage {
id: m.id,
sub_id,
role: role_str(m.role).to_string(),
ts: m.timestamp.to_rfc3339(),
text: m.text.clone(),
})
.collect()
}
fn message_order(messages: &[Message]) -> HashMap<Uuid, usize> {
messages
.iter()
.enumerate()
.map(|(i, m)| (m.id, i))
.collect()
}
fn role_str(role: MessageRole) -> &'static str {
match role {
MessageRole::System => "system",
MessageRole::User => "user",
MessageRole::Assistant => "assistant",
MessageRole::Tool => "tool",
}
}
pub(super) fn spawn_reconcile(storage: Arc<Storage>) {
tokio::task::spawn_blocking(move || reconcile(&storage));
}
fn reconcile(storage: &Storage) -> (usize, usize, usize) {
let indexed = match storage.cache().indexed_state() {
Ok(state) => state,
Err(err) => {
tracing::warn!(error = %format!("{err:#}"), "search index: cannot read its bookkeeping");
return (0, 0, 0);
}
};
let files = match storage.json().chat_files() {
Ok(files) => files,
Err(err) => {
tracing::warn!(error = %format!("{err:#}"), "search index: cannot list chat files");
return (0, 0, 0);
}
};
let (mut reindexed, mut forgotten, mut unchanged) = (0, 0, 0);
for file in &files {
let was = indexed.get(&file.id).copied();
if was == Some((file.mtime_ms, file.size)) {
unchanged += 1;
continue;
}
reconcile_file(
storage,
file,
was,
&mut reindexed,
&mut forgotten,
&mut unchanged,
);
}
for id in indexed.keys() {
if !files.iter().any(|f| f.id == *id) && forget(storage, *id) {
forgotten += 1;
}
}
tracing::info!(
indexed = reindexed,
forgotten,
unchanged,
"chat search index reconciled"
);
(reindexed, forgotten, unchanged)
}
fn reconcile_file(
storage: &Storage,
file: &crate::shared::storage::json::ChatFileInfo,
was: Option<(i64, u64)>,
reindexed: &mut usize,
forgotten: &mut usize,
unchanged: &mut usize,
) {
match storage.json().load_chat(file.id) {
Ok(Some(chat)) if chat.is_hidden => match write_indexed(storage, file, was, &[]) {
Ok(true) => *forgotten += 1,
Ok(false) => *unchanged += 1,
Err(err) => tracing::warn!(chat = %file.id, error = %format!("{err:#}"),
"search index: failed to drop a hidden chat"),
},
Ok(Some(chat)) => match write_indexed(storage, file, was, &indexed_messages(&chat)) {
Ok(true) => *reindexed += 1,
Ok(false) => {
*unchanged += 1;
tracing::debug!(chat = %file.id,
"search index: a fresher write won, leaving this chat alone");
}
Err(err) => tracing::warn!(chat = %file.id, error = %format!("{err:#}"),
"search index: failed to index a chat"),
},
Ok(None) => {}
Err(err) => tracing::warn!(chat = %file.id, error = %format!("{err:#}"),
"search index: skipped an unreadable chat file"),
}
}
fn write_indexed(
storage: &Storage,
file: &crate::shared::storage::json::ChatFileInfo,
was: Option<(i64, u64)>,
messages: &[IndexedMessage],
) -> anyhow::Result<bool> {
storage
.cache()
.index_chat_if_unchanged(file.id, was, file.mtime_ms, file.size, messages)
}
fn forget(storage: &Storage, id: Uuid) -> bool {
match storage.cache().forget_chat(id) {
Ok(()) => true,
Err(err) => {
tracing::warn!(chat = %id, error = %format!("{err:#}"),
"search index: failed to forget a chat");
false
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::entities::profile::Profile;
use crate::shared::paths::Paths;
fn storage() -> (tempfile::TempDir, Storage) {
let dir = tempfile::tempdir().unwrap();
let storage = Storage::open(Paths::with_root(dir.path())).unwrap();
(dir, storage)
}
fn chat_with(texts: &[&str]) -> Chat {
let profile = Profile::new("P", "sys");
let mut chat = Chat::from_profile(&profile, "заголовок");
for t in texts {
chat.push_message(Message::user(*t));
}
chat
}
#[test]
fn reconcile_indexes_new_files_and_skips_unchanged_ones() {
let (_d, storage) = storage();
let chat = chat_with(&["содержимое про кошек"]);
storage.json().save_chat(&chat).unwrap();
let (indexed, forgotten, unchanged) = reconcile(&storage);
assert_eq!((indexed, forgotten, unchanged), (1, 0, 0));
assert_eq!(
storage.cache().search_chats("\"кош\"").unwrap(),
vec![chat.id]
);
let (indexed, forgotten, unchanged) = reconcile(&storage);
assert_eq!(
(indexed, forgotten, unchanged),
(0, 0, 1),
"an unchanged file must not be re-parsed"
);
}
#[test]
fn reconcile_reindexes_a_changed_file_and_forgets_a_deleted_one() {
let (_d, storage) = storage();
let mut chat = chat_with(&["первая версия текста"]);
storage.json().save_chat(&chat).unwrap();
reconcile(&storage);
chat.messages.clear();
chat.push_message(Message::user("вторая версия текста"));
storage.json().save_chat(&chat).unwrap();
let (indexed, ..) = reconcile(&storage);
assert_eq!(indexed, 1);
assert_eq!(
storage.cache().search_chats("\"вторая\"").unwrap(),
vec![chat.id]
);
assert!(
storage
.cache()
.search_chats("\"первая\"")
.unwrap()
.is_empty()
);
std::fs::remove_file(_d.path().join("chats").join(format!("{}.json", chat.id))).unwrap();
let (_, forgotten, _) = reconcile(&storage);
assert_eq!(forgotten, 1);
assert!(
storage
.cache()
.search_chats("\"вторая\"")
.unwrap()
.is_empty()
);
assert!(storage.cache().indexed_state().unwrap().is_empty());
}
#[test]
fn reconcile_drops_a_hidden_chat() {
let (_d, storage) = storage();
let chat = chat_with(&["секретное содержимое"]);
storage.json().save_chat(&chat).unwrap();
reconcile(&storage);
assert!(
!storage
.cache()
.search_chats("\"секрет\"")
.unwrap()
.is_empty()
);
storage.json().hide_chat(chat.id).unwrap();
let (_, forgotten, _) = reconcile(&storage);
assert_eq!(forgotten, 1);
assert!(
storage
.cache()
.search_chats("\"секрет\"")
.unwrap()
.is_empty(),
"a hidden chat must not show up in results"
);
}
#[test]
fn a_hidden_chat_is_not_re_examined_on_every_pass() {
let (_d, storage) = storage();
let chat = chat_with(&["секретное содержимое"]);
storage.json().save_chat(&chat).unwrap();
reconcile(&storage);
storage.json().hide_chat(chat.id).unwrap();
assert_eq!(reconcile(&storage).1, 1, "the pass that drops it");
let (indexed, forgotten, unchanged) = reconcile(&storage);
assert_eq!(
(indexed, forgotten, unchanged),
(0, 0, 1),
"a hidden chat must count as unchanged, not be dropped again"
);
assert!(
storage
.cache()
.indexed_state()
.unwrap()
.contains_key(&chat.id),
"and it must stay in the bookkeeping — that is what stops the re-parse"
);
}
#[test]
fn reconcile_survives_a_corrupt_chat_file() {
let (_d, storage) = storage();
let good = chat_with(&["исправное содержимое"]);
storage.json().save_chat(&good).unwrap();
let broken = Uuid::new_v4();
std::fs::write(
_d.path().join("chats").join(format!("{broken}.json")),
"{ не json",
)
.unwrap();
let (indexed, ..) = reconcile(&storage);
assert_eq!(indexed, 1, "the healthy chat is still indexed");
assert_eq!(
storage.cache().search_chats("\"исправ\"").unwrap(),
vec![good.id]
);
}
#[test]
fn the_pass_does_not_clobber_a_write_that_landed_while_it_read() {
let (_d, storage) = storage();
let mut chat = chat_with(&[]);
storage.json().save_chat(&chat).unwrap();
let seen_by_the_pass = storage
.cache()
.indexed_state()
.unwrap()
.get(&chat.id)
.copied();
let file_as_read = storage.json().chat_file_info(chat.id).unwrap();
let messages_as_read = indexed_messages(&chat);
assert!(messages_as_read.is_empty());
chat.push_message(Message::user("живая переписка"));
storage.json().save_chat(&chat).unwrap();
let now = storage.json().chat_file_info(chat.id).unwrap();
storage
.cache()
.index_chat(chat.id, now.mtime_ms, now.size, &indexed_messages(&chat))
.unwrap();
let wrote =
write_indexed(&storage, &file_as_read, seen_by_the_pass, &messages_as_read).unwrap();
assert!(!wrote, "the stale snapshot must be refused");
assert_eq!(
storage.cache().search_chats("\"живая\"").unwrap(),
vec![chat.id],
"the live index survived the reconciliation"
);
}
#[test]
fn the_pass_writes_when_nothing_moved_underneath_it() {
let (_d, storage) = storage();
let chat = chat_with(&["содержимое для индексации"]);
storage.json().save_chat(&chat).unwrap();
let seen = storage
.cache()
.indexed_state()
.unwrap()
.get(&chat.id)
.copied();
let file = storage.json().chat_file_info(chat.id).unwrap();
let wrote = write_indexed(&storage, &file, seen, &indexed_messages(&chat)).unwrap();
assert!(wrote);
assert_eq!(
storage.cache().search_chats("\"индексац\"").unwrap(),
vec![chat.id]
);
}
#[test]
fn only_message_text_is_indexed() {
let profile = Profile::new("P", "sys");
let mut chat = Chat::from_profile(&profile, "t");
let mut msg = Message::assistant("видимый ответ");
msg.thoughts = Some("скрытые рассуждения".into());
chat.push_message(msg);
chat.push_message(Message::assistant(" "));
let indexed = indexed_messages(&chat);
assert_eq!(indexed.len(), 1, "a blank message carries nothing to index");
assert_eq!(indexed[0].text, "видимый ответ");
assert_eq!(indexed[0].role, "assistant");
assert!(!indexed[0].ts.is_empty());
}
}