use super::embedder::shared_embedder;
use crate::memory_core::embed::Embedder;
use crate::memory_core::palace::{Drawer, PalaceId};
use crate::memory_core::store::embed_ledger::{self, EmbedFailure};
use crate::memory_core::store::vector::{UsearchStore, VectorStore};
use crate::memory_core::timeouts;
use parking_lot::RwLock;
use std::path::PathBuf;
use std::sync::Arc;
use std::time::Duration;
use uuid::Uuid;
const DEFAULT_EMBED_ATTEMPTS: u32 = 3;
const DEFAULT_EMBED_RETRY_BASE_MS: u64 = 250;
#[derive(Debug, Clone, Copy)]
pub struct RetryPolicy {
pub attempts: u32,
pub base_delay: Duration,
}
impl Default for RetryPolicy {
fn default() -> Self {
Self {
attempts: DEFAULT_EMBED_ATTEMPTS,
base_delay: Duration::from_millis(DEFAULT_EMBED_RETRY_BASE_MS),
}
}
}
impl RetryPolicy {
fn delay_for(&self, attempt: u32) -> Duration {
let shift = attempt.saturating_sub(1).min(6);
self.base_delay.saturating_mul(1u32 << shift)
}
#[cfg(any(test, feature = "embedder-test-support"))]
pub fn instant(attempts: u32) -> Self {
Self {
attempts,
base_delay: Duration::ZERO,
}
}
}
#[derive(Debug, Clone)]
pub enum EmbedLoss {
EmbedderUnavailable(String),
Embed { reason: String, attempts: u32 },
Upsert { reason: String, attempts: u32 },
}
impl EmbedLoss {
pub fn is_drawer_specific(&self) -> bool {
!matches!(self, EmbedLoss::EmbedderUnavailable(_))
}
pub(super) fn attempts(&self) -> u32 {
match self {
EmbedLoss::EmbedderUnavailable(_) => 0,
EmbedLoss::Embed { attempts, .. } | EmbedLoss::Upsert { attempts, .. } => *attempts,
}
}
pub(super) fn reason(&self) -> String {
match self {
EmbedLoss::EmbedderUnavailable(r) => format!("embedder unavailable: {r}"),
EmbedLoss::Embed { reason, .. } => format!("embed failed: {reason}"),
EmbedLoss::Upsert { reason, .. } => format!("vector upsert failed: {reason}"),
}
}
}
pub(super) async fn with_retry<T, F, Fut>(
policy: &RetryPolicy,
mut op: F,
) -> Result<(T, u32), (String, u32)>
where
F: FnMut(u32) -> Fut,
Fut: std::future::Future<Output = Result<T, String>>,
{
let total = policy.attempts.max(1);
let mut last = String::from("no attempt was made");
for attempt in 1..=total {
match op(attempt).await {
Ok(v) => return Ok((v, attempt)),
Err(e) => last = e,
}
if attempt < total {
tokio::time::sleep(policy.delay_for(attempt)).await;
}
}
Err((last, total))
}
pub(super) async fn embed_and_store(
embedder: &Arc<dyn Embedder + Send + Sync>,
vector_store: &Arc<UsearchStore>,
id: Uuid,
content: &str,
policy: &RetryPolicy,
) -> Result<u32, EmbedLoss> {
let batch = [content.to_string()];
let embed_timeout = timeouts::embed_batch_timeout();
let (vector, embed_attempts) = with_retry(policy, |_| {
let embedder = embedder.clone();
let batch = &batch;
async move {
match tokio::time::timeout(embed_timeout, embedder.embed_batch(batch)).await {
Ok(Ok(v)) => v
.into_iter()
.next()
.ok_or_else(|| "embedder returned an empty batch".to_string()),
Ok(Err(e)) => Err(format!("{e:#}")),
Err(_) => Err(format!("embed_batch timed out after {embed_timeout:?}")),
}
}
})
.await
.map_err(|(reason, attempts)| EmbedLoss::Embed { reason, attempts })?;
let (_, upsert_attempts) = with_retry(policy, |_| {
let vector = vector.clone();
let vector_store = vector_store.clone();
async move {
vector_store
.upsert(id, vector)
.await
.map_err(|e| format!("{e:#}"))
}
})
.await
.map_err(|(reason, attempts)| EmbedLoss::Upsert { reason, attempts })?;
Ok(embed_attempts.max(upsert_attempts))
}
#[derive(Clone)]
pub(super) struct DeferredEmbedCtx {
pub palace_id: PalaceId,
pub vector_store: Arc<UsearchStore>,
pub drawers: Arc<RwLock<Vec<Drawer>>>,
pub data_dir: Option<PathBuf>,
}
pub(super) async fn run_deferred_embed(
ctx: DeferredEmbedCtx,
id: Uuid,
content: String,
policy: RetryPolicy,
) {
let embedder = match shared_embedder().await {
Ok(e) => e,
Err(e) => {
record_loss(&ctx, id, &EmbedLoss::EmbedderUnavailable(format!("{e:#}"))).await;
return;
}
};
embed_store_or_record(&ctx, &embedder, id, &content, &policy).await;
}
pub(super) async fn embed_store_or_record(
ctx: &DeferredEmbedCtx,
embedder: &Arc<dyn Embedder + Send + Sync>,
id: Uuid,
content: &str,
policy: &RetryPolicy,
) {
match embed_and_store(embedder, &ctx.vector_store, id, content, policy).await {
Ok(attempts) => {
tracing::info!(
palace = %ctx.palace_id, drawer = %id, attempts,
"deferred embed: vector backfill complete"
);
clear_ledger(ctx, id).await;
}
Err(loss) => record_loss(ctx, id, &loss).await,
}
}
pub(super) fn spawn(ctx: DeferredEmbedCtx, id: Uuid, content: String) {
tokio::spawn(run_deferred_embed(ctx, id, content, RetryPolicy::default()));
}
pub(super) async fn record_loss(ctx: &DeferredEmbedCtx, id: Uuid, loss: &EmbedLoss) {
let reason = loss.reason();
if !loss.is_drawer_specific() {
tracing::warn!(
palace = %ctx.palace_id, drawer = %id,
"#4906: deferred embed skipped — {reason}; drawer stored WITHOUT a \
vector and left unmarked (no embedder on this host). Run the \
vector backfill once an embedder is available."
);
return;
}
tracing::error!(
palace = %ctx.palace_id, drawer = %id, attempts = loss.attempts(),
"#4906: deferred embed FAILED — {reason}; drawer is durable but not \
vector-searchable. Recorded in the palace embed-failure ledger."
);
let still_present = ctx.drawers.read().iter().any(|d| d.id == id);
if !still_present {
return;
}
let Some(data_dir) = ctx.data_dir.as_ref() else {
tracing::warn!(
palace = %ctx.palace_id, drawer = %id,
"#4906: no data dir on this handle — failure could not be recorded durably"
);
return;
};
let entry = EmbedFailure {
drawer_id: id,
failed_at: chrono::Utc::now(),
attempts: loss.attempts(),
reason,
};
let dir = data_dir.clone();
let written = tokio::task::spawn_blocking(move || embed_ledger::record(&dir, entry)).await;
match written {
Ok(Ok(())) => {}
Ok(Err(e)) => tracing::error!(
palace = %ctx.palace_id, drawer = %id,
"#4906: could not write the embed-failure ledger: {e:#}"
),
Err(e) => tracing::error!(
palace = %ctx.palace_id, drawer = %id,
"#4906: embed-failure ledger write task failed: {e}"
),
}
}
async fn clear_ledger(ctx: &DeferredEmbedCtx, id: Uuid) {
let Some(data_dir) = ctx.data_dir.as_ref() else {
return;
};
let dir = data_dir.clone();
let cleared = tokio::task::spawn_blocking(move || {
let ids: std::collections::HashSet<Uuid> = std::iter::once(id).collect();
embed_ledger::clear(&dir, &ids)
})
.await;
if let Ok(Err(e)) = cleared {
tracing::warn!(palace = %ctx.palace_id, drawer = %id, "#4906: ledger clear failed: {e:#}");
}
}