use serde_json::{json, Map, Value};
use crate::error::{Error, Result};
use crate::vector::{Document, MetadataFilter, SearchResult, VectorStore};
#[derive(Debug, Clone)]
pub struct PineconeStore {
http: reqwest::Client,
host: String,
api_key: String,
namespace: String,
}
impl PineconeStore {
pub fn new(host: impl Into<String>, api_key: impl Into<String>) -> Self {
Self {
http: reqwest::Client::new(),
host: host.into().trim_end_matches('/').to_string(),
api_key: api_key.into(),
namespace: String::new(),
}
}
#[must_use]
pub fn with_namespace(mut self, namespace: impl Into<String>) -> Self {
self.namespace = namespace.into();
self
}
async fn post(&self, path: &str, body: Value) -> Result<Value> {
let response = self
.http
.post(format!("{}{path}", self.host))
.header("Api-Key", &self.api_key)
.json(&body)
.send()
.await?;
let status = response.status();
let body: Value = response.json().await.unwrap_or(Value::Null);
if !status.is_success() {
return Err(Error::VectorStore(format!(
"Pinecone returned {status}: {body}"
)));
}
Ok(body)
}
async fn query(
&self,
vector: Vec<f32>,
top_k: usize,
filter: Option<Value>,
) -> Result<Vec<SearchResult>> {
let mut body = json!({
"vector": vector,
"topK": top_k,
"includeMetadata": true,
"namespace": self.namespace,
});
if let Some(filter) = filter {
body["filter"] = filter;
}
let response = self.post("/query", body).await?;
let matches = response
.get("matches")
.and_then(Value::as_array)
.cloned()
.unwrap_or_default();
Ok(matches
.into_iter()
.map(|hit| {
let mut metadata = hit.get("metadata").cloned().unwrap_or(Value::Null);
let text = metadata
.get("_text")
.and_then(Value::as_str)
.map(String::from);
if let Value::Object(map) = &mut metadata {
map.remove("_text");
}
SearchResult {
id: hit
.get("id")
.and_then(Value::as_str)
.unwrap_or_default()
.to_string(),
score: hit.get("score").and_then(Value::as_f64).unwrap_or(0.0) as f32,
text,
metadata,
}
})
.collect())
}
}
fn to_pinecone_metadata(doc: &Document) -> Value {
let mut metadata = Map::new();
if let Some(text) = &doc.text {
metadata.insert("_text".into(), Value::String(text.clone()));
}
match &doc.metadata {
Value::Null => {}
Value::Object(map) => {
for (key, value) in map {
match value {
Value::String(_) | Value::Number(_) | Value::Bool(_) => {
metadata.insert(key.clone(), value.clone());
}
other => {
metadata.insert(key.clone(), Value::String(other.to_string()));
}
}
}
}
other => {
metadata.insert("_metadata".into(), Value::String(other.to_string()));
}
}
Value::Object(metadata)
}
#[async_trait::async_trait]
impl VectorStore for PineconeStore {
async fn upsert(&self, documents: Vec<Document>) -> Result<()> {
let vectors: Vec<Value> = documents
.iter()
.map(|doc| {
json!({
"id": doc.id,
"values": doc.vector,
"metadata": to_pinecone_metadata(doc),
})
})
.collect();
self.post(
"/vectors/upsert",
json!({ "vectors": vectors, "namespace": self.namespace }),
)
.await?;
Ok(())
}
async fn search(&self, vector: Vec<f32>, top_k: usize) -> Result<Vec<SearchResult>> {
self.query(vector, top_k, None).await
}
async fn search_filtered(
&self,
vector: Vec<f32>,
top_k: usize,
filter: &MetadataFilter,
) -> Result<Vec<SearchResult>> {
if filter.is_empty() {
return self.query(vector, top_k, None).await;
}
let conditions: serde_json::Map<String, Value> = filter
.equals
.iter()
.map(|(key, value)| {
let comparable = match value {
Value::String(_) | Value::Number(_) | Value::Bool(_) => value.clone(),
other => Value::String(other.to_string()),
};
(key.clone(), json!({ "$eq": comparable }))
})
.collect();
self.query(vector, top_k, Some(Value::Object(conditions)))
.await
}
async fn delete(&self, ids: &[String]) -> Result<()> {
self.post(
"/vectors/delete",
json!({ "ids": ids, "namespace": self.namespace }),
)
.await?;
Ok(())
}
}