queen_search 0.1.0

elasticsearch adpater for dramaverse queen
Documentation
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();
    }
}