use anyhow::Result;
use crate::index::Embedder;
use crate::storage::{MetadataStore, VectorStore};
use crate::types::Chunk;
pub trait ChunkSource {
fn name(&self) -> &str;
fn repo_key(&self) -> &str;
fn source_label(&self) -> &str;
fn fetch(&self) -> impl std::future::Future<Output = Result<Vec<Chunk>>> + Send;
}
pub struct SourceIndexReport {
pub indexed: usize,
pub unchanged: usize,
pub removed: usize,
}
pub async fn index_hashed_source<S: ChunkSource>(
source: &S,
vector_store: &mut VectorStore,
metadata_store: &MetadataStore,
embedder: &Embedder,
force: bool,
) -> Result<SourceIndexReport> {
let chunks = source.fetch().await?;
let repo_key = source.repo_key();
let now = chrono::Utc::now().timestamp().to_string();
if force {
metadata_store.clear_file_hashes(repo_key).ok();
}
let mut to_index = Vec::new();
let mut current_paths: std::collections::HashSet<String> = std::collections::HashSet::new();
let mut new_hashes: Vec<(String, String)> = Vec::new();
for c in &chunks {
current_paths.insert(c.file_path.clone());
let h = content_hash(&c.content);
let unchanged = metadata_store
.get_file_hash(repo_key, &c.file_path)
.ok()
.flatten()
.as_deref()
== Some(h.as_str());
if !unchanged {
to_index.push(c.clone());
}
new_hashes.push((c.file_path.clone(), h));
}
let removed: Vec<String> = metadata_store
.get_all_indexed_files(repo_key)
.unwrap_or_default()
.into_iter()
.filter(|p| !current_paths.contains(p))
.collect();
if !removed.is_empty() {
vector_store
.delete_by_file(&removed, Some(repo_key))
.await
.ok();
metadata_store
.delete_file_hashes(Some(repo_key), &removed)
.ok();
}
if to_index.is_empty() {
return Ok(SourceIndexReport {
indexed: 0,
unchanged: chunks.len(),
removed: removed.len(),
});
}
let embed_texts: Vec<String> = to_index.iter().map(|c| c.content.clone()).collect();
let embed_refs: Vec<&str> = embed_texts.iter().map(|s| s.as_str()).collect();
let embeddings = embedder.embed_batch(&embed_refs).await?;
let contexts = vec![None; to_index.len()];
let ids: Vec<String> = to_index.iter().map(|c| c.id.clone()).collect();
vector_store.delete(&ids).await.ok();
vector_store
.insert(
&to_index,
&embeddings,
&contexts,
repo_key,
source.source_label(),
&now,
)
.await?;
let hash_refs: Vec<(&str, &str)> = new_hashes
.iter()
.map(|(p, h)| (p.as_str(), h.as_str()))
.collect();
metadata_store
.set_file_hashes_bulk(repo_key, &hash_refs)
.ok();
Ok(SourceIndexReport {
indexed: to_index.len(),
unchanged: chunks.len() - to_index.len(),
removed: removed.len(),
})
}
pub trait WatermarkSource {
fn name(&self) -> &str;
fn repo(&self) -> &str;
fn source_label(&self) -> &str;
fn watermark_key(&self) -> String;
fn fetch_since(&self, since: Option<&str>) -> Result<(Vec<Chunk>, Option<String>)>;
}
pub async fn index_watermark_source<S: WatermarkSource>(
source: &S,
vector_store: &mut VectorStore,
metadata_store: &MetadataStore,
embedder: &Embedder,
force: bool,
) -> Result<usize> {
let stored = metadata_store.get_meta(&source.watermark_key())?;
let since = if force { None } else { stored.as_deref() };
let (chunks, new_watermark) = source.fetch_since(since)?;
if chunks.is_empty() {
return Ok(0);
}
let embed_texts: Vec<String> = chunks.iter().map(|c| c.content.clone()).collect();
let embed_refs: Vec<&str> = embed_texts.iter().map(|s| s.as_str()).collect();
let embeddings = embedder.embed_batch(&embed_refs).await?;
let contexts = vec![None; chunks.len()];
let now = chrono::Utc::now().timestamp().to_string();
if force {
let ids: Vec<String> = chunks.iter().map(|c| c.id.clone()).collect();
vector_store.delete(&ids).await.ok();
}
vector_store
.insert(
&chunks,
&embeddings,
&contexts,
source.repo(),
source.source_label(),
&now,
)
.await?;
if let Some(watermark) = new_watermark {
metadata_store.set_meta(&source.watermark_key(), &watermark)?;
}
Ok(chunks.len())
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ItemOutcome {
Indexed,
Unchanged,
Removed,
Absent,
}
pub async fn index_hashed_item(
repo_key: &str,
source_label: &str,
key: &str,
chunk: Option<&Chunk>,
vector_store: &mut VectorStore,
metadata_store: &MetadataStore,
embedder: &Embedder,
force: bool,
) -> Result<ItemOutcome> {
let Some(chunk) = chunk else {
let known = metadata_store
.get_file_hash(repo_key, key)
.ok()
.flatten()
.is_some();
if !known {
return Ok(ItemOutcome::Absent);
}
let keys = [key.to_string()];
vector_store
.delete_by_file(&keys, Some(repo_key))
.await
.ok();
metadata_store
.delete_file_hashes(Some(repo_key), &keys)
.ok();
return Ok(ItemOutcome::Removed);
};
let hash = content_hash(&chunk.content);
if !force {
let stored = metadata_store.get_file_hash(repo_key, key).ok().flatten();
if stored.as_deref() == Some(hash.as_str()) {
return Ok(ItemOutcome::Unchanged);
}
}
let embeddings = embedder.embed_batch(&[chunk.content.as_str()]).await?;
let now = chrono::Utc::now().timestamp().to_string();
let items = [chunk.clone()];
vector_store.delete(&[chunk.id.clone()]).await.ok();
vector_store
.insert(
&items,
&embeddings,
&vec![None; 1],
repo_key,
source_label,
&now,
)
.await?;
metadata_store
.set_file_hashes_bulk(repo_key, &[(key, hash.as_str())])
.ok();
Ok(ItemOutcome::Indexed)
}
pub fn content_hash(content: &str) -> String {
use sha2::{Digest, Sha256};
hex::encode(Sha256::digest(content.as_bytes()))
}