1use std::sync::Arc;
4
5use velesdb_core::{
6 Database as CoreDatabase, Filter, GatedRead, QueryOperationKind,
7 VectorCollection as CoreCollection,
8};
9
10use crate::types::{
11 IndividualSearchRequest, MobileAdvancedConfig, MobileCollectionDiagnostics,
12 MobileCollectionStats, MobileIndexInfo, MobileQueryLimits, MobileStreamingConfig,
13 SearchQuality, SearchResult, VelesError, VelesPoint,
14};
15
16#[derive(uniffi::Object)]
30pub struct VelesCollection {
31 pub(crate) inner: CoreCollection,
32 pub(crate) db: Arc<CoreDatabase>,
34 pub(crate) name: String,
36}
37
38fn to_mobile_results(results: Vec<velesdb_core::SearchResult>) -> Vec<SearchResult> {
42 results
43 .into_iter()
44 .map(|r| SearchResult {
45 id: r.point.id,
46 score: r.score,
47 payload: None,
48 })
49 .collect()
50}
51
52pub(crate) fn and_scope(caller: Option<Filter>, scope: Option<Filter>) -> Option<Filter> {
55 use velesdb_core::Condition;
56 match (caller, scope) {
57 (None, None) => None,
58 (Some(c), None) => Some(c),
59 (None, Some(s)) => Some(s),
60 (Some(c), Some(s)) => Some(Filter::new(Condition::And {
61 conditions: vec![c.condition, s.condition],
62 })),
63 }
64}
65
66pub(crate) fn deny_if_scoped(scope: Option<Filter>, context: &str) -> Result<(), VelesError> {
71 if scope.is_some() {
72 return Err(VelesError::database(format!(
73 "{context} cannot honor the governance scope filter returned by the observer \
74 (this entry point has no metadata-filtered leaf); refusing to run unscoped"
75 )));
76 }
77 Ok(())
78}
79
80#[uniffi::export]
81impl VelesCollection {
82 pub fn search(&self, vector: Vec<f32>, limit: u32) -> Result<Vec<SearchResult>, VelesError> {
93 let results = self.db.gated_search(
94 &self.name,
95 None,
96 None,
97 GatedRead::Dense {
98 query: &vector,
99 k: usize::try_from(limit).unwrap_or(usize::MAX),
100 ef: None,
101 quality: None,
102 filter: None,
103 },
104 )?;
105
106 Ok(to_mobile_results(results))
107 }
108
109 pub fn search_with_quality(
121 &self,
122 vector: Vec<f32>,
123 limit: u32,
124 quality: SearchQuality,
125 ) -> Result<Vec<SearchResult>, VelesError> {
126 let results = self.db.gated_search(
127 &self.name,
128 None,
129 None,
130 GatedRead::Dense {
131 query: &vector,
132 k: usize::try_from(limit).unwrap_or(usize::MAX),
133 ef: None,
134 quality: Some(quality.into()),
135 filter: None,
136 },
137 )?;
138
139 Ok(to_mobile_results(results))
140 }
141
142 pub fn upsert(&self, point: VelesPoint) -> Result<(), VelesError> {
148 let core_point = parse_point(point)?;
149 self.inner.upsert(vec![core_point])?;
150 Ok(())
151 }
152
153 pub fn upsert_batch(&self, points: Vec<VelesPoint>) -> Result<(), VelesError> {
159 let core_points: Result<Vec<velesdb_core::Point>, VelesError> =
160 points.into_iter().map(parse_point).collect();
161
162 self.inner.upsert(core_points?)?;
163 Ok(())
164 }
165
166 pub fn delete(&self, id: u64) -> Result<(), VelesError> {
168 self.inner.delete(&[id])?;
169 Ok(())
170 }
171
172 #[allow(clippy::cast_possible_truncation)]
174 pub fn count(&self) -> u64 {
175 self.inner.config().point_count as u64
176 }
177
178 #[allow(clippy::cast_possible_truncation)]
180 pub fn dimension(&self) -> u32 {
181 self.inner.config().dimension as u32
182 }
183
184 pub fn get(&self, ids: Vec<u64>) -> Vec<VelesPoint> {
194 self.inner
195 .get(&ids)
196 .into_iter()
197 .flatten()
198 .map(|p| VelesPoint {
199 id: p.id,
200 vector: p.vector,
201 payload: p.payload.map(|v| v.to_string()),
202 })
203 .collect()
204 }
205
206 pub fn get_by_id(&self, id: u64) -> Option<VelesPoint> {
216 self.inner
217 .get(&[id])
218 .into_iter()
219 .flatten()
220 .next()
221 .map(|p| VelesPoint {
222 id: p.id,
223 vector: p.vector,
224 payload: p.payload.map(|v| v.to_string()),
225 })
226 }
227
228 pub fn is_metadata_only(&self) -> bool {
230 self.inner.config().metadata_only
231 }
232
233 pub fn text_search(&self, query: String, limit: u32) -> Result<Vec<SearchResult>, VelesError> {
244 let results = self.db.gated_search(
245 &self.name,
246 None,
247 None,
248 GatedRead::Text {
249 query: &query,
250 k: usize::try_from(limit).unwrap_or(usize::MAX),
251 filter: None,
252 },
253 )?;
254
255 Ok(to_mobile_results(results))
256 }
257
258 pub fn hybrid_search(
271 &self,
272 vector: Vec<f32>,
273 text_query: String,
274 limit: u32,
275 vector_weight: f32,
276 ) -> Result<Vec<SearchResult>, VelesError> {
277 let results = self.db.gated_search(
278 &self.name,
279 None,
280 None,
281 GatedRead::Hybrid {
282 vector: &vector,
283 text: &text_query,
284 k: usize::try_from(limit).unwrap_or(usize::MAX),
285 alpha: Some(vector_weight),
286 filter: None,
287 },
288 )?;
289
290 Ok(to_mobile_results(results))
291 }
292
293 pub fn search_with_filter(
305 &self,
306 vector: Vec<f32>,
307 limit: u32,
308 filter_json: String,
309 ) -> Result<Vec<SearchResult>, VelesError> {
310 let filter: Filter = serde_json::from_str(&filter_json)
312 .map_err(|e| VelesError::database(format!("Invalid filter JSON: {e}")))?;
313
314 let results = self.db.gated_search(
315 &self.name,
316 None,
317 None,
318 GatedRead::Dense {
319 query: &vector,
320 k: usize::try_from(limit).unwrap_or(usize::MAX),
321 ef: None,
322 quality: None,
323 filter: Some(&filter),
324 },
325 )?;
326
327 Ok(to_mobile_results(results))
328 }
329
330 pub fn batch_search(
340 &self,
341 searches: Vec<IndividualSearchRequest>,
342 ) -> Result<Vec<Vec<SearchResult>>, VelesError> {
343 let query_refs: Vec<&[f32]> = searches.iter().map(|s| s.vector.as_slice()).collect();
344
345 let filters: Result<Vec<Option<Filter>>, VelesError> = searches
346 .iter()
347 .map(|s| {
348 s.filter
349 .as_ref()
350 .map(|f_json| {
351 serde_json::from_str(f_json).map_err(|e| {
352 VelesError::database(format!("Invalid filter JSON in batch: {e}"))
353 })
354 })
355 .transpose()
356 })
357 .collect();
358
359 let scope =
363 self.db
364 .authorize_read(&self.name, QueryOperationKind::VectorSearch, None, None)?;
365 let filters: Vec<Option<Filter>> = filters?
366 .into_iter()
367 .map(|f| and_scope(f, scope.clone()))
368 .collect();
369 let max_top_k = searches.iter().map(|s| s.top_k).max().unwrap_or(10);
370
371 let all_results = self.inner.search_batch_with_filters(
372 &query_refs,
373 usize::try_from(max_top_k).unwrap_or(usize::MAX),
374 &filters,
375 )?;
376
377 Ok(all_results
378 .into_iter()
379 .zip(searches)
380 .map(
381 |(results, s): (Vec<velesdb_core::SearchResult>, IndividualSearchRequest)| {
382 results
383 .into_iter()
384 .take(usize::try_from(s.top_k).unwrap_or(usize::MAX))
385 .map(|r| SearchResult {
386 id: r.point.id,
387 score: r.score,
388 payload: None,
389 })
390 .collect()
391 },
392 )
393 .collect())
394 }
395
396 pub fn text_search_with_filter(
404 &self,
405 query: String,
406 limit: u32,
407 filter_json: String,
408 ) -> Result<Vec<SearchResult>, VelesError> {
409 let filter: Filter = serde_json::from_str(&filter_json)
410 .map_err(|e| VelesError::database(format!("Invalid filter JSON: {e}")))?;
411
412 let results = self.db.gated_search(
413 &self.name,
414 None,
415 None,
416 GatedRead::Text {
417 query: &query,
418 k: usize::try_from(limit).unwrap_or(usize::MAX),
419 filter: Some(&filter),
420 },
421 )?;
422
423 Ok(to_mobile_results(results))
424 }
425
426 pub fn hybrid_search_with_filter(
436 &self,
437 vector: Vec<f32>,
438 text_query: String,
439 limit: u32,
440 vector_weight: f32,
441 filter_json: String,
442 ) -> Result<Vec<SearchResult>, VelesError> {
443 let filter: Filter = serde_json::from_str(&filter_json)
444 .map_err(|e| VelesError::database(format!("Invalid filter JSON: {e}")))?;
445
446 let results = self.db.gated_search(
447 &self.name,
448 None,
449 None,
450 GatedRead::Hybrid {
451 vector: &vector,
452 text: &text_query,
453 k: usize::try_from(limit).unwrap_or(usize::MAX),
454 alpha: Some(vector_weight),
455 filter: Some(&filter),
456 },
457 )?;
458
459 Ok(to_mobile_results(results))
460 }
461
462 pub fn query(
482 &self,
483 query_str: String,
484 params_json: Option<String>,
485 ) -> Result<Vec<SearchResult>, VelesError> {
486 let parsed = velesdb_core::velesql::Parser::parse(&query_str)
488 .map_err(|e| VelesError::database(format!("VelesQL parse error: {}", e.message)))?;
489
490 let params: std::collections::HashMap<String, serde_json::Value> = params_json
492 .map(|json| serde_json::from_str(&json))
493 .transpose()
494 .map_err(|e| VelesError::database(format!("Invalid params JSON: {e}")))?
495 .unwrap_or_default();
496
497 let results = self
502 .db
503 .execute_query(&parsed, ¶ms)
504 .map_err(|e| VelesError::database(format!("Query execution failed: {e}")))?;
505
506 Ok(results
507 .into_iter()
508 .map(|r| SearchResult {
509 id: r.point.id,
510 score: r.score,
511 payload: r.point.payload.as_ref().map(|p| p.to_string()),
512 })
513 .collect())
514 }
515
516 pub fn enable_streaming(
529 &self,
530 config: Option<MobileStreamingConfig>,
531 ) -> Result<(), VelesError> {
532 let core_config = config.map_or_else(velesdb_core::StreamingConfig::default, |c| {
533 velesdb_core::StreamingConfig::new(
534 usize::try_from(c.buffer_size).unwrap_or(usize::MAX),
535 usize::try_from(c.batch_size).unwrap_or(usize::MAX),
536 c.flush_interval_ms,
537 )
538 });
539 let rt = crate::streaming_runtime::stream_runtime()?;
542 let _guard = rt.enter();
543 self.inner.enable_streaming(core_config);
544 Ok(())
545 }
546
547 pub fn stream_insert(&self, points: Vec<VelesPoint>) -> Result<u64, VelesError> {
556 let core_points: Result<Vec<velesdb_core::Point>, VelesError> =
557 points.into_iter().map(parse_point).collect();
558
559 let queued = self.inner.stream_insert_batch(core_points?).map_err(|e| {
560 VelesError::database(format!(
561 "Stream insert failed (buffer full or not configured): {e}"
562 ))
563 })?;
564 Ok(u64::try_from(queued).unwrap_or(u64::MAX))
565 }
566
567 pub fn flush(&self) -> Result<(), VelesError> {
569 self.inner.flush()?;
570 Ok(())
571 }
572
573 pub fn compact_storage(&self) -> Result<u64, VelesError> {
577 Ok(u64::try_from(self.inner.compact_storage()?).unwrap_or(u64::MAX))
578 }
579
580 pub fn guard_rails(&self) -> MobileQueryLimits {
582 self.inner.guard_rails().limits().into()
583 }
584
585 pub fn apply_advanced_config(&self, config: MobileAdvancedConfig) -> Result<(), VelesError> {
590 self.inner.apply_advanced_config(
591 config.pq_rescore_oversampling.map(Some),
592 config.deferred_indexing.map(|c| Some(c.into())),
593 config.async_index_builder.map(|c| Some(c.into())),
594 )?;
595 Ok(())
596 }
597
598 pub fn all_ids(&self) -> Vec<u64> {
600 self.inner.all_ids()
601 }
602
603 pub fn create_index(&self, field_name: String) -> Result<(), VelesError> {
605 self.inner.create_index(&field_name)?;
606 Ok(())
607 }
608
609 pub fn has_secondary_index(&self, field_name: String) -> bool {
611 self.inner.has_secondary_index(&field_name)
612 }
613
614 pub fn create_property_index(&self, label: String, property: String) -> Result<(), VelesError> {
616 self.inner.create_property_index(&label, &property)?;
617 Ok(())
618 }
619
620 pub fn create_range_index(&self, label: String, property: String) -> Result<(), VelesError> {
622 self.inner.create_range_index(&label, &property)?;
623 Ok(())
624 }
625
626 pub fn has_property_index(&self, label: String, property: String) -> bool {
628 self.inner.has_property_index(&label, &property)
629 }
630
631 pub fn has_range_index(&self, label: String, property: String) -> bool {
633 self.inner.has_range_index(&label, &property)
634 }
635
636 pub fn list_indexes(&self) -> Vec<MobileIndexInfo> {
638 self.inner
639 .list_indexes()
640 .into_iter()
641 .map(MobileIndexInfo::from)
642 .collect()
643 }
644
645 pub fn drop_index(&self, label: String, property: String) -> Result<bool, VelesError> {
647 Ok(self.inner.drop_index(&label, &property)?)
648 }
649
650 pub fn indexes_memory_usage(&self) -> u64 {
652 u64::try_from(self.inner.indexes_memory_usage()).unwrap_or(u64::MAX)
653 }
654
655 pub fn analyze(&self) -> Result<MobileCollectionStats, VelesError> {
657 Ok(self.inner.analyze()?.into())
658 }
659
660 pub fn get_stats(&self) -> MobileCollectionStats {
662 self.inner.get_stats().into()
663 }
664
665 pub fn diagnostics(&self) -> MobileCollectionDiagnostics {
667 self.inner.diagnostics().into()
668 }
669}
670
671fn parse_point(p: VelesPoint) -> Result<velesdb_core::Point, VelesError> {
673 let payload = p
674 .payload
675 .map(|s| serde_json::from_str(&s))
676 .transpose()
677 .map_err(|e| VelesError::database(format!("Invalid JSON payload: {e}")))?;
678 Ok(velesdb_core::Point::new(p.id, p.vector, payload))
679}
680
681