use std::path::Path;
use std::sync::Arc;
use async_trait::async_trait;
use trusty_common::memory_core::palace::{Palace, PalaceId};
use trusty_common::memory_core::store::{KnowledgeGraph, OpenIntent};
use trusty_common::memory_core::{PalaceHandle, PalaceRegistry};
use trusty_common::palace_resolve::PalaceSource;
use super::apply::SnapshotView;
use super::discovery::{DiscoveredStore, STORE_DB_NAME};
use super::KuzuImportError;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Target {
pub palace: String,
pub source: &'static str,
}
pub fn target_palace(
store: &DiscoveredStore,
explicit: Option<&str>,
) -> Result<Target, KuzuImportError> {
if let Some(p) = explicit {
if trusty_common::palace_id::palace_id_is_valid(p) {
return Ok(Target {
palace: p.to_string(),
source: "--palace",
});
}
return Err(KuzuImportError::PalaceResolve(format!(
"invalid palace id {p:?}"
)));
}
let r = trusty_common::palace_resolve::resolve_palace(store.project_dir())
.map_err(|e| KuzuImportError::PalaceResolve(e.to_string()))?;
let source = match r.source {
PalaceSource::EnvOverride => "TRUSTY_MEMORY_PALACE",
PalaceSource::PinFile => "pin file",
PalaceSource::GitOwnerRepo => "git owner/repo",
PalaceSource::ParentDir => "parent/dir",
_ => "derived",
};
Ok(Target {
palace: r.id,
source,
})
}
pub fn check_env_override(walk: bool, env_palace: Option<&str>) -> Result<(), KuzuImportError> {
match env_palace {
Some(v) if walk => Err(KuzuImportError::EnvPalaceWithWalk(v.to_string())),
_ => Ok(()),
}
}
pub fn store_is_live(dir: &str) -> bool {
Path::new(dir).join(STORE_DB_NAME).exists()
}
pub fn open_palace_for_write(
data_root: &Path,
palace: &str,
) -> Result<Arc<PalaceHandle>, KuzuImportError> {
let registry = PalaceRegistry::new();
let id = PalaceId::new(palace);
let handle = match registry.open_palace(data_root, &id) {
Ok(h) => h,
Err(e) if PalaceRegistry::open_error_is_absent(&e) => registry
.create_palace(
data_root,
Palace {
id: id.clone(),
name: palace.to_string(),
description: Some("Imported from kuzu-memory".to_string()),
created_at: chrono::Utc::now(),
data_dir: data_root.join(palace),
},
)
.map_err(|e| KuzuImportError::Palace(format!("{e:#}")))?,
Err(e) => return Err(KuzuImportError::Palace(format!("{e:#}"))),
};
if handle.is_read_only() {
return Err(KuzuImportError::PalaceLocked(palace.to_string()));
}
if handle.drawer_load_degraded {
return Err(KuzuImportError::DrawersUnreadable(palace.to_string()));
}
Ok(handle)
}
pub fn open_snapshot_view(data_root: &Path, palace: &str) -> Result<SnapshotView, KuzuImportError> {
let live = data_root.join(palace).join("kg.redb");
if !live.exists() {
return Ok(SnapshotView {
drawers: Vec::new(),
kg: None,
snapshot_dir: None,
});
}
let tmp = tempfile::TempDir::with_prefix("trusty-kuzu-dry-run-")
.map_err(|e| KuzuImportError::Io(format!("create snapshot dir: {e}")))?;
std::fs::copy(&live, tmp.path().join("kg.redb"))
.map_err(|e| KuzuImportError::Io(format!("copy kg.redb for a dry run: {e}")))?;
let err = |e: anyhow::Error| KuzuImportError::Palace(format!("{e:#}"));
let kg =
KnowledgeGraph::open_with_intent(&tmp.path().join("kg.db"), OpenIntent::ReadOnlyClient)
.map_err(err)?;
let (drawers, skipped) = kg.load_drawers_with_skipped().map_err(err)?;
if skipped > 0 {
return Err(KuzuImportError::DrawersUnreadable(palace.to_string()));
}
Ok(SnapshotView {
drawers,
kg: Some(kg),
snapshot_dir: Some(tmp),
})
}
#[async_trait]
pub trait DaemonProbe: Send + Sync {
async fn live_daemon(&self) -> Option<String>;
}
pub struct SystemDaemonProbe;
#[async_trait]
impl DaemonProbe for SystemDaemonProbe {
async fn live_daemon(&self) -> Option<String> {
if let Ok(socket) = crate::socket_path() {
if crate::commands::daemon_guard::probe(&socket).await {
return Some(format!("serving {}", socket.display()));
}
}
let pids = crate::commands::stop::find_daemon_pids();
(!pids.is_empty()).then(|| format!("pid {pids:?}"))
}
}