use derive_new::new;
use elasticsearch::{http::transport::Transport, Elasticsearch, IndexParts, SearchParts};
use serde_json::{json, Value};
pub struct Message {
pub user: String,
pub content: String,
pub id: String,
}
#[derive(new)]
pub struct SearchEngine {
url: String,
index_name: String,
}
impl SearchEngine {
fn connect_es(&self) -> Result<Elasticsearch, Box<dyn std::error::Error>> {
let transport = Transport::single_node(&self.url)?;
Ok(Elasticsearch::new(transport))
}
pub async fn save(&self, msg: &Message) -> Result<(), Box<dyn std::error::Error>> {
let client = self.connect_es()?;
let response = client
.index(IndexParts::IndexId(&self.index_name, &msg.id))
.body(json! ({
"user": msg.user,
"content": msg.content,
}))
.send()
.await?;
let json: Value = response.json().await?;
log::info!("saved {json}");
Ok(())
}
pub async fn search(&self, txt: &str, user: &str) -> Result<Value, Box<dyn std::error::Error>> {
log::info!("searching....");
let client = self.connect_es()?;
let search_json = json!({
"query": {
"bool" : {
"must": [
{
"match": {
"user": user,
}
},
{
"match": {
"content": txt
}
}
]
}
}
});
let response = client
.search(SearchParts::Index(&[&self.index_name]))
.body(search_json)
.send()
.await?;
log::info!("search ok");
let json: Value = response.json().await?;
log::info!("{json}");
Ok(json)
}
}
#[cfg(test)]
mod tests {
use log4rs::config::Deserializers;
use crate::{SearchEngine, Message};
#[tokio::test]
async fn test() {
log4rs::init_file("log4rs.yml", Deserializers::default()).unwrap();
let se = SearchEngine::new("http://8.141.144.181:9043".to_owned(), "interlinked".to_owned());
let msg = Message {
user: "test_user".to_owned(),
content: "test_content".to_owned(),
id: "chat_2_msg_1".to_owned(),
};
se.save(&msg).await.unwrap();
se.save(&msg).await.unwrap();
se.search("test", "test_user").await.unwrap();
}
}