pub mod corpus;
pub mod registry;
use std::collections::{HashMap, HashSet};
use std::fs::File;
use std::future::Future;
use std::ops::ControlFlow;
use std::path::{Path, PathBuf};
use std::time::Duration;
use tokio_util::sync::CancellationToken;
use uuid::Uuid;
#[doc(hidden)]
pub fn merge_fresh_tail<S, E>(
candidates: Vec<(Uuid, S)>,
ops: Vec<(Uuid, Option<Vec<f32>>)>,
mut score: impl FnMut(&[f32]) -> Result<S, E>,
) -> Result<Vec<(Uuid, S)>, E>
where
S: PartialOrd,
{
if ops.is_empty() {
return Ok(candidates);
}
let mut deletes: HashSet<Uuid> = HashSet::new();
let mut upserts: HashMap<Uuid, S> = HashMap::new();
for (uuid, op) in ops {
match op {
None => {
deletes.insert(uuid);
}
Some(embedding) => {
upserts.insert(uuid, score(&embedding)?);
}
}
}
let mut merged: Vec<(Uuid, S)> = candidates
.into_iter()
.filter(|(uuid, _)| !deletes.contains(uuid) && !upserts.contains_key(uuid))
.collect();
merged.extend(upserts);
merged.sort_by(|a, b| {
b.1.partial_cmp(&a.1)
.unwrap_or(std::cmp::Ordering::Equal)
.then_with(|| a.0.cmp(&b.0))
});
Ok(merged)
}
const CHECKPOINT_LOCK_FILE: &str = ".bridge-checkpoint.lock";
#[doc(hidden)]
pub fn acquire_checkpoint_lock(dir: &Path, prefix: &str) -> Result<File, String> {
std::fs::create_dir_all(dir).map_err(|error| {
format!(
"create {prefix} checkpoint directory {}: {error}",
dir.display()
)
})?;
let lock_path = dir.join(CHECKPOINT_LOCK_FILE);
let lock = std::fs::OpenOptions::new()
.create(true)
.truncate(false)
.read(true)
.write(true)
.open(&lock_path)
.map_err(|error| format!("open {prefix} lock {}: {error}", lock_path.display()))?;
lock.lock()
.map_err(|error| format!("acquire {prefix} lock {}: {error}", lock_path.display()))?;
Ok(lock)
}
#[doc(hidden)]
pub async fn acquire_checkpoint_lock_async(
dir: PathBuf,
prefix: &'static str,
) -> Result<File, String> {
tokio::task::spawn_blocking(move || acquire_checkpoint_lock(&dir, prefix))
.await
.map_err(|error| format!("{prefix} lock task failed: {error}"))?
}
#[doc(hidden)]
pub fn rotation_watch_future<T, F, Fut>(
ann: &std::sync::Arc<T>,
started: &std::sync::atomic::AtomicBool,
ann_root: PathBuf,
interval: Duration,
shutdown: CancellationToken,
mut refresh: F,
) -> Option<impl Future<Output = ()> + Send + 'static>
where
T: Send + Sync + 'static,
F: FnMut(std::sync::Arc<T>, PathBuf) -> Fut + Send + 'static,
Fut: Future<Output = ()> + Send + 'static,
{
use std::sync::atomic::Ordering;
if started
.compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
.is_err()
{
return None;
}
let ann = std::sync::Arc::downgrade(ann);
let tick = move || {
let ann = ann.upgrade();
let ann_root = ann_root.clone();
let refresh = ann.map(|ann| refresh(ann, ann_root));
async move {
let Some(refresh) = refresh else {
return ControlFlow::Break(());
};
refresh.await;
ControlFlow::Continue(())
}
};
Some(rotation_watch_loop(interval, shutdown, tick))
}
#[doc(hidden)]
pub async fn rotation_watch_loop<F, Fut>(
interval: Duration,
shutdown: CancellationToken,
mut tick: F,
) where
F: FnMut() -> Fut,
Fut: Future<Output = ControlFlow<()>>,
{
let start = tokio::time::Instant::now() + interval;
let mut ticks = tokio::time::interval_at(start, interval);
ticks.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
loop {
tokio::select! {
_ = shutdown.cancelled() => break,
_ = ticks.tick() => {}
}
if tick().await.is_break() {
break;
}
}
}