use crate::backend::TantivyBackend;
use crate::indexer::mutations_from_step;
use crate::model::IndexDoc;
use anyhow::Result;
use mf_model::node_pool::NodePool;
use mf_transform::step::Step;
use std::sync::Arc;
use mf_state::transaction::Transaction;
#[derive(Debug, Clone)]
pub enum IndexEvent {
StepApplied {
pool_before: Option<Arc<NodePool>>,
pool_after: Arc<NodePool>,
step: Arc<dyn Step>,
},
TransactionCommitted {
pool_before: Option<Arc<NodePool>>,
pool_after: Arc<NodePool>,
steps: Vec<Arc<dyn Step>>,
},
Rebuild { pool: Arc<NodePool>, scope: RebuildScope },
}
#[derive(Debug, Clone, Copy)]
pub enum RebuildScope {
Full,
}
pub struct IndexService {
backend: Arc<TantivyBackend>,
}
impl IndexService {
pub fn new(backend: Arc<TantivyBackend>) -> Self {
Self { backend }
}
pub async fn handle(
&self,
event: IndexEvent,
) -> Result<()> {
match event {
IndexEvent::StepApplied { pool_before, pool_after, step } => {
let pool_b = pool_before.as_deref().unwrap_or(&pool_after);
let muts = mutations_from_step(pool_b, &pool_after, &step);
self.backend.apply(muts).await
},
IndexEvent::TransactionCommitted {
pool_before,
pool_after,
steps,
} => {
let pool_b = pool_before.as_deref().unwrap_or(&pool_after);
let mut all = Vec::new();
for s in &steps {
all.extend(mutations_from_step(pool_b, &pool_after, s));
}
self.backend.apply(all).await
},
IndexEvent::Rebuild { pool, scope: RebuildScope::Full } => {
let mut docs: Vec<IndexDoc> = Vec::new();
for shard in &pool.get_inner().nodes {
for node in shard.values() {
docs.push(IndexDoc::from_node(&pool, node));
}
}
self.backend.rebuild_all(docs).await
},
}
}
}
#[allow(dead_code)]
pub fn event_from_transaction(
pool_after: Arc<NodePool>,
tr: &Transaction,
) -> IndexEvent {
let steps: Vec<Arc<dyn Step>> = tr.steps.iter().cloned().collect();
IndexEvent::TransactionCommitted { pool_before: None, pool_after, steps }
}