forked from rnadigital/agentcloud
-
Notifications
You must be signed in to change notification settings - Fork 0
/
Copy pathprocessing_incoming_messages.rs
86 lines (84 loc) · 3.99 KB
/
processing_incoming_messages.rs
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
use mongodb::Database;
use std::sync::Arc;
use tokio::sync::{RwLock};
use qdrant_client::client::QdrantClient;
use serde_json::Value;
use crate::llm::models::EmbeddingModels;
use crate::mongo::queries::get_embedding_model_and_embedding_key;
use crate::qdrant::helpers::embed_payload;
use crate::qdrant::utils::Qdrant;
use crate::utils::conversions::convert_serde_value_to_hashmap_string;
pub async fn process_messages(
qdrant_conn: Arc<RwLock<QdrantClient>>,
mongo_conn: Arc<RwLock<Database>>,
message: String,
datasource_id: String,
) {
// initiate variables
let mongodb_connection = mongo_conn.read().await;
// let redis_connection = redis_connection_pool.lock().await;
match serde_json::from_str(message.as_str()) {
Ok::<Value, _>(message_data) => {
match get_embedding_model_and_embedding_key(&mongodb_connection, datasource_id.as_str())
.await
{
Ok((model_parameter_result, embedding_field)) => match model_parameter_result {
Some(model_parameters) => {
let vector_length = model_parameters.embeddingLength as u64;
let embedding_model_name = model_parameters.model;
let embedding_model_name_clone = embedding_model_name.clone();
let ds_clone = datasource_id.clone();
let qdrant = Qdrant::new(qdrant_conn, datasource_id);
if let Value::Object(data_obj) = message_data {
let mut metadata = convert_serde_value_to_hashmap_string(data_obj);
if let Some(text_field) = embedding_field {
println!("text field: {}", text_field.as_str());
let text = metadata.remove(text_field.as_str()).unwrap();
metadata.insert("page_content".to_string(), text.to_owned());
let mongo_conn_clone = Arc::clone(&mongo_conn);
match embed_payload(
mongo_conn_clone,
&metadata,
&text,
Some(ds_clone),
EmbeddingModels::from(embedding_model_name),
)
.await
{
Ok(point_struct) => {
if let Ok(_bulk_upload_result) =
qdrant
.upsert_data_point_blocking(
point_struct,
Some(vector_length),
Some(embedding_model_name_clone),
)
.await
{
// let _ = redis_connection.increment_count(&"some_key".to_string(), 1);
}
}
Err(e) => {
eprintln!("An error occurred while upserting point structs to Qdrant: {}", e);
}
}
}
}
}
None => {
eprintln!("Model mongo object returned None!");
}
},
Err(e) => {
println!("An error occurred: {}", e);
}
}
}
Err(e) => {
eprintln!(
"An error occurred while attempting to convert message to JSON: {}",
e
);
}
}
}