use std::sync::Arc;
use velesdb_core::{
Database as CoreDatabase, Filter, GatedRead, QueryOperationKind,
VectorCollection as CoreCollection,
};
use crate::types::{
IndividualSearchRequest, MobileAdvancedConfig, MobileCollectionDiagnostics,
MobileCollectionStats, MobileIndexInfo, MobileQueryLimits, MobileStreamingConfig,
SearchQuality, SearchResult, VelesError, VelesPoint,
};
#[derive(uniffi::Object)]
pub struct VelesCollection {
pub(crate) inner: CoreCollection,
pub(crate) db: Arc<CoreDatabase>,
pub(crate) name: String,
}
fn to_mobile_results(results: Vec<velesdb_core::SearchResult>) -> Vec<SearchResult> {
results
.into_iter()
.map(|r| SearchResult {
id: r.point.id,
score: r.score,
payload: None,
})
.collect()
}
pub(crate) fn and_scope(caller: Option<Filter>, scope: Option<Filter>) -> Option<Filter> {
use velesdb_core::Condition;
match (caller, scope) {
(None, None) => None,
(Some(c), None) => Some(c),
(None, Some(s)) => Some(s),
(Some(c), Some(s)) => Some(Filter::new(Condition::And {
conditions: vec![c.condition, s.condition],
})),
}
}
pub(crate) fn deny_if_scoped(scope: Option<Filter>, context: &str) -> Result<(), VelesError> {
if scope.is_some() {
return Err(VelesError::database(format!(
"{context} cannot honor the governance scope filter returned by the observer \
(this entry point has no metadata-filtered leaf); refusing to run unscoped"
)));
}
Ok(())
}
#[uniffi::export]
impl VelesCollection {
pub fn search(&self, vector: Vec<f32>, limit: u32) -> Result<Vec<SearchResult>, VelesError> {
let results = self.db.gated_search(
&self.name,
None,
None,
GatedRead::Dense {
query: &vector,
k: usize::try_from(limit).unwrap_or(usize::MAX),
ef: None,
quality: None,
filter: None,
},
)?;
Ok(to_mobile_results(results))
}
pub fn search_with_quality(
&self,
vector: Vec<f32>,
limit: u32,
quality: SearchQuality,
) -> Result<Vec<SearchResult>, VelesError> {
let results = self.db.gated_search(
&self.name,
None,
None,
GatedRead::Dense {
query: &vector,
k: usize::try_from(limit).unwrap_or(usize::MAX),
ef: None,
quality: Some(quality.into()),
filter: None,
},
)?;
Ok(to_mobile_results(results))
}
pub fn upsert(&self, point: VelesPoint) -> Result<(), VelesError> {
let core_point = parse_point(point)?;
self.inner.upsert(vec![core_point])?;
Ok(())
}
pub fn upsert_batch(&self, points: Vec<VelesPoint>) -> Result<(), VelesError> {
let core_points: Result<Vec<velesdb_core::Point>, VelesError> =
points.into_iter().map(parse_point).collect();
self.inner.upsert(core_points?)?;
Ok(())
}
pub fn delete(&self, id: u64) -> Result<(), VelesError> {
self.inner.delete(&[id])?;
Ok(())
}
#[allow(clippy::cast_possible_truncation)]
pub fn count(&self) -> u64 {
self.inner.config().point_count as u64
}
#[allow(clippy::cast_possible_truncation)]
pub fn dimension(&self) -> u32 {
self.inner.config().dimension as u32
}
pub fn get(&self, ids: Vec<u64>) -> Vec<VelesPoint> {
self.inner
.get(&ids)
.into_iter()
.flatten()
.map(|p| VelesPoint {
id: p.id,
vector: p.vector,
payload: p.payload.map(|v| v.to_string()),
})
.collect()
}
pub fn get_by_id(&self, id: u64) -> Option<VelesPoint> {
self.inner
.get(&[id])
.into_iter()
.flatten()
.next()
.map(|p| VelesPoint {
id: p.id,
vector: p.vector,
payload: p.payload.map(|v| v.to_string()),
})
}
pub fn is_metadata_only(&self) -> bool {
self.inner.config().metadata_only
}
pub fn text_search(&self, query: String, limit: u32) -> Result<Vec<SearchResult>, VelesError> {
let results = self.db.gated_search(
&self.name,
None,
None,
GatedRead::Text {
query: &query,
k: usize::try_from(limit).unwrap_or(usize::MAX),
filter: None,
},
)?;
Ok(to_mobile_results(results))
}
pub fn hybrid_search(
&self,
vector: Vec<f32>,
text_query: String,
limit: u32,
vector_weight: f32,
) -> Result<Vec<SearchResult>, VelesError> {
let results = self.db.gated_search(
&self.name,
None,
None,
GatedRead::Hybrid {
vector: &vector,
text: &text_query,
k: usize::try_from(limit).unwrap_or(usize::MAX),
alpha: Some(vector_weight),
filter: None,
},
)?;
Ok(to_mobile_results(results))
}
pub fn search_with_filter(
&self,
vector: Vec<f32>,
limit: u32,
filter_json: String,
) -> Result<Vec<SearchResult>, VelesError> {
let filter: Filter = serde_json::from_str(&filter_json)
.map_err(|e| VelesError::database(format!("Invalid filter JSON: {e}")))?;
let results = self.db.gated_search(
&self.name,
None,
None,
GatedRead::Dense {
query: &vector,
k: usize::try_from(limit).unwrap_or(usize::MAX),
ef: None,
quality: None,
filter: Some(&filter),
},
)?;
Ok(to_mobile_results(results))
}
pub fn batch_search(
&self,
searches: Vec<IndividualSearchRequest>,
) -> Result<Vec<Vec<SearchResult>>, VelesError> {
let query_refs: Vec<&[f32]> = searches.iter().map(|s| s.vector.as_slice()).collect();
let filters: Result<Vec<Option<Filter>>, VelesError> = searches
.iter()
.map(|s| {
s.filter
.as_ref()
.map(|f_json| {
serde_json::from_str(f_json).map_err(|e| {
VelesError::database(format!("Invalid filter JSON in batch: {e}"))
})
})
.transpose()
})
.collect();
let scope =
self.db
.authorize_read(&self.name, QueryOperationKind::VectorSearch, None, None)?;
let filters: Vec<Option<Filter>> = filters?
.into_iter()
.map(|f| and_scope(f, scope.clone()))
.collect();
let max_top_k = searches.iter().map(|s| s.top_k).max().unwrap_or(10);
let all_results = self.inner.search_batch_with_filters(
&query_refs,
usize::try_from(max_top_k).unwrap_or(usize::MAX),
&filters,
)?;
Ok(all_results
.into_iter()
.zip(searches)
.map(
|(results, s): (Vec<velesdb_core::SearchResult>, IndividualSearchRequest)| {
results
.into_iter()
.take(usize::try_from(s.top_k).unwrap_or(usize::MAX))
.map(|r| SearchResult {
id: r.point.id,
score: r.score,
payload: None,
})
.collect()
},
)
.collect())
}
pub fn text_search_with_filter(
&self,
query: String,
limit: u32,
filter_json: String,
) -> Result<Vec<SearchResult>, VelesError> {
let filter: Filter = serde_json::from_str(&filter_json)
.map_err(|e| VelesError::database(format!("Invalid filter JSON: {e}")))?;
let results = self.db.gated_search(
&self.name,
None,
None,
GatedRead::Text {
query: &query,
k: usize::try_from(limit).unwrap_or(usize::MAX),
filter: Some(&filter),
},
)?;
Ok(to_mobile_results(results))
}
pub fn hybrid_search_with_filter(
&self,
vector: Vec<f32>,
text_query: String,
limit: u32,
vector_weight: f32,
filter_json: String,
) -> Result<Vec<SearchResult>, VelesError> {
let filter: Filter = serde_json::from_str(&filter_json)
.map_err(|e| VelesError::database(format!("Invalid filter JSON: {e}")))?;
let results = self.db.gated_search(
&self.name,
None,
None,
GatedRead::Hybrid {
vector: &vector,
text: &text_query,
k: usize::try_from(limit).unwrap_or(usize::MAX),
alpha: Some(vector_weight),
filter: Some(&filter),
},
)?;
Ok(to_mobile_results(results))
}
pub fn query(
&self,
query_str: String,
params_json: Option<String>,
) -> Result<Vec<SearchResult>, VelesError> {
let parsed = velesdb_core::velesql::Parser::parse(&query_str)
.map_err(|e| VelesError::database(format!("VelesQL parse error: {}", e.message)))?;
let params: std::collections::HashMap<String, serde_json::Value> = params_json
.map(|json| serde_json::from_str(&json))
.transpose()
.map_err(|e| VelesError::database(format!("Invalid params JSON: {e}")))?
.unwrap_or_default();
let results = self
.db
.execute_query(&parsed, ¶ms)
.map_err(|e| VelesError::database(format!("Query execution failed: {e}")))?;
Ok(results
.into_iter()
.map(|r| SearchResult {
id: r.point.id,
score: r.score,
payload: r.point.payload.as_ref().map(|p| p.to_string()),
})
.collect())
}
pub fn enable_streaming(
&self,
config: Option<MobileStreamingConfig>,
) -> Result<(), VelesError> {
let core_config = config.map_or_else(velesdb_core::StreamingConfig::default, |c| {
velesdb_core::StreamingConfig::new(
usize::try_from(c.buffer_size).unwrap_or(usize::MAX),
usize::try_from(c.batch_size).unwrap_or(usize::MAX),
c.flush_interval_ms,
)
});
let rt = crate::streaming_runtime::stream_runtime()?;
let _guard = rt.enter();
self.inner.enable_streaming(core_config);
Ok(())
}
pub fn stream_insert(&self, points: Vec<VelesPoint>) -> Result<u64, VelesError> {
let core_points: Result<Vec<velesdb_core::Point>, VelesError> =
points.into_iter().map(parse_point).collect();
let queued = self.inner.stream_insert_batch(core_points?).map_err(|e| {
VelesError::database(format!(
"Stream insert failed (buffer full or not configured): {e}"
))
})?;
Ok(u64::try_from(queued).unwrap_or(u64::MAX))
}
pub fn flush(&self) -> Result<(), VelesError> {
self.inner.flush()?;
Ok(())
}
pub fn compact_storage(&self) -> Result<u64, VelesError> {
Ok(u64::try_from(self.inner.compact_storage()?).unwrap_or(u64::MAX))
}
pub fn guard_rails(&self) -> MobileQueryLimits {
self.inner.guard_rails().limits().into()
}
pub fn apply_advanced_config(&self, config: MobileAdvancedConfig) -> Result<(), VelesError> {
self.inner.apply_advanced_config(
config.pq_rescore_oversampling.map(Some),
config.deferred_indexing.map(|c| Some(c.into())),
config.async_index_builder.map(|c| Some(c.into())),
)?;
Ok(())
}
pub fn all_ids(&self) -> Vec<u64> {
self.inner.all_ids()
}
pub fn create_index(&self, field_name: String) -> Result<(), VelesError> {
self.inner.create_index(&field_name)?;
Ok(())
}
pub fn has_secondary_index(&self, field_name: String) -> bool {
self.inner.has_secondary_index(&field_name)
}
pub fn create_property_index(&self, label: String, property: String) -> Result<(), VelesError> {
self.inner.create_property_index(&label, &property)?;
Ok(())
}
pub fn create_range_index(&self, label: String, property: String) -> Result<(), VelesError> {
self.inner.create_range_index(&label, &property)?;
Ok(())
}
pub fn has_property_index(&self, label: String, property: String) -> bool {
self.inner.has_property_index(&label, &property)
}
pub fn has_range_index(&self, label: String, property: String) -> bool {
self.inner.has_range_index(&label, &property)
}
pub fn list_indexes(&self) -> Vec<MobileIndexInfo> {
self.inner
.list_indexes()
.into_iter()
.map(MobileIndexInfo::from)
.collect()
}
pub fn drop_index(&self, label: String, property: String) -> Result<bool, VelesError> {
Ok(self.inner.drop_index(&label, &property)?)
}
pub fn indexes_memory_usage(&self) -> u64 {
u64::try_from(self.inner.indexes_memory_usage()).unwrap_or(u64::MAX)
}
pub fn analyze(&self) -> Result<MobileCollectionStats, VelesError> {
Ok(self.inner.analyze()?.into())
}
pub fn get_stats(&self) -> MobileCollectionStats {
self.inner.get_stats().into()
}
pub fn diagnostics(&self) -> MobileCollectionDiagnostics {
self.inner.diagnostics().into()
}
}
fn parse_point(p: VelesPoint) -> Result<velesdb_core::Point, VelesError> {
let payload = p
.payload
.map(|s| serde_json::from_str(&s))
.transpose()
.map_err(|e| VelesError::database(format!("Invalid JSON payload: {e}")))?;
Ok(velesdb_core::Point::new(p.id, p.vector, payload))
}