diff --git a/apps/labrinth/.env.docker-compose b/apps/labrinth/.env.docker-compose index 81bd6934a9..57e7cd8e61 100644 --- a/apps/labrinth/.env.docker-compose +++ b/apps/labrinth/.env.docker-compose @@ -21,10 +21,10 @@ SEARCH_BACKEND=typesense MEILISEARCH_READ_ADDR=http://localhost:7700 MEILISEARCH_WRITE_ADDRS=http://localhost:7700 MEILISEARCH_KEY=modrinth -ELASTICSEARCH_URL=http://localhost:9200 +ELASTICSEARCH_URL=http://elasticsearch0:9200 ELASTICSEARCH_INDEX_PREFIX=labrinth -ELASTICSEARCH_USERNAME=elastic -ELASTICSEARCH_PASSWORD=elastic +ELASTICSEARCH_USERNAME= +ELASTICSEARCH_PASSWORD= SEARCH_INDEX_CHUNK_SIZE=5000 SEARCH_INCREMENTAL_INDEX_BATCH_DELAY_SECONDS=5 SEARCH_INCREMENTAL_INDEX_BATCH_MAX_SIZE=1000 diff --git a/apps/labrinth/src/env.rs b/apps/labrinth/src/env.rs index 20ce98f546..88978d006a 100644 --- a/apps/labrinth/src/env.rs +++ b/apps/labrinth/src/env.rs @@ -237,6 +237,11 @@ vars! { SEARCH_TYPESENSE_DEFAULT_BUCKETING: Json = Json(crate::search::backend::typesense::Bucketing::Buckets(5)); SEARCH_TYPESENSE_DEFAULT_MAX_CANDIDATES: usize = 24usize; + ELASTICSEARCH_URL: String = "http://localhost:9200"; + ELASTICSEARCH_INDEX_PREFIX: String = "labrinth"; + ELASTICSEARCH_USERNAME: String = ""; + ELASTICSEARCH_PASSWORD: String = ""; + ELASTICSEARCH_BULK_BATCH_SIZE: usize = 1000usize; // storage STORAGE_BACKEND: crate::file_hosting::FileHostKind = crate::file_hosting::FileHostKind::Local; diff --git a/apps/labrinth/src/search/backend/elasticsearch/filter.rs b/apps/labrinth/src/search/backend/elasticsearch/filter.rs new file mode 100644 index 0000000000..d21d44f1e6 --- /dev/null +++ b/apps/labrinth/src/search/backend/elasticsearch/filter.rs @@ -0,0 +1,393 @@ +use eyre::{Result, eyre}; +use serde_json::{Value, json}; + +use crate::search::filter::{ + FilterComparison, FilterCondition, FilterExpr, FilterLiteral, + FilterPredicate, +}; + +const MAX_DNF_CLAUSES: usize = 64; +const MAX_FILTER_DEPTH: usize = 64; +const MAX_FILTER_NODES: usize = 1024; + +#[derive(Clone, Copy, PartialEq, Eq)] +enum FilterScope { + Project, + Version, + Mixed, +} + +pub(super) struct ElasticsearchFilter { + pub query: Value, + pub has_version_filter: bool, +} + +pub(super) fn serialize_filter( + filter: &FilterExpr, +) -> Result { + let (nodes, depth) = filter_complexity(filter); + if nodes > MAX_FILTER_NODES { + return Err(eyre!("search filter has too many expressions")); + } + if depth > MAX_FILTER_DEPTH { + return Err(eyre!("search filter is nested too deeply")); + } + + let mut inner_hits_index = 0; + let query = plan(filter, &mut inner_hits_index)?; + Ok(ElasticsearchFilter { + query, + has_version_filter: inner_hits_index != 0, + }) +} + +fn plan(filter: &FilterExpr, inner_hits_index: &mut usize) -> Result { + match filter_scope(filter) { + FilterScope::Project => lower(filter), + FilterScope::Version => { + version_query(lower(filter)?, inner_hits_index) + } + FilterScope::Mixed => plan_mixed(filter, inner_hits_index), + } +} + +fn plan_mixed( + filter: &FilterExpr, + inner_hits_index: &mut usize, +) -> Result { + match filter { + FilterExpr::Or(expressions) => expressions + .iter() + .map(|expression| plan(expression, inner_hits_index)) + .collect::>>() + .map(or_query), + FilterExpr::And(expressions) + if expressions.iter().all(|expression| { + filter_scope(expression) != FilterScope::Mixed + }) => + { + plan_partitioned_and(expressions, inner_hits_index) + } + _ => { + let clauses = to_dnf(filter)?; + clauses + .into_iter() + .map(|clause| plan_clause(clause, inner_hits_index)) + .collect::>>() + .map(or_query) + } + } +} + +fn plan_partitioned_and( + expressions: &[FilterExpr], + inner_hits_index: &mut usize, +) -> Result { + let mut project = Vec::new(); + let mut version = Vec::new(); + for expression in expressions { + match filter_scope(expression) { + FilterScope::Project => project.push(lower(expression)?), + FilterScope::Version => version.push(lower(expression)?), + FilterScope::Mixed => { + return Err(eyre!("could not partition mixed search filter")); + } + } + } + if !version.is_empty() { + project.push(version_query( + and_query(version), + inner_hits_index, + )?); + } + Ok(and_query(project)) +} + +fn plan_clause( + predicates: Vec<&FilterPredicate>, + inner_hits_index: &mut usize, +) -> Result { + let mut project = Vec::new(); + let mut version = Vec::new(); + for predicate in predicates { + let query = predicate_query(predicate)?; + if is_version_filter_field(predicate.field.as_str()) { + version.push(query); + } else { + project.push(query); + } + } + if !version.is_empty() { + project.push(version_query( + and_query(version), + inner_hits_index, + )?); + } + Ok(and_query(project)) +} + +fn lower(filter: &FilterExpr) -> Result { + match filter { + FilterExpr::And(expressions) => expressions + .iter() + .map(lower) + .collect::>>() + .map(and_query), + FilterExpr::Or(expressions) => expressions + .iter() + .map(lower) + .collect::>>() + .map(or_query), + FilterExpr::Predicate(predicate) => predicate_query(predicate), + FilterExpr::Not(_) => { + Err(eyre!("search filter contains an unnormalized negation")) + } + } +} + +fn predicate_query(predicate: &FilterPredicate) -> Result { + let field = exact_field(predicate.field.as_str()); + match &predicate.condition { + FilterCondition::Compare { comparison, value } => { + let value = literal_value(value)?; + Ok(match comparison { + FilterComparison::Equal => { + json!({"term": {(field): {"value": value}}}) + } + FilterComparison::NotEqual => not_query(json!({ + "term": {(field): {"value": value}} + })), + FilterComparison::GreaterThan => { + json!({"range": {(field): {"gt": value}}}) + } + FilterComparison::GreaterThanOrEqual => { + json!({"range": {(field): {"gte": value}}}) + } + FilterComparison::LessThan => { + json!({"range": {(field): {"lt": value}}}) + } + FilterComparison::LessThanOrEqual => { + json!({"range": {(field): {"lte": value}}}) + } + }) + } + FilterCondition::In { values, negated } => { + let values = values + .iter() + .map(literal_value) + .collect::>>()?; + let query = json!({"terms": {(field): values}}); + Ok(if *negated { not_query(query) } else { query }) + } + FilterCondition::Exists { negated } => { + let query = json!({"exists": {"field": field}}); + Ok(if *negated { not_query(query) } else { query }) + } + } +} + +fn literal_value(literal: &FilterLiteral) -> Result { + match literal { + FilterLiteral::String(value) => Ok(Value::String(value.clone())), + FilterLiteral::Number(value) => serde_json::from_str(value) + .map_err(|error| eyre!("invalid numeric filter literal: {error}")), + FilterLiteral::Bool(value) => Ok(Value::Bool(*value)), + } +} + +fn exact_field(field: &str) -> &str { + match field { + "name" => "name.keyword", + "author" => "author.keyword", + "summary" => "summary.keyword", + "slug" => "slug.keyword", + _ => field, + } +} + +fn version_query( + query: Value, + inner_hits_index: &mut usize, +) -> Result { + let name = format!("matching_versions_{}", *inner_hits_index); + *inner_hits_index += 1; + Ok(json!({ + "has_child": { + "type": "version", + "score_mode": "none", + "query": query, + "inner_hits": { + "name": name, + "size": 1, + "_source": ["version_id", "version_published_timestamp"], + "sort": [ + {"version_published_timestamp": {"order": "desc"}} + ] + } + } + })) +} + +fn and_query(queries: Vec) -> Value { + match queries.len() { + 0 => json!({"match_all": {}}), + 1 => queries.into_iter().next().unwrap_or_default(), + _ => json!({"bool": {"filter": queries}}), + } +} + +fn or_query(queries: Vec) -> Value { + match queries.len() { + 0 => json!({"match_none": {}}), + 1 => queries.into_iter().next().unwrap_or_default(), + _ => json!({ + "bool": { + "should": queries, + "minimum_should_match": 1 + } + }), + } +} + +fn not_query(query: Value) -> Value { + json!({ + "bool": { + "must": [{"match_all": {}}], + "must_not": [query] + } + }) +} + +fn filter_scope(filter: &FilterExpr) -> FilterScope { + match filter { + FilterExpr::Predicate(predicate) => { + if is_version_filter_field(predicate.field.as_str()) { + FilterScope::Version + } else { + FilterScope::Project + } + } + FilterExpr::And(expressions) | FilterExpr::Or(expressions) => { + let mut scopes = expressions.iter().map(filter_scope); + let Some(first) = scopes.next() else { + return FilterScope::Project; + }; + if scopes.all(|scope| scope == first) { + first + } else { + FilterScope::Mixed + } + } + FilterExpr::Not(expression) => filter_scope(expression), + } +} + +fn is_version_filter_field(field: &str) -> bool { + matches!( + field, + "categories" + | "project_types" + | "environment" + | "game_versions" + | "client_side" + | "server_side" + ) +} + +fn to_dnf(filter: &FilterExpr) -> Result>> { + match filter { + FilterExpr::Predicate(predicate) => Ok(vec![vec![predicate]]), + FilterExpr::Or(expressions) => { + let mut clauses = Vec::new(); + for expression in expressions.iter() { + clauses.extend(to_dnf(expression)?); + if clauses.len() > MAX_DNF_CLAUSES { + return Err(eyre!( + "search filter has too many boolean clauses" + )); + } + } + Ok(clauses) + } + FilterExpr::And(expressions) => { + let mut clauses = vec![Vec::new()]; + for expression in expressions.iter() { + let right = to_dnf(expression)?; + if clauses.len().saturating_mul(right.len()) + > MAX_DNF_CLAUSES + { + return Err(eyre!( + "search filter has too many boolean clauses" + )); + } + clauses = clauses + .into_iter() + .flat_map(|left| { + right.iter().map(move |right| { + let mut clause = left.clone(); + clause.extend(right); + clause + }) + }) + .collect(); + } + Ok(clauses) + } + FilterExpr::Not(_) => { + Err(eyre!("search filter contains an unnormalized negation")) + } + } +} + +fn filter_complexity(filter: &FilterExpr) -> (usize, usize) { + match filter { + FilterExpr::Predicate(_) => (1, 1), + FilterExpr::And(expressions) | FilterExpr::Or(expressions) => { + expressions.iter().map(filter_complexity).fold( + (1, 1), + |(nodes, depth), (child_nodes, child_depth)| { + (nodes + child_nodes, depth.max(child_depth + 1)) + }, + ) + } + FilterExpr::Not(expression) => { + let (nodes, depth) = filter_complexity(expression); + (nodes + 1, depth + 1) + } + } +} + +#[cfg(test)] +mod tests { + use super::serialize_filter; + use crate::search::filter::{normalize, parse_expression}; + use serde_json::Value; + + fn serialize(input: &str) -> Value { + let filter = normalize(parse_expression(input).unwrap()); + serialize_filter(&filter).unwrap().query + } + + #[test] + fn correlated_version_filters_use_one_join() { + let query = serialize( + "categories = fabric AND game_versions = 1.21", + ); + assert_eq!(query.to_string().matches("has_child").count(), 1); + } + + #[test] + fn project_filters_do_not_use_a_join() { + let query = serialize("license = MIT"); + assert_eq!(query.to_string().matches("has_child").count(), 0); + assert_eq!(query["term"]["license"]["value"], "MIT"); + } + + #[test] + fn mixed_boolean_filters_preserve_version_correlation() { + let query = serialize( + "(license = MIT OR categories = fabric) AND game_versions = 1.21", + ); + assert_eq!(query.to_string().matches("has_child").count(), 2); + } +} diff --git a/apps/labrinth/src/search/backend/elasticsearch/mod.rs b/apps/labrinth/src/search/backend/elasticsearch/mod.rs new file mode 100644 index 0000000000..2e3e900217 --- /dev/null +++ b/apps/labrinth/src/search/backend/elasticsearch/mod.rs @@ -0,0 +1,1136 @@ +//! Search implementation backed by an Elasticsearch cluster. +//! +//! Projects and versions share an index and use an Elasticsearch join field. +//! This keeps version filters correlated without duplicating every version +//! into its project document. + +use async_trait::async_trait; +use eyre::{Result, eyre}; +use itertools::Itertools; +use reqwest::{Method, Response, StatusCode}; +use serde::Serialize; +use serde_json::{Map, Value, json}; +use tracing::{debug, info}; +use xredis::RedisPool; + +use crate::database::PgPool; +use crate::env::ENV; +use crate::routes::ApiError; +use crate::search::backend::{ + SearchIndex, combined_search_filters, parse_search_index, + parse_search_request, +}; +use crate::search::filter::{ + FilterExpr, from_legacy_v2_facets_json, normalize, parse_expression, +}; +use crate::search::indexing::index_local; +use crate::search::{ + ResultSearchProject, SearchBackend, SearchIndexUpdate, SearchRequest, + SearchResults, TasksCancelFilter, UploadSearchProject, UploadSearchVersion, +}; +use crate::util::error::Context; + +use self::filter::{ElasticsearchFilter, serialize_filter}; + +mod filter; + +const DELETE_FILTER_ID_BATCH_SIZE: usize = 1024; + +#[derive(Debug, Clone)] +pub struct ElasticsearchConfig { + pub url: String, + pub username: String, + pub password: String, + pub index_prefix: String, + pub meta_namespace: String, + pub index_chunk_size: i64, + pub bulk_batch_size: usize, +} + +impl ElasticsearchConfig { + pub fn new(meta_namespace: Option) -> Self { + Self { + url: ENV.ELASTICSEARCH_URL.clone(), + username: ENV.ELASTICSEARCH_USERNAME.clone(), + password: ENV.ELASTICSEARCH_PASSWORD.clone(), + index_prefix: ENV.ELASTICSEARCH_INDEX_PREFIX.clone(), + meta_namespace: meta_namespace.unwrap_or_default(), + index_chunk_size: ENV.SEARCH_INDEX_CHUNK_SIZE, + bulk_batch_size: ENV.ELASTICSEARCH_BULK_BATCH_SIZE, + } + } + + fn alias_name(&self) -> String { + if self.meta_namespace.is_empty() { + format!("{}_projects", self.index_prefix) + } else { + format!( + "{}_{}_projects", + self.meta_namespace, self.index_prefix + ) + } + } + + fn next_index_name(&self, alias: &str, use_alt: bool) -> String { + if use_alt { + format!("{alias}__alt") + } else { + format!("{alias}__current") + } + } +} + +struct ElasticsearchClient { + client: reqwest::Client, + base_url: String, + username: String, + password: String, +} + +impl ElasticsearchClient { + fn new(config: &ElasticsearchConfig) -> Self { + Self { + client: reqwest::Client::new(), + base_url: config.url.trim_end_matches('/').to_string(), + username: config.username.clone(), + password: config.password.clone(), + } + } + + fn request(&self, method: Method, path: &str) -> reqwest::RequestBuilder { + let request = self + .client + .request(method, format!("{}{}", self.base_url, path)); + if self.username.is_empty() { + request + } else { + request.basic_auth(&self.username, Some(&self.password)) + } + } + + async fn get_alias_target(&self, alias: &str) -> Result> { + let response = self + .request(Method::GET, &format!("/_alias/{alias}")) + .send() + .await + .wrap_err("failed to get Elasticsearch alias")?; + if response.status() == StatusCode::NOT_FOUND { + return Ok(None); + } + let body = response_json(response, "get Elasticsearch alias").await?; + Ok(body + .as_object() + .and_then(|indices| indices.keys().next()) + .cloned()) + } + + async fn index_exists(&self, index: &str) -> Result { + let response = self + .request(Method::HEAD, &format!("/{index}")) + .send() + .await + .wrap_err("failed to check Elasticsearch index existence")?; + Ok(response.status().is_success()) + } + + async fn create_index(&self, index: &str, schema: &Value) -> Result<()> { + let response = self + .request(Method::PUT, &format!("/{index}")) + .json(schema) + .send() + .await + .wrap_err("failed to create Elasticsearch index")?; + response_json(response, "create Elasticsearch index").await?; + Ok(()) + } + + async fn delete_index_if_exists(&self, index: &str) -> Result<()> { + let response = self + .request(Method::DELETE, &format!("/{index}")) + .send() + .await + .wrap_err("failed to delete Elasticsearch index")?; + if response.status() == StatusCode::NOT_FOUND { + return Ok(()); + } + response_json(response, "delete Elasticsearch index").await?; + Ok(()) + } + + async fn swap_alias( + &self, + alias: &str, + old_index: Option<&str>, + new_index: &str, + ) -> Result<()> { + let mut actions = Vec::new(); + if let Some(old_index) = old_index { + actions.push(json!({ + "remove": {"index": old_index, "alias": alias} + })); + } + actions.push(json!({ + "add": { + "index": new_index, + "alias": alias, + "is_write_index": true + } + })); + + let response = self + .request(Method::POST, "/_aliases") + .json(&json!({"actions": actions})) + .send() + .await + .wrap_err("failed to swap Elasticsearch alias")?; + response_json(response, "swap Elasticsearch alias").await?; + Ok(()) + } + + async fn bulk(&self, index: &str, body: String) -> Result<()> { + let response = self + .request( + Method::POST, + &format!("/{index}/_bulk?refresh=false"), + ) + .header("Content-Type", "application/x-ndjson") + .body(body) + .send() + .await + .wrap_err("failed to execute Elasticsearch bulk request")?; + let body = + response_json(response, "execute Elasticsearch bulk request") + .await?; + if body["errors"].as_bool() == Some(true) { + let failures = body["items"] + .as_array() + .into_iter() + .flatten() + .filter_map(|item| { + item.as_object()? + .values() + .next()? + .get("error") + .cloned() + }) + .unique() + .take(10) + .map(|error| error.to_string()) + .join("; "); + return Err(eyre!( + "Elasticsearch bulk request contained failures: {failures}" + )); + } + Ok(()) + } + + async fn delete_by_query( + &self, + index: &str, + query: &Value, + ) -> Result<()> { + let response = self + .request( + Method::POST, + &format!( + "/{index}/_delete_by_query?conflicts=proceed&refresh=false" + ), + ) + .json(&json!({"query": query})) + .send() + .await + .wrap_err("failed to delete Elasticsearch documents")?; + let body = + response_json(response, "delete Elasticsearch documents").await?; + if body["failures"] + .as_array() + .is_some_and(|failures| !failures.is_empty()) + { + return Err(eyre!( + "Elasticsearch delete-by-query contained failures: {}", + body["failures"] + )); + } + Ok(()) + } + + async fn refresh(&self, index: &str) -> Result<()> { + let response = self + .request(Method::POST, &format!("/{index}/_refresh")) + .send() + .await + .wrap_err("failed to refresh Elasticsearch index")?; + response_json(response, "refresh Elasticsearch index").await?; + Ok(()) + } + + async fn force_merge(&self, index: &str) -> Result<()> { + let response = self + .request( + Method::POST, + &format!( + "/{index}/_forcemerge?max_num_segments=1&flush=true" + ), + ) + .send() + .await + .wrap_err("failed to force-merge Elasticsearch index")?; + response_json(response, "force-merge Elasticsearch index").await?; + Ok(()) + } +} + +async fn response_json( + response: Response, + operation: &str, +) -> Result { + let status = response.status(); + let body = response + .text() + .await + .wrap_err_with(|| format!("failed to read response for {operation}"))?; + let json = serde_json::from_str(&body).unwrap_or_else(|_| { + json!({ + "unparsed_response": body + }) + }); + if !status.is_success() { + return Err(eyre!("{operation} failed ({status}): {json}")); + } + Ok(json) +} + +pub struct Elasticsearch { + pub config: ElasticsearchConfig, + client: ElasticsearchClient, +} + +impl Elasticsearch { + pub fn new(config: ElasticsearchConfig) -> Self { + let client = ElasticsearchClient::new(&config); + Self { config, client } + } + + fn index_schema() -> Value { + json!({ + "settings": { + "number_of_shards": 3, + "number_of_replicas": 1, + "refresh_interval": "30s", + "index.mapping.total_fields.limit": 5000, + "analysis": { + "char_filter": { + "hyphen_separator": { + "type": "pattern_replace", + "pattern": "-", + "replacement": " " + }, + "strip_symbols": { + "type": "pattern_replace", + "pattern": r"[^\p{L}\p{N}\s]", + "replacement": "" + } + }, + "filter": { + "typesense_stemmer": { + "type": "stemmer", + "language": "light_english" + } + }, + "analyzer": { + "typesense_text": { + "type": "custom", + "char_filter": ["strip_symbols"], + "tokenizer": "whitespace", + "filter": ["lowercase"] + }, + "typesense_hyphen_text": { + "type": "custom", + "char_filter": [ + "hyphen_separator", + "strip_symbols" + ], + "tokenizer": "whitespace", + "filter": ["lowercase"] + }, + "typesense_stemmed_text": { + "type": "custom", + "char_filter": ["strip_symbols"], + "tokenizer": "whitespace", + "filter": ["lowercase", "typesense_stemmer"] + } + } + } + }, + "mappings": { + "dynamic_templates": [ + { + "strings_as_keywords": { + "match_mapping_type": "string", + "mapping": { + "type": "keyword", + "ignore_above": 8191 + } + } + } + ], + "properties": { + "document_type": { + "type": "join", + "relations": {"project": "version"}, + "eager_global_ordinals": true + }, + "version_id": {"type": "keyword"}, + "project_id": {"type": "keyword"}, + "project_types": {"type": "keyword"}, + "all_project_types": {"type": "keyword"}, + "slug": { + "type": "text", + "analyzer": "typesense_text", + "index_options": "docs", + "norms": false, + "index_prefixes": { + "min_chars": 1, + "max_chars": 10 + }, + "fields": { + "keyword": {"type": "keyword", "ignore_above": 8191} + } + }, + "author": { + "type": "text", + "analyzer": "typesense_hyphen_text", + "index_options": "docs", + "norms": false, + "index_prefixes": { + "min_chars": 1, + "max_chars": 10 + }, + "fields": { + "keyword": {"type": "keyword", "ignore_above": 8191} + } + }, + "indexed_author": { + "type": "text", + "analyzer": "typesense_text", + "index_options": "docs", + "norms": false, + "index_prefixes": { + "min_chars": 1, + "max_chars": 10 + } + }, + "name": { + "type": "text", + "analyzer": "typesense_hyphen_text", + "index_options": "docs", + "norms": false, + "index_prefixes": { + "min_chars": 1, + "max_chars": 10 + }, + "fields": { + "keyword": {"type": "keyword", "ignore_above": 8191} + } + }, + "indexed_name": { + "type": "text", + "analyzer": "typesense_stemmed_text", + "index_options": "docs", + "norms": false, + "index_prefixes": { + "min_chars": 1, + "max_chars": 10 + } + }, + "summary": { + "type": "text", + "analyzer": "typesense_text", + "index_options": "docs", + "norms": false, + "index_prefixes": { + "min_chars": 1, + "max_chars": 10 + }, + "fields": { + "keyword": {"type": "keyword", "ignore_above": 8191} + } + }, + "categories": {"type": "keyword"}, + "project_categories": {"type": "keyword"}, + "display_categories": {"type": "keyword"}, + "license": {"type": "keyword"}, + "open_source": {"type": "boolean"}, + "environment": {"type": "keyword"}, + "game_versions": {"type": "keyword"}, + "client_side": {"type": "keyword"}, + "server_side": {"type": "keyword"}, + "dependency_project_ids": {"type": "keyword"}, + "compatible_dependency_project_ids": {"type": "keyword"}, + "downloads": {"type": "integer"}, + "log_downloads": {"type": "double"}, + "follows": {"type": "integer"}, + "created_timestamp": {"type": "long"}, + "modified_timestamp": {"type": "long"}, + "version_published_timestamp": {"type": "long"}, + "date_created": {"type": "date"}, + "date_modified": {"type": "date"}, + "project_loader_fields": {"type": "object", "enabled": false}, + "minecraft_java_server": { + "properties": { + "verified_plays_2w": {"type": "long"}, + "is_online": {"type": "boolean"}, + "ping": { + "properties": { + "data": { + "properties": { + "players_online": {"type": "integer"} + } + } + } + } + } + } + } + } + }) + } + + fn text_query(query: &str) -> Value { + if query.is_empty() { + return json!({"match_all": {}}); + } + + let tokens = query.split_whitespace().collect_vec(); + if tokens.is_empty() { + return json!({"match_all": {}}); + } + let fields = [ + ("name", 15), + ("indexed_name", 15), + ("slug", 10), + ("author", 3), + ("indexed_author", 3), + ("summary", 1), + ]; + let token_queries = + |token: &str, prefix: bool, exact_boost: bool| { + let mut queries = Vec::with_capacity(fields.len() * 3); + for (field, weight) in fields { + if exact_boost && field != "indexed_name" { + queries.push(json!({ + "constant_score": { + "filter": { + "match": { + (field): { + "query": token, + "fuzziness": 0 + } + } + }, + "boost": weight + 1 + } + })); + } + if prefix { + queries.push(json!({ + "constant_score": { + "filter": { + "match_bool_prefix": { + (field): {"query": token} + } + }, + "boost": weight + } + })); + } + queries.push(json!({ + "constant_score": { + "filter": { + "match": { + (field): { + "query": token, + "fuzziness": "AUTO:4,7", + "prefix_length": 1, + "max_expansions": 2 + } + } + }, + "boost": weight + } + })); + } + queries + }; + let queries_by_token = tokens + .iter() + .enumerate() + .map(|(index, token)| { + // Typesense caps candidates globally, while Elasticsearch + // expands them per field and shard. These bounds keep broad + // prefixes and stem-only matches out of the top relevance + // bucket while preserving short autocomplete queries. + let is_last = index == tokens.len() - 1; + let prefix = + is_last && (tokens.len() == 1 || token.chars().count() < 6); + let exact_boost = + tokens.len() > 1 || token.chars().count() >= 6; + token_queries(token, prefix, exact_boost) + }) + .collect_vec(); + let scoring_query = json!({ + "dis_max": { + "queries": queries_by_token.iter().flatten().collect_vec(), + "tie_breaker": 0 + } + }); + if tokens.len() == 1 { + scoring_query + } else { + json!({ + "bool": { + "must": [scoring_query], + "filter": queries_by_token + .into_iter() + .map(|queries| { + json!({ + "dis_max": { + "queries": queries, + "tie_breaker": 0 + } + }) + }) + .collect_vec() + } + }) + } + } + + fn sort(index: SearchIndex) -> Vec { + let descending = |field: &str| { + json!({(field): {"order": "desc", "missing": "_last"}}) + }; + let mut sort = match index { + SearchIndex::Relevance => vec![ + json!({"_score": {"order": "desc"}}), + descending("log_downloads"), + descending("version_published_timestamp"), + ], + SearchIndex::Downloads => vec![ + descending("log_downloads"), + descending("version_published_timestamp"), + ], + SearchIndex::Follows => vec![ + descending("follows"), + descending("version_published_timestamp"), + ], + SearchIndex::Updated => vec![ + descending("modified_timestamp"), + descending("version_published_timestamp"), + ], + SearchIndex::Newest => vec![ + descending("created_timestamp"), + descending("version_published_timestamp"), + ], + SearchIndex::MinecraftJavaServerVerifiedPlays2w => vec![ + json!({"_score": {"order": "desc"}}), + descending("minecraft_java_server.verified_plays_2w"), + descending("minecraft_java_server.is_online"), + ], + SearchIndex::MinecraftJavaServerPlayersOnline => vec![ + json!({"_score": {"order": "desc"}}), + descending("minecraft_java_server.is_online"), + descending( + "minecraft_java_server.ping.data.players_online", + ), + ], + }; + sort.push(json!({"project_id": {"order": "asc"}})); + sort + } + + fn build_filter( + info: &SearchRequest, + ) -> Result, ApiError> { + let facet_part = if let Some(facets_json) = info.facets.as_deref() { + from_legacy_v2_facets_json(facets_json) + .wrap_request_err("failed to parse facets")? + } else { + None + }; + + let filter_part = combined_search_filters(info) + .filter(|filter| !filter.trim().is_empty()) + .map(|filter| parse_expression(&filter)) + .transpose() + .wrap_request_err("failed to parse filters")?; + + FilterExpr::and([facet_part, filter_part].into_iter().flatten()) + .map(normalize) + .map(|filter| { + serialize_filter(&filter) + .wrap_request_err("failed to build search filter") + }) + .transpose() + } + + async fn existing_write_indices(&self) -> Result> { + let alias = self.config.alias_name(); + let mut indices = self + .client + .get_alias_target(&alias) + .await? + .into_iter() + .collect_vec(); + + for index in [ + self.config.next_index_name(&alias, false), + self.config.next_index_name(&alias, true), + ] { + if !indices.contains(&index) + && self.client.index_exists(&index).await? + { + indices.push(index); + } + } + Ok(indices) + } + + async fn import_projects( + &self, + indices: &[String], + documents: &[UploadSearchProject], + ) -> Result<()> { + let batch_size = self.config.bulk_batch_size.max(1); + for documents in documents.chunks(batch_size) { + let body = projects_to_bulk(documents)?; + for index in indices { + info!( + index, + document_count = documents.len(), + content_length_bytes = body.len(), + "sending Elasticsearch project bulk request" + ); + self.client.bulk(index, body.clone()).await?; + } + } + Ok(()) + } + + async fn import_versions( + &self, + indices: &[String], + documents: &[UploadSearchVersion], + ) -> Result<()> { + let batch_size = self.config.bulk_batch_size.max(1); + for documents in documents.chunks(batch_size) { + let body = versions_to_bulk(documents)?; + for index in indices { + info!( + index, + document_count = documents.len(), + content_length_bytes = body.len(), + "sending Elasticsearch version bulk request" + ); + self.client.bulk(index, body.clone()).await?; + } + } + Ok(()) + } + + async fn delete_ids( + &self, + field: &str, + ids: &[String], + ) -> Result<()> { + let indices = self.existing_write_indices().await?; + for ids in ids.chunks(DELETE_FILTER_ID_BATCH_SIZE) { + let query = json!({"terms": {(field): ids}}); + for index in &indices { + self.client.delete_by_query(index, &query).await?; + } + } + Ok(()) + } + + async fn refresh_write_indices(&self) -> Result<()> { + for index in self.existing_write_indices().await? { + self.client.refresh(&index).await?; + } + Ok(()) + } +} + +#[async_trait] +impl SearchBackend for Elasticsearch { + async fn search_for_project_raw( + &self, + info: &SearchRequest, + ) -> Result { + let parsed = parse_search_request(info)?; + let search_sort = + parse_search_index(parsed.index, info.new_filters.as_deref())?; + let filter = Self::build_filter(info)?; + + let mut filters = + vec![json!({"term": {"document_type": "project"}})]; + if let Some(filter) = &filter { + filters.push(filter.query.clone()); + } + let query = json!({ + "bool": { + "must": [Self::text_query(parsed.query)], + "filter": filters + } + }); + let body = json!({ + "from": parsed.offset, + "size": parsed.hits_per_page, + "track_total_hits": true, + "query": query, + "sort": Self::sort(search_sort.index) + }); + + let alias = self.config.alias_name(); + let response = self + .client + .request(Method::POST, &format!("/{alias}/_search")) + .json(&body) + .send() + .await + .wrap_internal_err("failed to execute Elasticsearch search")?; + let body = response_json(response, "execute Elasticsearch search") + .await + .map_err(ApiError::Internal)?; + + let total_hits = body["hits"]["total"]["value"] + .as_u64() + .unwrap_or_default() as usize; + let hits = body["hits"]["hits"] + .as_array() + .into_iter() + .flatten() + .filter_map(|hit| { + let mut document = hit["_source"].clone(); + let object = document.as_object_mut()?; + object.remove("document_type"); + if filter + .as_ref() + .is_some_and(|filter| filter.has_version_filter) + { + if let Some(version_id) = matching_version_id(hit) { + object.insert( + "version_id".to_string(), + Value::String(version_id), + ); + } + } + + let metadata = info.show_metadata.then(|| { + json!({ + "score": hit["_score"], + "sort": hit["sort"] + }) + }); + let mut result: ResultSearchProject = + serde_json::from_value::(document) + .ok()? + .into(); + result.search_metadata = metadata; + Some(result) + }) + .collect(); + + Ok(SearchResults { + hits, + page: parsed.page, + hits_per_page: parsed.hits_per_page, + total_hits, + }) + } + + async fn rebuild_index( + &self, + ro_pool: PgPool, + redis: RedisPool, + ) -> Result<()> { + info!("starting Elasticsearch project indexing"); + let alias = self.config.alias_name(); + let current = self.client.get_alias_target(&alias).await?; + let use_alt = + !current.as_deref().is_some_and(|name| name.ends_with("__alt")); + let next = self.config.next_index_name(&alias, use_alt); + + info!(index = next, "creating Elasticsearch shadow index"); + self.client.delete_index_if_exists(&next).await?; + self.client + .create_index(&next, &Self::index_schema()) + .await?; + + let mut cursor = 0_i64; + let mut chunk_index = 0_usize; + let mut total_projects = 0_usize; + let mut total_versions = 0_usize; + + loop { + info!("fetching index chunk {chunk_index}"); + chunk_index += 1; + let (documents, next_cursor) = index_local( + &ro_pool, + &redis, + cursor, + self.config.index_chunk_size, + ) + .await + .wrap_err("failed to fetch projects from local DB")?; + if documents.projects.is_empty() { + info!( + "no more documents; indexed {total_projects} projects and {total_versions} versions in {chunk_index} chunks" + ); + break; + } + + total_projects += documents.projects.len(); + total_versions += documents.versions.len(); + cursor = next_cursor; + self.import_projects( + std::slice::from_ref(&next), + &documents.projects, + ) + .await?; + self.import_versions( + std::slice::from_ref(&next), + &documents.versions, + ) + .await?; + } + + self.client.refresh(&next).await?; + info!("force-merging Elasticsearch shadow index"); + self.client.force_merge(&next).await?; + info!("swapping Elasticsearch index alias"); + self.client + .swap_alias(&alias, current.as_deref(), &next) + .await?; + if let Some(old) = current { + self.client.delete_index_if_exists(&old).await?; + } + info!("Elasticsearch indexing complete"); + Ok(()) + } + + async fn apply_update( + &self, + update: SearchIndexUpdate<'_>, + ) -> Result<()> { + let removed_project_ids = update + .removed_projects + .iter() + .map(ToString::to_string) + .collect::>(); + if !removed_project_ids.is_empty() { + self.delete_ids("project_id", &removed_project_ids).await?; + } + + let version_ids = update + .removed_versions + .iter() + .map(ToString::to_string) + .collect::>(); + if !version_ids.is_empty() { + self.delete_ids("version_id", &version_ids).await?; + } + + let indices = self.existing_write_indices().await?; + if !update.projects.is_empty() { + debug!( + ?indices, + num_documents = update.projects.len(), + "replacing Elasticsearch project documents" + ); + self.import_projects(&indices, update.projects).await?; + } + if !update.versions.is_empty() { + debug!( + ?indices, + num_documents = update.versions.len(), + "replacing Elasticsearch version documents" + ); + self.import_versions(&indices, update.versions).await?; + } + self.refresh_write_indices().await?; + debug!("done applying Elasticsearch search index update"); + Ok(()) + } + + async fn tasks(&self) -> Result { + let response = self + .client + .request(Method::GET, "/_tasks?detailed=true") + .send() + .await + .wrap_err("failed to get Elasticsearch tasks")?; + response_json(response, "get Elasticsearch tasks").await + } + + async fn tasks_cancel(&self, filter: &TasksCancelFilter) -> Result<()> { + match filter { + TasksCancelFilter::All => { + let response = self + .client + .request(Method::POST, "/_tasks/_cancel") + .send() + .await + .wrap_err("failed to cancel Elasticsearch tasks")?; + response_json(response, "cancel Elasticsearch tasks").await?; + } + TasksCancelFilter::AllEnqueued => { + // Elasticsearch executes operations immediately and does not + // expose an enqueued-task state. + } + TasksCancelFilter::Indexes { indexes } => { + let tasks = self.tasks().await?; + for (_node_id, node) in tasks["nodes"] + .as_object() + .into_iter() + .flatten() + { + for (task_id, task) in node["tasks"] + .as_object() + .into_iter() + .flatten() + { + let description = + task["description"].as_str().unwrap_or_default(); + if indexes + .iter() + .any(|index| description.contains(index)) + { + let response = self + .client + .request( + Method::POST, + &format!( + "/_tasks/{task_id}/_cancel" + ), + ) + .send() + .await + .wrap_err( + "failed to cancel Elasticsearch task", + )?; + response_json( + response, + "cancel Elasticsearch task", + ) + .await?; + } + } + } + } + } + Ok(()) + } +} + +fn matching_version_id(hit: &Value) -> Option { + hit["inner_hits"] + .as_object()? + .values() + .filter_map(|inner_hits| { + let source = inner_hits["hits"]["hits"] + .as_array()? + .first()? + .get("_source")?; + Some(( + source["version_published_timestamp"].as_i64()?, + source["version_id"].as_str()?.to_string(), + )) + }) + .max_by_key(|(published, _)| *published) + .map(|(_, version_id)| version_id) +} + +fn projects_to_bulk(documents: &[UploadSearchProject]) -> Result { + let mut output = String::new(); + for document in documents { + let id = format!("project:{}", document.project_id); + push_json_line( + &mut output, + &json!({ + "index": { + "_id": id, + "routing": document.project_id + } + }), + )?; + + let mut source = serde_json::to_value(document) + .wrap_err("failed to serialize `UploadSearchProject`")?; + let object = source + .as_object_mut() + .ok_or_else(|| eyre!("project search document is not an object"))?; + object.insert( + "document_type".to_string(), + Value::String("project".to_string()), + ); + add_server_online_field(object); + push_json_line(&mut output, &source)?; + } + Ok(output) +} + +fn versions_to_bulk(documents: &[UploadSearchVersion]) -> Result { + let mut output = String::new(); + for document in documents { + let id = format!("version:{}", document.version_id); + push_json_line( + &mut output, + &json!({ + "index": { + "_id": id, + "routing": document.project_id + } + }), + )?; + + let mut source = serde_json::to_value(document) + .wrap_err("failed to serialize `UploadSearchVersion`")?; + source + .as_object_mut() + .ok_or_else(|| eyre!("version search document is not an object"))? + .insert( + "document_type".to_string(), + json!({ + "name": "version", + "parent": format!("project:{}", document.project_id) + }), + ); + push_json_line(&mut output, &source)?; + } + Ok(output) +} + +fn add_server_online_field(object: &mut Map) { + let Some(server) = object + .get_mut("minecraft_java_server") + .and_then(Value::as_object_mut) + else { + return; + }; + let is_online = server + .get("ping") + .and_then(Value::as_object) + .and_then(|ping| ping.get("data")) + .is_some_and(|data| !data.is_null()); + server.insert("is_online".to_string(), Value::Bool(is_online)); +} + +fn push_json_line( + output: &mut String, + value: &T, +) -> Result<()> { + output.push_str(&serde_json::to_string(value)?); + output.push('\n'); + Ok(()) +} diff --git a/apps/labrinth/src/search/backend/mod.rs b/apps/labrinth/src/search/backend/mod.rs index 544cf2a557..cb38ae6782 100644 --- a/apps/labrinth/src/search/backend/mod.rs +++ b/apps/labrinth/src/search/backend/mod.rs @@ -1,8 +1,10 @@ mod common; +pub mod elasticsearch; pub mod typesense; pub use common::{ ParsedSearchRequest, SearchIndex, SearchSort, combined_search_filters, parse_search_index, parse_search_request, }; +pub use elasticsearch::{Elasticsearch, ElasticsearchConfig}; pub use typesense::{Typesense, TypesenseConfig}; diff --git a/apps/labrinth/src/search/mod.rs b/apps/labrinth/src/search/mod.rs index f98f9f4cb1..690efe149a 100644 --- a/apps/labrinth/src/search/mod.rs +++ b/apps/labrinth/src/search/mod.rs @@ -189,6 +189,7 @@ pub enum TasksCancelFilter { #[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] pub enum SearchBackendKind { Typesense, + Elasticsearch, } #[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, strum::EnumIter)] @@ -224,6 +225,7 @@ impl FromStr for SearchBackendKind { fn from_str(s: &str) -> Result { Ok(match s { "typesense" => SearchBackendKind::Typesense, + "elasticsearch" => SearchBackendKind::Elasticsearch, _ => return Err(InvalidSearchBackendKind), }) } @@ -438,5 +440,9 @@ pub fn backend(meta_namespace: Option) -> Box { let config = backend::TypesenseConfig::new(meta_namespace); Box::new(backend::Typesense::new(config)) } + SearchBackendKind::Elasticsearch => { + let config = backend::ElasticsearchConfig::new(meta_namespace); + Box::new(backend::Elasticsearch::new(config)) + } } } diff --git a/docker-compose.yml b/docker-compose.yml index def52d2862..64112bdbc0 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -32,6 +32,102 @@ services: interval: 3s timeout: 5s retries: 3 + elasticsearch0: + image: docker.elastic.co/elasticsearch/elasticsearch:9.4.4 + container_name: labrinth-elasticsearch0 + restart: on-failure + networks: + - elasticsearch-mesh + ports: + - '127.0.0.1:9200:9200' + volumes: + - elasticsearch0-data:/usr/share/elasticsearch/data + environment: + node.name: elasticsearch0 + cluster.name: labrinth-elasticsearch + discovery.seed_hosts: elasticsearch0,elasticsearch1,elasticsearch2 + cluster.initial_master_nodes: elasticsearch0,elasticsearch1,elasticsearch2 + bootstrap.memory_lock: 'true' + xpack.security.enabled: 'false' + xpack.security.enrollment.enabled: 'false' + ES_JAVA_OPTS: -Xms512m -Xmx512m + ulimits: + memlock: + soft: -1 + hard: -1 + healthcheck: + test: + [ + 'CMD-SHELL', + 'curl --fail http://localhost:9200/_cluster/health?wait_for_status=yellow', + ] + interval: 5s + timeout: 5s + retries: 30 + elasticsearch1: + image: docker.elastic.co/elasticsearch/elasticsearch:9.4.4 + container_name: labrinth-elasticsearch1 + restart: on-failure + networks: + - elasticsearch-mesh + ports: + - '127.0.0.1:9201:9200' + volumes: + - elasticsearch1-data:/usr/share/elasticsearch/data + environment: + node.name: elasticsearch1 + cluster.name: labrinth-elasticsearch + discovery.seed_hosts: elasticsearch0,elasticsearch1,elasticsearch2 + cluster.initial_master_nodes: elasticsearch0,elasticsearch1,elasticsearch2 + bootstrap.memory_lock: 'true' + xpack.security.enabled: 'false' + xpack.security.enrollment.enabled: 'false' + ES_JAVA_OPTS: -Xms512m -Xmx512m + ulimits: + memlock: + soft: -1 + hard: -1 + healthcheck: + test: + [ + 'CMD-SHELL', + 'curl --fail http://localhost:9200/_cluster/health?wait_for_status=yellow', + ] + interval: 5s + timeout: 5s + retries: 30 + elasticsearch2: + image: docker.elastic.co/elasticsearch/elasticsearch:9.4.4 + container_name: labrinth-elasticsearch2 + restart: on-failure + networks: + - elasticsearch-mesh + ports: + - '127.0.0.1:9202:9200' + volumes: + - elasticsearch2-data:/usr/share/elasticsearch/data + environment: + node.name: elasticsearch2 + cluster.name: labrinth-elasticsearch + discovery.seed_hosts: elasticsearch0,elasticsearch1,elasticsearch2 + cluster.initial_master_nodes: elasticsearch0,elasticsearch1,elasticsearch2 + bootstrap.memory_lock: 'true' + xpack.security.enabled: 'false' + xpack.security.enrollment.enabled: 'false' + ES_JAVA_OPTS: -Xms512m -Xmx512m + ulimits: + memlock: + soft: -1 + hard: -1 + healthcheck: + test: + [ + 'CMD-SHELL', + 'curl --fail http://localhost:9200/_cluster/health?wait_for_status=yellow', + ] + interval: 5s + timeout: 5s + retries: 30 meilisearch0: image: getmeili/meilisearch:v1.12.0 container_name: labrinth-meilisearch0 @@ -398,7 +494,7 @@ services: depends_on: postgres_db: condition: service_healthy - meilisearch: + meilisearch0: condition: service_healthy elasticsearch0: condition: service_healthy @@ -481,6 +577,8 @@ services: volumes: - ./apps/labrinth/nginx/meili-lb.conf:/etc/nginx/conf.d/default.conf:ro networks: + elasticsearch-mesh: + driver: bridge meilisearch-mesh: driver: bridge redis-cluster-mesh: