1use crate::storage::{
5 self, OcrAnnotationEntry, SideChatMetadata, ThreadMetadata, ThreadStorage, WorkspaceMetadata,
6};
7use chrono::Utc;
8use regex::{Regex, RegexBuilder};
9use std::cmp::Ordering;
10use std::collections::HashSet;
11
12use crate::services::brain;
13use crate::{settings, thread};
14
15pub type ExplorerResult<T> = std::result::Result<T, String>;
16
17#[derive(Clone)]
18pub struct ExplorerThread {
19 pub id: String,
20 pub title: String,
21 pub created_at: String,
22 pub updated_at: String,
23 pub image_hash: String,
24 pub pinned_at: Option<String>,
25 pub workspace_id: Option<String>,
26}
27
28impl From<ThreadMetadata> for ExplorerThread {
29 fn from(thread: ThreadMetadata) -> Self {
30 Self {
31 id: thread.id,
32 title: thread.title,
33 created_at: thread.created_at.to_rfc3339(),
34 updated_at: thread.updated_at.to_rfc3339(),
35 image_hash: thread.image_hash,
36 pinned_at: thread.pinned_at.map(|value| value.to_rfc3339()),
37 workspace_id: None,
38 }
39 }
40}
41
42pub struct ExplorerWorkspace {
43 pub id: String,
44 pub name: String,
45 pub created_at: String,
46 pub directories: Vec<String>,
47 pub total_thread_count: u32,
48 pub threads: Vec<ExplorerThread>,
49}
50
51impl From<WorkspaceMetadata> for ExplorerWorkspace {
52 fn from(workspace: WorkspaceMetadata) -> Self {
53 Self {
54 id: workspace.id.clone(),
55 name: workspace.name,
56 created_at: workspace.created_at.to_rfc3339(),
57 directories: workspace.directories,
58 total_thread_count: workspace.threads.len().min(u32::MAX as usize) as u32,
59 threads: workspace
60 .threads
61 .into_values()
62 .map(|thread| ExplorerThread {
63 workspace_id: Some(workspace.id.clone()),
64 ..thread.into()
65 })
66 .collect(),
67 }
68 }
69}
70
71pub struct UnassignedThreadsPage {
72 pub threads: Vec<ExplorerThread>,
73 pub total: u32,
74}
75
76pub struct ExplorerSideChatThread {
77 pub id: String,
78 pub title: String,
79 pub created_at: String,
80 pub updated_at: String,
81}
82
83impl From<SideChatMetadata> for ExplorerSideChatThread {
84 fn from(sidechat: SideChatMetadata) -> Self {
85 Self {
86 id: sidechat.id,
87 title: sidechat.title,
88 created_at: sidechat.created_at.to_rfc3339(),
89 updated_at: sidechat.updated_at.to_rfc3339(),
90 }
91 }
92}
93
94pub struct ExplorerJobSnapshot {
95 pub job_id: String,
96 pub thread_id: String,
97 pub kind: String,
98 pub status: String,
99}
100
101pub struct ThreadSearchResult {
102 pub thread_id: String,
103 pub thread_title: String,
104 pub thread_created_at: String,
105 pub thread_updated_at: String,
106 pub workspace_title: Option<String>,
107 pub result_kind: String,
108 pub result_index: u32,
109 pub snippet: String,
110 pub score: u32,
111}
112
113fn active_storage() -> ExplorerResult<ThreadStorage> {
114 storage::thread_store().map_err(|error| error.to_string())
115}
116
117#[derive(Clone, Copy)]
118enum ThreadOrdering {
119 Created,
120 Updated,
121}
122
123impl ThreadOrdering {
124 fn parse(value: &str) -> ExplorerResult<Self> {
125 match value {
126 "created" => Ok(Self::Created),
127 "updated" => Ok(Self::Updated),
128 _ => Err(format!("Unsupported thread ordering: {value}")),
129 }
130 }
131}
132
133#[derive(Clone, Copy)]
134enum WorkspaceOrdering {
135 Created,
136 Updated,
137}
138
139impl WorkspaceOrdering {
140 fn parse(value: &str) -> ExplorerResult<Self> {
141 match value {
142 "created" => Ok(Self::Created),
143 "updated" => Ok(Self::Updated),
144 _ => Err(format!("Unsupported workspace ordering: {value}")),
145 }
146 }
147}
148
149fn compare_threads(
150 left: &ThreadMetadata,
151 right: &ThreadMetadata,
152 ordering: ThreadOrdering,
153) -> Ordering {
154 match (&left.pinned_at, &right.pinned_at) {
155 (Some(left_pinned), Some(right_pinned)) => right_pinned.cmp(left_pinned),
156 (Some(_), None) => Ordering::Less,
157 (None, Some(_)) => Ordering::Greater,
158 (None, None) => match ordering {
159 ThreadOrdering::Created => right.created_at.cmp(&left.created_at),
160 ThreadOrdering::Updated => right.updated_at.cmp(&left.updated_at),
161 },
162 }
163}
164
165fn sort_threads(mut threads: Vec<ThreadMetadata>, ordering: ThreadOrdering) -> Vec<ThreadMetadata> {
166 threads.sort_by(|left, right| compare_threads(left, right, ordering));
167 threads
168}
169
170fn latest_workspace_activity(workspace: &WorkspaceMetadata) -> Option<chrono::DateTime<Utc>> {
171 workspace
172 .threads
173 .values()
174 .map(|thread| &thread.updated_at)
175 .max()
176 .cloned()
177}
178
179pub fn list_workspaces(
180 workspace_ordering: String,
181 thread_ordering: String,
182) -> ExplorerResult<Vec<ExplorerWorkspace>> {
183 let workspace_ordering = WorkspaceOrdering::parse(&workspace_ordering)?;
184 let thread_ordering = ThreadOrdering::parse(&thread_ordering)?;
185 let mut workspaces = active_storage()?
186 .list_workspaces()
187 .map_err(|error| error.to_string())?;
188 workspaces.sort_by(|left, right| match workspace_ordering {
189 WorkspaceOrdering::Created => right.created_at.cmp(&left.created_at),
190 WorkspaceOrdering::Updated => latest_workspace_activity(right)
191 .unwrap_or(right.created_at)
192 .cmp(&latest_workspace_activity(left).unwrap_or(left.created_at)),
193 });
194 Ok(workspaces
195 .into_iter()
196 .map(|workspace| {
197 let threads = sort_threads(
198 workspace.threads.values().cloned().collect(),
199 thread_ordering,
200 )
201 .into_iter()
202 .map(|thread| ExplorerThread {
203 workspace_id: Some(workspace.id.clone()),
204 ..thread.into()
205 })
206 .collect();
207 ExplorerWorkspace {
208 threads,
209 ..workspace.into()
210 }
211 })
212 .collect())
213}
214
215pub fn list_unassigned_threads(
216 offset: u32,
217 limit: u32,
218 thread_ordering: String,
219) -> ExplorerResult<UnassignedThreadsPage> {
220 let ordering = ThreadOrdering::parse(&thread_ordering)?;
221 let threads = active_storage()?
222 .list_unassigned_threads()
223 .map_err(|error| error.to_string())?;
224 let total = threads.len().min(u32::MAX as usize) as u32;
225 let threads = sort_threads(threads, ordering)
226 .into_iter()
227 .skip(offset as usize)
228 .take(limit as usize)
229 .map(Into::into)
230 .collect();
231 Ok(UnassignedThreadsPage { threads, total })
232}
233
234pub fn list_sidechat_threads() -> ExplorerResult<Vec<ExplorerSideChatThread>> {
235 let mut sidechats = active_storage()?
236 .list_sidechat_threads()
237 .map_err(|error| error.to_string())?;
238 sidechats.sort_by_key(|sidechat| std::cmp::Reverse(sidechat.created_at));
239 Ok(sidechats.into_iter().map(Into::into).collect())
240}
241
242pub fn rename_sidechat(
243 sidechat_id: String,
244 title: String,
245) -> ExplorerResult<ExplorerSideChatThread> {
246 let storage = active_storage()?;
247 let mut sidechat = storage
248 .load_sidechat(&sidechat_id)
249 .map_err(|error| error.to_string())?;
250 sidechat.metadata.title = title;
251 sidechat.metadata.updated_at = Utc::now();
252 storage
253 .update_sidechat_metadata(&sidechat.metadata)
254 .map_err(|error| error.to_string())?;
255 Ok(sidechat.metadata.into())
256}
257
258pub fn delete_sidechat(sidechat_id: String) -> ExplorerResult<()> {
259 active_storage()?
260 .delete_sidechat(&sidechat_id)
261 .map_err(|error| error.to_string())
262}
263
264pub fn fork_sidechat(sidechat_id: String) -> ExplorerResult<ExplorerSideChatThread> {
265 active_storage()?
266 .fork_sidechat_latest(&sidechat_id)
267 .map(Into::into)
268 .map_err(|error| error.to_string())
269}
270
271pub fn group_threads(first_id: String, second_id: String) -> ExplorerResult<ExplorerWorkspace> {
272 active_storage()?
273 .group_threads(&first_id, &second_id)
274 .map(Into::into)
275 .map_err(|error| error.to_string())
276}
277
278pub fn create_workspace(
279 name: String,
280 directories: Vec<String>,
281) -> ExplorerResult<ExplorerWorkspace> {
282 active_storage()?
283 .create_workspace(&name, &directories)
284 .map(Into::into)
285 .map_err(|error| error.to_string())
286}
287
288pub fn update_workspace(
289 workspace_id: String,
290 name: String,
291 directories: Vec<String>,
292) -> ExplorerResult<ExplorerWorkspace> {
293 active_storage()?
294 .update_workspace(&workspace_id, &name, &directories)
295 .map(Into::into)
296 .map_err(|error| error.to_string())
297}
298
299pub fn delete_workspace(workspace_id: String) -> ExplorerResult<()> {
300 active_storage()?
301 .delete_workspace(&workspace_id)
302 .map_err(|error| error.to_string())
303}
304
305pub fn set_thread_workspace(thread_id: String, workspace_id: Option<String>) -> ExplorerResult<()> {
306 active_storage()?
307 .set_thread_workspace(&thread_id, workspace_id.as_deref())
308 .map_err(|error| error.to_string())
309}
310
311pub fn rename_thread(thread_id: String, title: String) -> ExplorerResult<ExplorerThread> {
312 let storage = active_storage()?;
313 let mut metadata = storage
314 .load_thread(&thread_id)
315 .map_err(|error| error.to_string())?
316 .metadata;
317 metadata.title = title;
318 storage
319 .update_thread_metadata(&metadata)
320 .map_err(|error| error.to_string())?;
321 let workspace_id = storage
322 .get_thread_workspace_id(&metadata.id)
323 .map_err(|error| error.to_string())?;
324 Ok(ExplorerThread {
325 workspace_id,
326 ..metadata.into()
327 })
328}
329
330pub fn toggle_pin_thread(thread_id: String) -> ExplorerResult<ExplorerThread> {
331 let storage = active_storage()?;
332 let mut metadata = storage
333 .load_thread(&thread_id)
334 .map_err(|error| error.to_string())?
335 .metadata;
336 metadata.pinned_at = metadata.pinned_at.is_none().then(Utc::now);
337 storage
338 .update_thread_metadata(&metadata)
339 .map_err(|error| error.to_string())?;
340 let workspace_id = storage
341 .get_thread_workspace_id(&metadata.id)
342 .map_err(|error| error.to_string())?;
343 Ok(ExplorerThread {
344 workspace_id,
345 ..metadata.into()
346 })
347}
348
349pub fn delete_thread(thread_id: String) -> ExplorerResult<()> {
350 active_storage()?
351 .delete_thread(&thread_id)
352 .map_err(|error| error.to_string())
353}
354
355pub fn delete_threads(thread_ids: Vec<String>) -> ExplorerResult<()> {
356 active_storage()?
357 .delete_threads(&thread_ids)
358 .map_err(|error| error.to_string())
359}
360
361pub fn fork_thread(thread_id: String) -> ExplorerResult<Option<ExplorerThread>> {
362 let storage = active_storage()?;
363 let metadata = storage
364 .fork_thread_latest(&thread_id)
365 .map_err(|error| error.to_string())?;
366 let workspace_id = storage
367 .get_thread_workspace_id(&metadata.id)
368 .map_err(|error| error.to_string())?;
369 Ok(Some(ExplorerThread {
370 workspace_id,
371 ..metadata.into()
372 }))
373}
374
375enum SearchPlan {
376 Regex(Option<Regex>),
377 Tokens(Vec<String>),
378 Literal(String),
379}
380
381fn search_plan(query: &str) -> SearchPlan {
382 let regex_source = query.strip_prefix("re:").map(str::to_string).or_else(|| {
383 (query.starts_with('/') && query.ends_with('/') && query.len() > 2)
384 .then(|| query[1..query.len() - 1].to_string())
385 });
386 if let Some(source) = regex_source {
387 return SearchPlan::Regex(
388 RegexBuilder::new(&source)
389 .case_insensitive(true)
390 .build()
391 .ok(),
392 );
393 }
394
395 let mut seen = HashSet::new();
396 let mut tokens = query
397 .split_whitespace()
398 .map(|part| part.trim_start_matches(['-', '+']))
399 .map(|part| {
400 part.to_lowercase()
401 .chars()
402 .filter(|character| character.is_ascii_alphanumeric())
403 .collect::<String>()
404 })
405 .filter(|part| part.len() >= 2)
406 .filter(|part| seen.insert(part.clone()))
407 .collect::<Vec<_>>();
408 tokens.sort_by_key(|right| std::cmp::Reverse(right.len()));
409 tokens.truncate(8);
410 if tokens.is_empty() {
411 SearchPlan::Literal(query.to_lowercase())
412 } else {
413 SearchPlan::Tokens(tokens)
414 }
415}
416
417fn matching_index(content: &str, plan: &SearchPlan) -> Option<(usize, u32)> {
418 match plan {
419 SearchPlan::Regex(regex) => regex
420 .as_ref()
421 .and_then(|regex| regex.find(content))
422 .map(|matched| (matched.start(), 100)),
423 SearchPlan::Tokens(tokens) => {
424 let lower = content.to_lowercase();
425 let mut first = None;
426 for token in tokens {
427 let index = lower.find(token)?;
428 first.get_or_insert(index);
429 }
430 first.map(|index| (index, tokens.len() as u32 * 10))
431 }
432 SearchPlan::Literal(query) => content.to_lowercase().find(query).map(|index| (index, 50)),
433 }
434}
435
436fn char_boundary_at_or_before(value: &str, mut index: usize) -> usize {
437 index = index.min(value.len());
438 while !value.is_char_boundary(index) {
439 index -= 1;
440 }
441 index
442}
443
444fn snippet(content: &str, match_index: usize, query_length: usize) -> String {
445 let start = char_boundary_at_or_before(content, match_index.saturating_sub(40));
446 let end = char_boundary_at_or_before(
447 content,
448 match_index
449 .saturating_add(query_length)
450 .saturating_add(40)
451 .min(content.len()),
452 );
453 let mut value = content[start..end].replace('\n', " ");
454 if start > 0 {
455 value.insert_str(0, "...");
456 }
457 if end < content.len() {
458 value.push_str("...");
459 }
460 value
461}
462
463struct ThreadSearchContext<'a> {
464 workspace: Option<&'a WorkspaceMetadata>,
465 thread: &'a ThreadMetadata,
466 plan: &'a SearchPlan,
467 query_length: usize,
468}
469
470fn push_search_result(
471 results: &mut Vec<ThreadSearchResult>,
472 context: &ThreadSearchContext<'_>,
473 result_kind: &str,
474 result_index: usize,
475 content: &str,
476 score_bonus: u32,
477) {
478 let Some((match_index, score)) = matching_index(content, context.plan) else {
479 return;
480 };
481 results.push(ThreadSearchResult {
482 thread_id: context.thread.id.clone(),
483 thread_title: context.thread.title.clone(),
484 thread_created_at: context.thread.created_at.to_rfc3339(),
485 thread_updated_at: context.thread.updated_at.to_rfc3339(),
486 workspace_title: context.workspace.map(|workspace| workspace.name.clone()),
487 result_kind: result_kind.to_string(),
488 result_index: result_index.min(u32::MAX as usize) as u32,
489 snippet: snippet(content, match_index, context.query_length),
490 score: score.saturating_add(score_bonus),
491 });
492}
493
494pub fn search_threads(query: String, limit: u32) -> ExplorerResult<Vec<ThreadSearchResult>> {
495 let query = query.trim().to_string();
496 if query.is_empty() {
497 return Ok(Vec::new());
498 }
499 let storage = active_storage()?;
500 let plan = search_plan(&query);
501 let mut results = Vec::new();
502 let workspaces = storage
503 .list_workspaces()
504 .map_err(|error| error.to_string())?;
505 let unassigned = storage
506 .list_unassigned_threads()
507 .map_err(|error| error.to_string())?;
508 let entries = workspaces
509 .iter()
510 .flat_map(|workspace| {
511 workspace
512 .threads
513 .values()
514 .map(move |thread| (Some(workspace), thread))
515 })
516 .chain(unassigned.iter().map(|thread| (None, thread)));
517 for (workspace, thread) in entries {
518 let context = ThreadSearchContext {
519 workspace,
520 thread,
521 plan: &plan,
522 query_length: query.len(),
523 };
524 push_search_result(&mut results, &context, "title", 0, &thread.title, 60);
525 let messages = match storage.load_messages(&thread.id) {
526 Ok(messages) => messages,
527 Err(_) => continue,
528 };
529 for (message_index, message) in messages.iter().enumerate() {
530 push_search_result(
531 &mut results,
532 &context,
533 "message",
534 message_index,
535 message.content(),
536 20,
537 );
538 }
539 if let Ok(annotations) = storage.get_ocr_annotations(&thread.id) {
540 let mut ocr_index = 0;
541 for annotation in annotations.into_values() {
542 let regions = match annotation {
543 OcrAnnotationEntry::EmptyState(regions) => regions,
544 OcrAnnotationEntry::Model(model) => model.ocr_data,
545 };
546 for region in regions {
547 push_search_result(&mut results, &context, "ocr", ocr_index, ®ion.text, 10);
548 ocr_index += 1;
549 }
550 }
551 }
552 }
553 results.sort_by(|left, right| {
554 right.score.cmp(&left.score).then_with(|| {
555 let left_time = left.thread_updated_at.parse::<chrono::DateTime<Utc>>().ok();
556 let right_time = right
557 .thread_updated_at
558 .parse::<chrono::DateTime<Utc>>()
559 .ok();
560 right_time.cmp(&left_time)
561 })
562 });
563 results.truncate(limit as usize);
564 Ok(results)
565}
566
567pub async fn suggest_thread_title(thread_id: String) -> ExplorerResult<String> {
568 let config = tokio::task::spawn_blocking(settings::load_config)
569 .await
570 .map_err(|error| format!("Settings load task failed: {error}"))??;
571 let candidates = brain()
572 .build_model_attempt_plan(config.model, config.effort)
573 .await?;
574 let title = brain()
575 .suggest_thread_title(thread_id.clone(), candidates)
576 .await?;
577 let persisted_title = title.clone();
578 tokio::task::spawn_blocking(move || {
579 let storage = active_storage()?;
580 let mut metadata = storage
581 .load_thread(&thread_id)
582 .map_err(|error| error.to_string())?
583 .metadata;
584 metadata.title = persisted_title;
585 storage
586 .update_thread_metadata(&metadata)
587 .map_err(|error| error.to_string())
588 })
589 .await
590 .map_err(|error| error.to_string())??;
591 Ok(title)
592}
593
594pub fn get_jobs_snapshot() -> ExplorerResult<Vec<ExplorerJobSnapshot>> {
595 let brain_jobs = thread::get_thread_jobs_snapshot()?;
596 let ocr_jobs = thread::ocr::get_ocr_jobs_snapshot()?;
597 Ok(brain_jobs
598 .into_iter()
599 .map(|job| ExplorerJobSnapshot {
600 job_id: job.job_id,
601 thread_id: job.thread_id,
602 kind: "brain".to_string(),
603 status: job.status,
604 })
605 .chain(ocr_jobs.into_iter().map(|job| ExplorerJobSnapshot {
606 job_id: job.job_id,
607 thread_id: job.thread_id,
608 kind: "ocr".to_string(),
609 status: job.status,
610 }))
611 .collect())
612}