mirror of
https://github.com/meilisearch/meilisearch.git
synced 2024-11-30 00:55:00 +08:00
factorize the code between the two documents batcher
This commit is contained in:
parent
a9523146a3
commit
e64ba122e1
@ -34,6 +34,19 @@ mod segment {
|
|||||||
|
|
||||||
const SEGMENT_API_KEY: &str = "vHi89WrNDckHSQssyUJqLvIyp2QFITSC";
|
const SEGMENT_API_KEY: &str = "vHi89WrNDckHSQssyUJqLvIyp2QFITSC";
|
||||||
|
|
||||||
|
pub fn extract_user_agents(request: &HttpRequest) -> Vec<String> {
|
||||||
|
request
|
||||||
|
.headers()
|
||||||
|
.get(USER_AGENT)
|
||||||
|
.map(|header| header.to_str().ok())
|
||||||
|
.flatten()
|
||||||
|
.unwrap_or("unknown")
|
||||||
|
.split(";")
|
||||||
|
.map(str::trim)
|
||||||
|
.map(ToString::to_string)
|
||||||
|
.collect()
|
||||||
|
}
|
||||||
|
|
||||||
pub struct SegmentAnalytics {
|
pub struct SegmentAnalytics {
|
||||||
user: User,
|
user: User,
|
||||||
opt: Opt,
|
opt: Opt,
|
||||||
@ -196,7 +209,7 @@ mod segment {
|
|||||||
query: &SearchQuery,
|
query: &SearchQuery,
|
||||||
request: &HttpRequest,
|
request: &HttpRequest,
|
||||||
) {
|
) {
|
||||||
let user_agent = SearchBatcher::extract_user_agents(request);
|
let user_agent = extract_user_agents(request);
|
||||||
let sorted = query.sort.is_some() as usize;
|
let sorted = query.sort.is_some() as usize;
|
||||||
let sort_with_geo_point = query
|
let sort_with_geo_point = query
|
||||||
.sort
|
.sort
|
||||||
@ -275,6 +288,36 @@ mod segment {
|
|||||||
search_batcher.max_offset = search_batcher.max_offset.max(max_offset);
|
search_batcher.max_offset = search_batcher.max_offset.max(max_offset);
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn batch_documents(
|
||||||
|
&'static self,
|
||||||
|
batcher: &'static Mutex<DocumentsBatcher>,
|
||||||
|
documents_query: &UpdateDocumentsQuery,
|
||||||
|
index_creation: bool,
|
||||||
|
request: &HttpRequest,
|
||||||
|
) {
|
||||||
|
let user_agents = extract_user_agents(request);
|
||||||
|
let primary_key = documents_query.primary_key.clone();
|
||||||
|
let content_type = request
|
||||||
|
.headers()
|
||||||
|
.get(CONTENT_TYPE)
|
||||||
|
.map(|s| s.to_str().unwrap_or("unkown"))
|
||||||
|
.unwrap()
|
||||||
|
.to_string();
|
||||||
|
|
||||||
|
tokio::spawn(async move {
|
||||||
|
let mut lock = batcher.lock().await;
|
||||||
|
for user_agent in user_agents {
|
||||||
|
lock.user_agents.insert(user_agent);
|
||||||
|
}
|
||||||
|
lock.content_types.insert(content_type);
|
||||||
|
if let Some(primary_key) = primary_key {
|
||||||
|
lock.primary_keys.insert(primary_key);
|
||||||
|
}
|
||||||
|
lock.index_creation |= index_creation;
|
||||||
|
// drop the lock here
|
||||||
|
});
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
#[async_trait::async_trait]
|
#[async_trait::async_trait]
|
||||||
@ -333,30 +376,12 @@ mod segment {
|
|||||||
index_creation: bool,
|
index_creation: bool,
|
||||||
request: &HttpRequest,
|
request: &HttpRequest,
|
||||||
) {
|
) {
|
||||||
let user_agents = request
|
self.batch_documents(
|
||||||
.headers()
|
&self.documents_added_batcher,
|
||||||
.get(USER_AGENT)
|
documents_query,
|
||||||
.map(|header| header.to_str().unwrap_or("unknown").to_string());
|
index_creation,
|
||||||
let primary_key = documents_query.primary_key.clone();
|
request,
|
||||||
let content_type = request
|
)
|
||||||
.headers()
|
|
||||||
.get(CONTENT_TYPE)
|
|
||||||
.map(|s| s.to_str().unwrap_or("unkown"))
|
|
||||||
.unwrap()
|
|
||||||
.to_string();
|
|
||||||
|
|
||||||
tokio::spawn(async move {
|
|
||||||
let mut lock = self.documents_added_batcher.lock().await;
|
|
||||||
for user_agent in user_agents {
|
|
||||||
lock.user_agents.insert(user_agent);
|
|
||||||
}
|
|
||||||
lock.content_types.insert(content_type);
|
|
||||||
if let Some(primary_key) = primary_key {
|
|
||||||
lock.primary_keys.insert(primary_key);
|
|
||||||
}
|
|
||||||
lock.index_creation |= index_creation;
|
|
||||||
// drop the lock here
|
|
||||||
});
|
|
||||||
}
|
}
|
||||||
|
|
||||||
fn update_documents(
|
fn update_documents(
|
||||||
@ -365,30 +390,12 @@ mod segment {
|
|||||||
index_creation: bool,
|
index_creation: bool,
|
||||||
request: &HttpRequest,
|
request: &HttpRequest,
|
||||||
) {
|
) {
|
||||||
let user_agents = request
|
self.batch_documents(
|
||||||
.headers()
|
&self.documents_updated_batcher,
|
||||||
.get(USER_AGENT)
|
documents_query,
|
||||||
.map(|header| header.to_str().unwrap_or("unknown").to_string());
|
index_creation,
|
||||||
let primary_key = documents_query.primary_key.clone();
|
request,
|
||||||
let content_type = request
|
)
|
||||||
.headers()
|
|
||||||
.get(CONTENT_TYPE)
|
|
||||||
.map(|s| s.to_str().unwrap_or("unkown"))
|
|
||||||
.unwrap()
|
|
||||||
.to_string();
|
|
||||||
|
|
||||||
tokio::spawn(async move {
|
|
||||||
let mut lock = self.documents_updated_batcher.lock().await;
|
|
||||||
for user_agent in user_agents {
|
|
||||||
lock.user_agents.insert(user_agent);
|
|
||||||
}
|
|
||||||
lock.content_types.insert(content_type);
|
|
||||||
if let Some(primary_key) = primary_key {
|
|
||||||
lock.primary_keys.insert(primary_key);
|
|
||||||
}
|
|
||||||
lock.index_creation |= index_creation;
|
|
||||||
// drop the lock here
|
|
||||||
});
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@ -435,19 +442,6 @@ mod segment {
|
|||||||
}
|
}
|
||||||
|
|
||||||
impl SearchBatcher {
|
impl SearchBatcher {
|
||||||
pub fn extract_user_agents(request: &HttpRequest) -> Vec<String> {
|
|
||||||
request
|
|
||||||
.headers()
|
|
||||||
.get(USER_AGENT)
|
|
||||||
.map(|header| header.to_str().ok())
|
|
||||||
.flatten()
|
|
||||||
.unwrap_or("unknown")
|
|
||||||
.split(";")
|
|
||||||
.map(str::trim)
|
|
||||||
.map(ToString::to_string)
|
|
||||||
.collect()
|
|
||||||
}
|
|
||||||
|
|
||||||
pub fn into_event(mut self, user: &User, event_name: &str) -> Option<Track> {
|
pub fn into_event(mut self, user: &User, event_name: &str) -> Option<Track> {
|
||||||
if self.total_received == 0 {
|
if self.total_received == 0 {
|
||||||
None
|
None
|
||||||
|
Loading…
Reference in New Issue
Block a user