From e1f0456525317291f9dc4215be46529b00afc2e8 Mon Sep 17 00:00:00 2001 From: Lewis Date: Sun, 10 May 2026 19:03:20 +0300 Subject: [PATCH] feat(search): search author/repo/createdAt with nice filters Lewis: May this revision serve well! --- crates/search/src/lib.rs | 140 +++++++++++++++++++++++++++++-------- crates/types/src/search.rs | 57 +++++++++++++-- 2 files changed, 161 insertions(+), 36 deletions(-) diff --git a/crates/search/src/lib.rs b/crates/search/src/lib.rs index 7738262..7f10ea3 100644 --- a/crates/search/src/lib.rs +++ b/crates/search/src/lib.rs @@ -1,4 +1,5 @@ use std::future::Future; +use std::ops::Bound; use std::pin::Pin; use std::sync::Arc; use std::time::Duration; @@ -7,11 +8,14 @@ use bobbin_runtime::Clock; use bobbin_types::search::{SearchDoc, SearchSink}; use jacquard_common::DefaultStr; use jacquard_common::types::nsid::Nsid; -use jacquard_common::types::string::AtUri; +use jacquard_common::types::string::{AtUri, Did}; use tantivy::collector::TopDocs; -use tantivy::query::{BooleanQuery, Occur, Query, QueryParser, QueryParserError, TermQuery}; +use tantivy::query::{ + BooleanQuery, Occur, Query, QueryParser, QueryParserError, RangeQuery, TermQuery, +}; use tantivy::schema::{ - Field, IndexRecordOption, STORED, STRING, Schema, TEXT, TextFieldIndexing, TextOptions, Value, + FAST, Field, INDEXED, IndexRecordOption, STORED, STRING, Schema, TEXT, TextFieldIndexing, + TextOptions, Value, }; use tantivy::{Index, IndexReader, IndexWriter, ReloadPolicy, TantivyDocument, TantivyError, Term}; use thiserror::Error; @@ -113,6 +117,28 @@ struct Fields { nsid: Field, title: Field, body: Field, + author: Field, + created_at: Field, + repo: Field, +} + +#[derive(Clone, Debug, Default)] +pub struct SearchFilters { + pub nsid: Option>, + pub author: Option>, + pub repo: Option>, + pub since: Option, + pub until: Option, +} + +impl SearchFilters { + pub fn is_empty(&self) -> bool { + self.nsid.is_none() + && self.author.is_none() + && self.repo.is_none() + && self.since.is_none() + && self.until.is_none() + } } enum WriteOp { @@ -145,6 +171,9 @@ impl SearchIndex { .set_stored(); let title = sb.add_text_field("title", title_opts); let body = sb.add_text_field("body", TEXT); + let author = sb.add_text_field("author", STRING | STORED); + let created_at = sb.add_i64_field("created_at", INDEXED | FAST | STORED); + let repo = sb.add_text_field("repo", STRING | STORED); let schema = sb.build(); let index = Index::create_in_ram(schema); let writer: IndexWriter = index.writer(heap_bytes)?; @@ -159,6 +188,9 @@ impl SearchIndex { nsid, title, body, + author, + created_at, + repo, }; let (tx, rx) = mpsc::channel(WRITE_QUEUE_CAPACITY); let writer_reader = reader.clone(); @@ -187,7 +219,7 @@ impl SearchIndex { pub async fn search( &self, q: &str, - nsid_filter: Option<&Nsid>, + filters: SearchFilters, cursor: SearchCursor, limit: u32, ) -> Result { @@ -200,11 +232,10 @@ impl SearchIndex { } let inner = self.inner.clone(); let q_owned = trimmed.to_owned(); - let nsid_owned = nsid_filter.cloned(); let limit_usize = limit as usize; let offset = cursor.offset() as usize; match tokio::task::spawn_blocking(move || { - inner.search_blocking(&q_owned, nsid_owned.as_ref(), offset, limit_usize) + inner.search_blocking(&q_owned, &filters, offset, limit_usize) }) .await { @@ -222,7 +253,7 @@ pub trait SearchReader: Send + Sync + 'static { fn search<'a>( &'a self, query: &'a str, - nsid_filter: Option<&'a Nsid>, + filters: SearchFilters, cursor: SearchCursor, limit: u32, ) -> SearchReadFuture<'a>; @@ -232,11 +263,11 @@ impl SearchReader for SearchIndex { fn search<'a>( &'a self, query: &'a str, - nsid_filter: Option<&'a Nsid>, + filters: SearchFilters, cursor: SearchCursor, limit: u32, ) -> SearchReadFuture<'a> { - Box::pin(SearchIndex::search(self, query, nsid_filter, cursor, limit)) + Box::pin(SearchIndex::search(self, query, filters, cursor, limit)) } } @@ -244,23 +275,47 @@ impl Inner { fn search_blocking( &self, q: &str, - nsid_filter: Option<&Nsid>, + filters: &SearchFilters, offset: usize, limit: usize, ) -> Result { let parsed = self.parser.parse_query(q)?; - let final_query: Box = match nsid_filter { - None => parsed, - Some(n) => { + let final_query: Box = if filters.is_empty() { + parsed + } else { + let mut clauses: Vec<(Occur, Box)> = Vec::with_capacity(5); + clauses.push((Occur::Must, parsed)); + if let Some(n) = &filters.nsid { let term = Term::from_field_text(self.fields.nsid, n.as_ref()); - Box::new(BooleanQuery::new(vec![ - (Occur::Must, parsed), - ( - Occur::Must, - Box::new(TermQuery::new(term, IndexRecordOption::Basic)), - ), - ])) + clauses.push(( + Occur::Must, + Box::new(TermQuery::new(term, IndexRecordOption::Basic)), + )); + } + if let Some(a) = &filters.author { + let term = Term::from_field_text(self.fields.author, a.as_ref()); + clauses.push(( + Occur::Must, + Box::new(TermQuery::new(term, IndexRecordOption::Basic)), + )); + } + if let Some(r) = &filters.repo { + let term = Term::from_field_text(self.fields.repo, r.as_ref()); + clauses.push(( + Occur::Must, + Box::new(TermQuery::new(term, IndexRecordOption::Basic)), + )); + } + if filters.since.is_some() || filters.until.is_some() { + let lower = filters.since.map_or(Bound::Unbounded, |s| { + Bound::Included(Term::from_field_i64(self.fields.created_at, s)) + }); + let upper = filters.until.map_or(Bound::Unbounded, |u| { + Bound::Excluded(Term::from_field_i64(self.fields.created_at, u)) + }); + clauses.push((Occur::Must, Box::new(RangeQuery::new(lower, upper)))); } + Box::new(BooleanQuery::new(clauses)) }; let collector = TopDocs::with_limit(limit + 1) .and_offset(offset) @@ -367,6 +422,15 @@ fn apply_upsert(writer: &mut IndexWriter, fields: Fields, doc: SearchDoc) { td.add_text(fields.nsid, doc.nsid.as_ref()); td.add_text(fields.title, &doc.title); td.add_text(fields.body, &doc.body); + if let Some(author) = &doc.author { + td.add_text(fields.author, author.as_ref()); + } + if let Some(ts) = doc.created_at { + td.add_i64(fields.created_at, ts); + } + if let Some(repo) = &doc.repo { + td.add_text(fields.repo, repo.as_ref()); + } if let Err(e) = writer.add_document(td) { warn!(?e, "search add_document failed"); } @@ -439,6 +503,9 @@ mod tests { nsid: nsid(nsid_s), title: title.to_owned(), body: body.to_owned(), + author: None, + created_at: None, + repo: None, } } @@ -466,7 +533,7 @@ mod tests { idx.flush().await; let page = idx - .search("barnacle", None, SearchCursor::Start, 10) + .search("barnacle", SearchFilters::default(), SearchCursor::Start, 10) .await .unwrap(); assert_eq!(page.hits.len(), 1); @@ -499,7 +566,10 @@ mod tests { let only_strings = idx .search( "anemone", - Some(&nsid("sh.tangled.string")), + SearchFilters { + nsid: Some(nsid("sh.tangled.string")), + ..SearchFilters::default() + }, SearchCursor::Start, 10, ) @@ -520,13 +590,13 @@ mod tests { idx.flush().await; let abalone_hits = idx - .search("abalone", None, SearchCursor::Start, 10) + .search("abalone", SearchFilters::default(), SearchCursor::Start, 10) .await .unwrap(); assert!(abalone_hits.hits.is_empty(), "old title must be evicted"); let limpet_hits = idx - .search("limpet", None, SearchCursor::Start, 10) + .search("limpet", SearchFilters::default(), SearchCursor::Start, 10) .await .unwrap(); assert_eq!(limpet_hits.hits.len(), 1); @@ -541,7 +611,7 @@ mod tests { idx.remove(&uri).await; idx.flush().await; let hits = idx - .search("whelk", None, SearchCursor::Start, 10) + .search("whelk", SearchFilters::default(), SearchCursor::Start, 10) .await .unwrap(); assert!(hits.hits.is_empty()); @@ -563,19 +633,29 @@ mod tests { idx.flush().await; let page1 = idx - .search("anemone", None, SearchCursor::Start, 2) + .search("anemone", SearchFilters::default(), SearchCursor::Start, 2) .await .unwrap(); assert_eq!(page1.hits.len(), 2); let next = page1.next.expect("more pages"); let page2 = idx - .search("anemone", None, SearchCursor::At(next), 2) + .search( + "anemone", + SearchFilters::default(), + SearchCursor::At(next), + 2, + ) .await .unwrap(); assert_eq!(page2.hits.len(), 2); let next2 = page2.next.expect("more pages"); let page3 = idx - .search("anemone", None, SearchCursor::At(next2), 2) + .search( + "anemone", + SearchFilters::default(), + SearchCursor::At(next2), + 2, + ) .await .unwrap(); assert_eq!(page3.hits.len(), 1); @@ -594,7 +674,7 @@ mod tests { .await; idx.flush().await; let page = idx - .search(" ", None, SearchCursor::Start, 10) + .search(" ", SearchFilters::default(), SearchCursor::Start, 10) .await .unwrap(); assert!(page.hits.is_empty()); @@ -619,7 +699,7 @@ mod tests { .await; tokio::time::sleep(BATCH_INTERVAL * 3).await; let page = idx - .search("auto", None, SearchCursor::Start, 10) + .search("auto", SearchFilters::default(), SearchCursor::Start, 10) .await .unwrap(); assert_eq!(page.hits.len(), 1); diff --git a/crates/types/src/search.rs b/crates/types/src/search.rs index 2ad570f..a962377 100644 --- a/crates/types/src/search.rs +++ b/crates/types/src/search.rs @@ -2,9 +2,10 @@ use alloc::string::String; use alloc::vec::Vec; use core::future::Future; -use jacquard_common::DefaultStr; +use jacquard_common::types::ident::AtIdentifier; use jacquard_common::types::nsid::Nsid; -use jacquard_common::types::string::AtUri; +use jacquard_common::types::string::{AtUri, Did}; +use jacquard_common::{DefaultStr, IntoStatic}; use serde::Serialize; use crate::edges::{ExtractError, Record}; @@ -24,6 +25,9 @@ pub struct SearchDoc { pub nsid: Nsid, pub title: String, pub body: String, + pub author: Option>, + pub created_at: Option, + pub repo: Option>, } pub trait SearchSink: Send + Sync { @@ -116,17 +120,29 @@ impl SearchableRecord { } } +fn author_of(source: &AtUri) -> Option> { + match source.authority() { + AtIdentifier::Did(d) => Some(d.into_static()), + AtIdentifier::Handle(_) => None, + } +} + fn doc( source: &AtUri, nsid: &'static str, title: &str, body_parts: Vec, + created_at: Option, + repo: Option>, ) -> SearchDoc { SearchDoc { uri: source.clone(), nsid: nsid_static(nsid), title: title.to_owned(), body: body_parts.join(" "), + author: author_of(source), + created_at, + repo, } } @@ -146,7 +162,7 @@ fn profile_doc(source: &AtUri, r: &Profile) -> SearchDoc if let Some(p) = &r.pronouns { parts.push(p.as_str().to_owned()); } - doc(source, "sh.tangled.actor.profile", &title, parts) + doc(source, "sh.tangled.actor.profile", &title, parts, None, None) } fn repo_doc(source: &AtUri, r: &RepoRecord) -> SearchDoc { @@ -157,7 +173,14 @@ fn repo_doc(source: &AtUri, r: &RepoRecord) -> SearchDoc if let Some(topics) = &r.topics { parts.extend(topics.iter().map(|t| t.as_str().to_owned())); } - doc(source, "sh.tangled.repo", r.name.as_str(), parts) + doc( + source, + "sh.tangled.repo", + r.name.as_str(), + parts, + Some(r.created_at.timestamp()), + r.repo_did.clone(), + ) } fn issue_doc(source: &AtUri, r: &Issue) -> SearchDoc { @@ -166,7 +189,14 @@ fn issue_doc(source: &AtUri, r: &Issue) -> SearchDoc { .as_ref() .map(|b| Vec::from([b.as_str().to_owned()])) .unwrap_or_default(); - doc(source, "sh.tangled.repo.issue", r.title.as_str(), body) + doc( + source, + "sh.tangled.repo.issue", + r.title.as_str(), + body, + Some(r.created_at.timestamp()), + r.repo_did.clone(), + ) } fn issue_comment_doc(source: &AtUri, r: &IssueCommentRecord) -> SearchDoc { @@ -175,6 +205,8 @@ fn issue_comment_doc(source: &AtUri, r: &IssueCommentRecord, r: &Pull) -> SearchDoc { .as_ref() .map(|b| Vec::from([b.as_str().to_owned()])) .unwrap_or_default(); - doc(source, "sh.tangled.repo.pull", r.title.as_str(), body) + doc( + source, + "sh.tangled.repo.pull", + r.title.as_str(), + body, + Some(r.created_at.timestamp()), + r.target.repo_did.clone(), + ) } fn pull_comment_doc(source: &AtUri, r: &PullCommentRecord) -> SearchDoc { @@ -193,6 +232,8 @@ fn pull_comment_doc(source: &AtUri, r: &PullCommentRecord, r: &TangledString) -> Sear r.description.as_str().to_owned(), r.contents.as_str().to_owned(), ]), + Some(r.created_at.timestamp()), + None, ) } @@ -217,6 +260,8 @@ fn label_definition_doc( "sh.tangled.label.definition", r.name.as_str(), Vec::new(), + Some(r.created_at.timestamp()), + None, ) } -- 2.51.2