use crate::{
CollectionHandle, Error, JsonObject, Point, PointId, Result, ScoredPoint, Store, WriteResult,
};
use serde::Serialize;
use serde_json::Value;
#[cfg(feature = "fastembed")]
use std::sync::Mutex;
#[cfg(feature = "fastembed")]
pub use fastembed::{EmbeddingModel as FastEmbedModel, TextInitOptions as FastEmbedInitOptions};
pub trait Embedder {
fn model_id(&self) -> &str;
fn embed(&self, input: &[String]) -> Result<Vec<Vec<f32>>>;
}
#[cfg(feature = "fastembed")]
pub struct FastEmbedder {
model: Mutex<fastembed::TextEmbedding>,
model_id: String,
}
#[cfg(feature = "fastembed")]
impl FastEmbedder {
pub fn try_new() -> Result<Self> {
Self::try_with_model(FastEmbedModel::default())
}
pub fn try_with_model(model: FastEmbedModel) -> Result<Self> {
Self::try_from_options(FastEmbedInitOptions::new(model))
}
pub fn try_from_options(options: FastEmbedInitOptions) -> Result<Self> {
let model_id = format!("fastembed/{}@5.17.3", options.model_name);
let model = fastembed::TextEmbedding::try_new(options)
.map_err(|error| Error::Embedding(error.to_string()))?;
Ok(Self {
model: Mutex::new(model),
model_id,
})
}
}
#[cfg(feature = "fastembed")]
impl Embedder for FastEmbedder {
fn model_id(&self) -> &str {
&self.model_id
}
fn embed(&self, input: &[String]) -> Result<Vec<Vec<f32>>> {
self.model
.lock()
.map_err(|_| Error::Embedding("FastEmbed model lock was poisoned".into()))?
.embed(input, None)
.map_err(|error| Error::Embedding(error.to_string()))
}
}
#[derive(Clone, Debug, PartialEq)]
pub struct Document {
pub id: PointId,
pub text: String,
pub metadata: JsonObject,
}
impl Document {
pub fn new(id: impl Into<PointId>, text: impl Into<String>) -> Self {
Self {
id: id.into(),
text: text.into(),
metadata: JsonObject::new(),
}
}
pub fn with_metadata(mut self, metadata: impl Serialize) -> Result<Self> {
match serde_json::to_value(metadata)? {
Value::Object(metadata) => {
self.metadata = metadata;
Ok(self)
}
_ => Err(Error::Invalid(
"document metadata must serialize to a JSON object".into(),
)),
}
}
}
#[derive(Clone, Debug)]
pub struct TextCollection<E> {
collection: CollectionHandle,
embedder: E,
model_id: String,
}
impl Store {
pub fn text_collection<E: Embedder>(
&self,
name: impl Into<String>,
embedder: E,
) -> Result<TextCollection<E>> {
TextCollection::new(self.collection(name), embedder)
}
}
impl<E: Embedder> TextCollection<E> {
fn new(collection: CollectionHandle, embedder: E) -> Result<Self> {
let model_id = embedder.model_id().trim().to_owned();
if model_id.is_empty() {
return Err(Error::Invalid(
"embedding model identity must not be empty".into(),
));
}
if let Ok(existing) = collection.advanced() {
let actual = existing.info()?.config.vector_space;
if actual.as_deref() != Some(model_id.as_str()) {
return Err(Error::Invalid(format!(
"collection uses vector space {actual:?}, expected {model_id:?}"
)));
}
}
Ok(Self {
collection,
embedder,
model_id,
})
}
pub fn upsert_documents(
&self,
documents: impl IntoIterator<Item = Document>,
) -> Result<WriteResult> {
let documents: Vec<Document> = documents.into_iter().collect();
if documents.is_empty() {
return Err(Error::Invalid("document batch must not be empty".into()));
}
let input: Vec<String> = documents
.iter()
.map(|document| document.text.clone())
.collect();
let vectors = self.embedder.embed(&input)?;
if vectors.len() != documents.len() {
return Err(Error::Invalid(format!(
"embedder returned {} vectors for {} documents",
vectors.len(),
documents.len()
)));
}
let points = documents
.into_iter()
.zip(vectors)
.map(|(document, vector)| {
let mut payload = document.metadata;
if payload
.insert("document".into(), Value::String(document.text))
.is_some()
{
return Err(Error::Invalid(
"document metadata reserves the key \"document\"".into(),
));
}
Ok(Point {
id: document.id,
vector,
payload,
})
})
.collect::<Result<Vec<_>>>()?;
self.collection
.upsert_with_vector_space(points, Some(&self.model_id))
}
pub fn search_text(&self, text: impl Into<String>, limit: usize) -> Result<Vec<ScoredPoint>> {
let vectors = self.embedder.embed(&[text.into()])?;
let mut vectors = vectors.into_iter();
let vector = vectors
.next()
.ok_or_else(|| Error::Invalid("embedder returned no query vector".into()))?;
if vectors.next().is_some() {
return Err(Error::Invalid(
"embedder returned multiple vectors for one query".into(),
));
}
self.collection.search(vector, limit)
}
pub fn vectors(&self) -> &CollectionHandle {
&self.collection
}
}