use super::deferred_embed::{RetryPolicy, embed_and_store};
use super::embedder::{shared_embedder, shared_embedder_initialized};
use super::handle::PalaceHandle;
use crate::memory_core::store::embed_ledger::{self, EmbedFailure};
use crate::memory_core::store::vector::{AliasScan, UnaliasOutcome};
use anyhow::{Context, Result};
use std::collections::HashSet;
use uuid::Uuid;
#[derive(Debug, Clone)]
pub enum AliasAudit {
Measured {
key_rows: usize,
distinct_vector_ids: usize,
aliased_drawer_ids: Vec<Uuid>,
unnameable_keys: Vec<String>,
},
Unavailable {
reason: String,
},
}
impl AliasAudit {
pub fn from_scan(scan: anyhow::Result<AliasScan>) -> Self {
match scan {
Ok(s) => Self::Measured {
key_rows: s.key_rows,
distinct_vector_ids: s.distinct_vector_ids,
aliased_drawer_ids: s.aliased_drawer_ids,
unnameable_keys: s.unnameable_keys,
},
Err(e) => Self::Unavailable {
reason: format!("{e:#}"),
},
}
}
pub fn is_clean(&self) -> bool {
matches!(
self,
Self::Measured {
key_rows,
distinct_vector_ids,
aliased_drawer_ids,
unnameable_keys,
} if aliased_drawer_ids.is_empty()
&& unnameable_keys.is_empty()
&& key_rows == distinct_vector_ids
)
}
pub fn unnameable_keys(&self) -> Option<&[String]> {
match self {
Self::Measured {
unnameable_keys, ..
} => Some(unnameable_keys),
Self::Unavailable { .. } => None,
}
}
pub fn aliased_drawer_ids(&self) -> Option<&[Uuid]> {
match self {
Self::Measured {
aliased_drawer_ids, ..
} => Some(aliased_drawer_ids),
Self::Unavailable { .. } => None,
}
}
pub fn counts(&self) -> Option<(usize, usize)> {
match self {
Self::Measured {
key_rows,
distinct_vector_ids,
..
} => Some((*key_rows, *distinct_vector_ids)),
Self::Unavailable { .. } => None,
}
}
pub fn unavailable_reason(&self) -> Option<&str> {
match self {
Self::Unavailable { reason } => Some(reason),
Self::Measured { .. } => None,
}
}
}
#[derive(Debug, Clone)]
pub struct EmbedHealth {
pub palace_id: String,
pub drawer_count: usize,
pub vector_count: usize,
pub missing_vector_ids: Vec<Uuid>,
pub recorded_failures: Vec<EmbedFailure>,
pub embedder_ready: bool,
pub alias_audit: AliasAudit,
}
impl EmbedHealth {
pub fn is_healthy(&self) -> bool {
self.missing_vector_ids.is_empty() && self.alias_audit.is_clean()
}
}
#[derive(Debug, Clone, Copy)]
pub struct VectorBackfillOptions {
pub dry_run: bool,
pub limit: Option<usize>,
pub retry: RetryPolicy,
}
impl Default for VectorBackfillOptions {
fn default() -> Self {
Self {
dry_run: true,
limit: None,
retry: RetryPolicy::default(),
}
}
}
#[derive(Debug, Clone)]
pub struct VectorBackfillReport {
pub palace_id: String,
pub dry_run: bool,
pub drawer_count: usize,
pub vector_count: usize,
pub missing: usize,
pub attempted: usize,
pub repaired: usize,
pub still_failing: usize,
pub still_missing_ids: Vec<Uuid>,
pub alias_audit: AliasAudit,
}
#[derive(Debug, Clone, Copy)]
pub struct AliasRepairOptions {
pub dry_run: bool,
}
impl Default for AliasRepairOptions {
fn default() -> Self {
Self { dry_run: true }
}
}
#[derive(Debug, Clone)]
pub enum AliasRepairOutcome {
Clean,
Planned,
Repaired,
Partial {
still_aliased: Vec<Uuid>,
not_freed: Vec<Uuid>,
unparsed_keys: Vec<String>,
},
Unavailable { reason: String },
}
impl AliasRepairOutcome {
pub fn as_str(&self) -> &'static str {
match self {
Self::Clean => "clean",
Self::Planned => "planned",
Self::Repaired => "repaired",
Self::Partial { .. } => "partial",
Self::Unavailable { .. } => "unavailable",
}
}
pub fn is_success(&self) -> bool {
matches!(self, Self::Clean | Self::Repaired)
}
}
#[derive(Debug, Clone)]
pub struct AliasRepairReport {
pub palace_id: String,
pub dry_run: bool,
pub before: AliasAudit,
pub freed_ids: Vec<Uuid>,
pub unnameable_keys: Vec<String>,
pub after: Option<AliasAudit>,
pub outcome: AliasRepairOutcome,
}
impl AliasRepairReport {
pub fn reembed_required(&self) -> bool {
!self.dry_run && !self.freed_ids.is_empty()
}
}
pub(crate) fn classify_repair(
after: &AliasAudit,
freed: &UnaliasOutcome,
expected: &[Uuid],
) -> AliasRepairOutcome {
let Some(still) = after.aliased_drawer_ids() else {
return AliasRepairOutcome::Unavailable {
reason: after
.unavailable_reason()
.unwrap_or("post-repair alias audit unavailable")
.to_string(),
};
};
let freed_set: HashSet<Uuid> = freed.freed.iter().copied().collect();
let not_freed: Vec<Uuid> = expected
.iter()
.copied()
.filter(|id| !freed_set.contains(id))
.collect();
let mut still_aliased = still.to_vec();
still_aliased.sort();
if still_aliased.is_empty()
&& not_freed.is_empty()
&& freed.unparsed_keys.is_empty()
&& after.is_clean()
{
return AliasRepairOutcome::Repaired;
}
AliasRepairOutcome::Partial {
still_aliased,
not_freed,
unparsed_keys: freed.unparsed_keys.clone(),
}
}
impl PalaceHandle {
pub fn repair_aliases(&self, opts: AliasRepairOptions) -> Result<AliasRepairReport> {
let before = AliasAudit::from_scan(self.vector_store.alias_audit());
let mut report = AliasRepairReport {
palace_id: self.id.as_str().to_string(),
dry_run: opts.dry_run,
before: before.clone(),
freed_ids: Vec::new(),
unnameable_keys: before.unnameable_keys().unwrap_or_default().to_vec(),
after: None,
outcome: AliasRepairOutcome::Clean,
};
let Some(aliased) = before.aliased_drawer_ids() else {
let reason = before
.unavailable_reason()
.unwrap_or("alias audit unavailable")
.to_string();
tracing::error!(
palace = %self.id,
"#5005: refusing to repair aliases — the audit could not run: {reason}"
);
report.outcome = AliasRepairOutcome::Unavailable { reason };
return Ok(report);
};
let mut expected: Vec<Uuid> = aliased.to_vec();
expected.sort();
if before.is_clean() {
return Ok(report);
}
if opts.dry_run {
report.freed_ids = expected;
report.outcome = AliasRepairOutcome::Planned;
return Ok(report);
}
if self.is_read_only() {
anyhow::bail!(
"palace '{}' is read-only: the HTTP daemon holds the write lock — run the \
alias repair through the daemon (`palace_unalias`) or stop it first",
self.id
);
}
let freed = self.vector_store.unalias().context("repair_aliases")?;
report.freed_ids = freed.freed.clone();
report.freed_ids.sort();
let after = AliasAudit::from_scan(self.vector_store.alias_audit());
report.after = Some(after.clone());
report.outcome = classify_repair(&after, &freed, &expected);
match &report.outcome {
AliasRepairOutcome::Repaired => tracing::warn!(
palace = %self.id, freed = report.freed_ids.len(),
"#5005: alias repair freed every aliased drawer and verified the palace \
clean; those drawers now need a re-embed"
),
other => tracing::error!(
palace = %self.id, freed = report.freed_ids.len(), outcome = other.as_str(),
"#5005: alias repair is INCOMPLETE — do not treat this palace as repaired"
),
}
Ok(report)
}
pub fn embed_health(&self) -> EmbedHealth {
let vector_ids: HashSet<Uuid> = self.vector_store.all_ids().into_iter().collect();
let now = chrono::Utc::now();
let live: Vec<Uuid> = self
.drawers
.read()
.iter()
.filter(|d| !d.is_expired_at(now) || d.is_tier_c())
.map(|d| d.id)
.collect();
let missing_vector_ids: Vec<Uuid> = live
.iter()
.copied()
.filter(|id| !vector_ids.contains(id))
.collect();
let recorded_failures = self
.data_dir
.as_ref()
.map(|d| embed_ledger::load(d))
.unwrap_or_default();
let alias_audit = AliasAudit::from_scan(self.vector_store.alias_audit());
if let Some(reason) = alias_audit.unavailable_reason() {
tracing::error!(palace = %self.id, "#5005: alias audit failed: {reason}");
}
EmbedHealth {
palace_id: self.id.as_str().to_string(),
drawer_count: live.len(),
vector_count: vector_ids.len(),
missing_vector_ids,
recorded_failures,
embedder_ready: shared_embedder_initialized(),
alias_audit,
}
}
pub async fn backfill_missing_vectors(
&self,
opts: VectorBackfillOptions,
) -> Result<VectorBackfillReport> {
let health = self.embed_health();
let mut report = VectorBackfillReport {
palace_id: health.palace_id.clone(),
dry_run: opts.dry_run,
drawer_count: health.drawer_count,
vector_count: health.vector_count,
missing: health.missing_vector_ids.len(),
attempted: 0,
repaired: 0,
still_failing: 0,
still_missing_ids: health.missing_vector_ids.clone(),
alias_audit: health.alias_audit.clone(),
};
if opts.dry_run || health.missing_vector_ids.is_empty() {
return Ok(report);
}
if self.is_read_only() {
anyhow::bail!(
"palace '{}' is read-only: the HTTP daemon holds the write lock — run the \
vector backfill through the daemon (`palace_reembed`) or stop it first",
self.id
);
}
let embedder = shared_embedder()
.await
.context("acquire shared embedder for vector backfill")?;
let targets: Vec<Uuid> = match opts.limit {
Some(n) => health.missing_vector_ids.iter().copied().take(n).collect(),
None => health.missing_vector_ids.clone(),
};
let mut repaired: HashSet<Uuid> = HashSet::new();
for id in targets {
let content = {
let drawers = self.drawers.read();
drawers
.iter()
.find(|d| d.id == id)
.map(|d| d.content.clone())
};
let Some(content) = content else {
continue;
};
report.attempted += 1;
match embed_and_store(&embedder, &self.vector_store, id, &content, &opts.retry).await {
Ok(attempts) => {
tracing::info!(
palace = %self.id, drawer = %id, attempts,
"#4906: vector backfill repaired a drawer"
);
repaired.insert(id);
}
Err(loss) => {
report.still_failing += 1;
tracing::error!(
palace = %self.id, drawer = %id,
"#4906: vector backfill could not repair this drawer: {}",
loss.reason()
);
self.record_backfill_failure(id, &loss).await;
}
}
}
report.repaired = repaired.len();
report.still_missing_ids.retain(|id| !repaired.contains(id));
if let Some(data_dir) = self.data_dir.as_ref()
&& !repaired.is_empty()
{
let dir = data_dir.clone();
let cleared =
tokio::task::spawn_blocking(move || embed_ledger::clear(&dir, &repaired)).await;
if let Ok(Err(e)) = cleared {
tracing::warn!(palace = %self.id, "#4906: ledger clear after backfill failed: {e:#}");
}
}
Ok(report)
}
async fn record_backfill_failure(&self, id: Uuid, loss: &super::deferred_embed::EmbedLoss) {
let Some(data_dir) = self.data_dir.as_ref() else {
return;
};
if !loss.is_drawer_specific() {
return;
}
let entry = EmbedFailure {
drawer_id: id,
failed_at: chrono::Utc::now(),
attempts: loss.attempts(),
reason: loss.reason(),
};
let dir = data_dir.clone();
let written = tokio::task::spawn_blocking(move || embed_ledger::record(&dir, entry)).await;
if let Ok(Err(e)) = written {
tracing::error!(palace = %self.id, drawer = %id, "#4906: ledger write failed: {e:#}");
}
}
}