use std::collections::{BTreeMap, BTreeSet};
use std::path::{Path, PathBuf};
use std::sync::{Arc, RwLock};
use evorule_reactor::{Fact, FactId, FactsLog, FactsLogError};
use evorule_tcb::JsonValue;
use serde::{Deserialize, Serialize};
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SharedFact {
pub fact_id: FactId,
pub path: String,
pub value: JsonValue,
pub source_session_id: u64,
pub version: u64,
}
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
pub struct SharedFactsMetadata {
pub next_fact_id: u64,
pub fact_sources: BTreeMap<u64, u64>,
pub used_at_startup: BTreeMap<u64, Vec<u64>>,
pub rolled_up: BTreeSet<u64>,
}
#[derive(Debug, Clone)]
pub struct SharedFactsLog {
inner: Arc<RwLock<SharedFactsLogInner>>,
}
struct SharedFactsLogInner {
facts_log: FactsLog,
next_fact_id: u64,
fact_sources: BTreeMap<FactId, u64>,
used_at_startup: BTreeMap<u64, Vec<FactId>>,
rolled_up: BTreeSet<FactId>,
metadata_path: Option<PathBuf>,
}
impl std::fmt::Debug for SharedFactsLogInner {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("SharedFactsLogInner")
.field("facts_log", &self.facts_log)
.field("next_fact_id", &self.next_fact_id)
.field("fact_sources_len", &self.fact_sources.len())
.finish()
}
}
impl SharedFactsLogInner {
fn persist_metadata_locked(&self) {
if let Some(ref meta_path) = self.metadata_path {
if let Err(e) = write_metadata_atomic(meta_path, self) {
tracing::warn!("SharedFactsLog metadata persist failed: {e}");
}
}
}
}
impl SharedFactsLog {
pub fn new() -> Self {
Self {
inner: Arc::new(RwLock::new(SharedFactsLogInner {
facts_log: FactsLog::new(),
next_fact_id: 1,
fact_sources: BTreeMap::new(),
used_at_startup: BTreeMap::new(),
rolled_up: BTreeSet::new(),
metadata_path: None,
})),
}
}
pub fn with_wal<P: AsRef<Path>>(path: P) -> Result<Self, FactsLogError> {
let facts_log = FactsLog::with_wal(path)?;
Ok(Self {
inner: Arc::new(RwLock::new(SharedFactsLogInner {
facts_log,
next_fact_id: 1,
fact_sources: BTreeMap::new(),
used_at_startup: BTreeMap::new(),
rolled_up: BTreeSet::new(),
metadata_path: None,
})),
})
}
pub fn recover<P1: AsRef<Path>, P2: AsRef<Path>>(
wal_path: P1,
metadata_path: P2,
) -> Result<Self, FactsLogError> {
let facts_log = if wal_path.as_ref().exists() {
FactsLog::recover(&wal_path)?
} else {
FactsLog::with_wal(&wal_path)?
};
let tmp_path = metadata_path.as_ref().with_extension("tmp");
if tmp_path.exists() {
match std::fs::remove_file(&tmp_path) {
Ok(()) => tracing::info!(tmp = %tmp_path.display(), "已清理孤儿 metadata.tmp"),
Err(e) => tracing::warn!(tmp = %tmp_path.display(), error = %e, "清理孤儿 metadata.tmp 失败"),
}
}
let metadata = match std::fs::read_to_string(&metadata_path) {
Ok(content) => serde_json::from_str::<SharedFactsMetadata>(&content)
.map_err(|e| FactsLogError::WalError(format!("metadata parse error: {e}")))?,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => SharedFactsMetadata::default(),
Err(e) => return Err(FactsLogError::WalError(format!("metadata read error: {e}"))),
};
let fact_sources: BTreeMap<FactId, u64> = metadata
.fact_sources
.into_iter()
.map(|(k, v)| (FactId(k), v))
.collect();
let used_at_startup: BTreeMap<u64, Vec<FactId>> = metadata
.used_at_startup
.into_iter()
.map(|(k, v)| (k, v.into_iter().map(FactId).collect()))
.collect();
let rolled_up: BTreeSet<FactId> = metadata.rolled_up.into_iter().map(FactId).collect();
let orphans = detect_orphan_facts(&facts_log, &fact_sources);
if !orphans.is_empty() {
for id in &orphans {
tracing::warn!(
fact_id = id,
"孤立共享 fact:WAL 中存在但 fact_sources 缺失(崩溃窗口产物,因果映射永久丢失,该段审计链不可归因)"
);
}
tracing::warn!(
count = orphans.len(),
"孤立共享 fact 检测汇总:恢复的审计链中存在不可归因段,建议人工核对 WAL 与 metadata 文件"
);
}
Ok(Self {
inner: Arc::new(RwLock::new(SharedFactsLogInner {
facts_log,
next_fact_id: metadata.next_fact_id.max(1),
fact_sources,
used_at_startup,
rolled_up,
metadata_path: Some(metadata_path.as_ref().to_path_buf()),
})),
})
}
pub fn verify_causal_consistency(&self) -> Vec<u64> {
let inner = self.inner.read().unwrap_or_else(|e| e.into_inner());
detect_orphan_facts(&inner.facts_log, &inner.fact_sources)
}
pub fn append(
&self,
path: &str,
value: JsonValue,
source_session_id: u64,
) -> Result<u64, FactsLogError> {
let mut inner = self.inner.write().unwrap_or_else(|e| e.into_inner());
let fact_id = FactId(inner.next_fact_id);
inner.next_fact_id = inner
.next_fact_id
.checked_add(1)
.ok_or(FactsLogError::VersionOverflow)?;
let fact = Fact::PayloadUpdate {
id: fact_id,
path: path.to_string(),
value,
};
let version = inner.facts_log.append(fact)?;
inner.fact_sources.insert(fact_id, source_session_id);
inner.persist_metadata_locked();
Ok(version)
}
pub fn facts_by_path_prefix(&self, prefix: &str) -> Vec<SharedFact> {
let inner = self.inner.read().unwrap_or_else(|e| e.into_inner());
inner
.facts_log
.facts_by_path_prefix(prefix)
.into_iter()
.filter(|(_, fact)| matches!(fact, Fact::PayloadUpdate { .. }))
.filter(|(_, fact)| {
if let Fact::PayloadUpdate { id, .. } = fact {
!inner.rolled_up.contains(id)
} else {
false
}
})
.map(|(version, fact)| {
if let Fact::PayloadUpdate { id, path, value } = fact {
let source_session_id = *inner.fact_sources.get(&id).unwrap_or(&0);
SharedFact {
fact_id: id,
path,
value,
source_session_id,
version,
}
} else {
unreachable!()
}
})
.collect()
}
pub fn source_session_id(&self, fact_id: FactId) -> Option<u64> {
let inner = self.inner.read().unwrap_or_else(|e| e.into_inner());
inner.fact_sources.get(&fact_id).copied()
}
pub fn fact_by_id(&self, fact_id: FactId) -> Option<SharedFact> {
let inner = self.inner.read().unwrap_or_else(|e| e.into_inner());
let source_session_id = *inner.fact_sources.get(&fact_id)?;
for (version, fact) in inner.facts_log.history_with_versions() {
if let Fact::PayloadUpdate { id, path, value } = fact {
if id == fact_id {
return Some(SharedFact {
fact_id: id,
path,
value,
source_session_id,
version,
});
}
}
}
None
}
pub fn version(&self) -> u64 {
let inner = self.inner.read().unwrap_or_else(|e| e.into_inner());
inner.facts_log.version()
}
pub fn history_len(&self) -> usize {
let inner = self.inner.read().unwrap_or_else(|e| e.into_inner());
inner.facts_log.history_len()
}
pub fn record_used_at_startup(&self, session_id: u64, fact_ids: &[FactId]) {
let mut inner = self.inner.write().unwrap_or_else(|e| e.into_inner());
inner.used_at_startup.insert(session_id, fact_ids.to_vec());
inner.persist_metadata_locked();
}
pub fn get_used_at_startup(&self, session_id: u64) -> Option<Vec<FactId>> {
let inner = self.inner.read().unwrap_or_else(|e| e.into_inner());
inner.used_at_startup.get(&session_id).cloned()
}
pub fn get_sessions_using_fact(&self, fact_id: FactId) -> Vec<u64> {
let inner = self.inner.read().unwrap_or_else(|e| e.into_inner());
inner
.used_at_startup
.iter()
.filter(|(_, facts)| facts.contains(&fact_id))
.map(|(session_id, _)| *session_id)
.collect()
}
pub fn mark_as_rollup(&self, fact_ids: &[FactId]) {
let mut inner = self.inner.write().unwrap_or_else(|e| e.into_inner());
for id in fact_ids {
inner.rolled_up.insert(*id);
}
inner.persist_metadata_locked();
}
pub fn reset(&self) {
let mut inner = self.inner.write().unwrap_or_else(|e| e.into_inner());
inner.facts_log.reset();
inner.next_fact_id = 1;
inner.fact_sources.clear();
inner.used_at_startup.clear();
inner.rolled_up.clear();
}
}
fn write_metadata_atomic(path: &Path, inner: &SharedFactsLogInner) -> Result<(), FactsLogError> {
let metadata = SharedFactsMetadata {
next_fact_id: inner.next_fact_id,
fact_sources: inner.fact_sources.iter().map(|(k, v)| (k.0, *v)).collect(),
used_at_startup: inner
.used_at_startup
.iter()
.map(|(k, v)| (*k, v.iter().map(|f| f.0).collect()))
.collect(),
rolled_up: inner.rolled_up.iter().map(|f| f.0).collect(),
};
let json = serde_json::to_string(&metadata)
.map_err(|e| FactsLogError::WalError(format!("metadata serialize error: {e}")))?;
let tmp_path = path.with_extension("tmp");
std::fs::write(&tmp_path, json)
.map_err(|e| FactsLogError::WalError(format!("metadata write error: {e}")))?;
std::fs::rename(&tmp_path, path)
.map_err(|e| FactsLogError::WalError(format!("metadata rename error: {e}")))?;
Ok(())
}
impl Default for SharedFactsLog {
fn default() -> Self {
Self::new()
}
}
fn detect_orphan_facts(facts_log: &FactsLog, fact_sources: &BTreeMap<FactId, u64>) -> Vec<u64> {
let wal_fact_ids: BTreeSet<u64> = facts_log
.read_from(0)
.iter()
.filter_map(|f| match f {
Fact::PayloadUpdate { id, .. } => Some(id.0),
_ => None,
})
.collect();
wal_fact_ids
.iter()
.filter(|id| !fact_sources.contains_key(&FactId(**id)))
.copied()
.collect()
}
#[cfg(test)]
mod tests {
#![allow(clippy::unwrap_used)]
#![allow(clippy::panic, clippy::expect_used)]
use super::*;
#[test]
fn test_shared_facts_log_new() {
let log = SharedFactsLog::new();
assert_eq!(log.version(), 0);
assert_eq!(log.history_len(), 0);
}
#[test]
fn test_shared_facts_log_append() {
let log = SharedFactsLog::new();
let version = log
.append("shared.research.note1", JsonValue::string("hello"), 100)
.unwrap();
assert_eq!(version, 1); assert_eq!(log.history_len(), 1);
}
#[test]
fn test_shared_facts_log_facts_by_path_prefix() {
let log = SharedFactsLog::new();
log.append("shared.research.note1", JsonValue::string("v1"), 100)
.unwrap();
log.append("shared.research.note2", JsonValue::string("v2"), 100)
.unwrap();
log.append("shared.other.data", JsonValue::string("v3"), 200)
.unwrap();
let result = log.facts_by_path_prefix("shared.research");
assert_eq!(result.len(), 2);
assert_eq!(result[0].path, "shared.research.note1");
assert_eq!(result[1].path, "shared.research.note2");
assert_eq!(result[0].source_session_id, 100);
assert_eq!(result[1].source_session_id, 100);
}
#[test]
fn test_shared_facts_log_facts_by_path_prefix_no_matches() {
let log = SharedFactsLog::new();
log.append("shared.research.note1", JsonValue::string("v1"), 100)
.unwrap();
let result = log.facts_by_path_prefix("shared.other");
assert!(result.is_empty());
}
#[test]
fn test_shared_facts_log_facts_by_path_prefix_empty() {
let log = SharedFactsLog::new();
let result = log.facts_by_path_prefix("any_prefix");
assert!(result.is_empty());
}
#[test]
fn test_shared_facts_log_source_session_id() {
let log = SharedFactsLog::new();
log.append("shared.note", JsonValue::string("v1"), 100)
.unwrap();
let result = log.facts_by_path_prefix("shared");
assert_eq!(result.len(), 1);
assert_eq!(result[0].source_session_id, 100);
}
#[test]
fn test_shared_facts_log_reset() {
let log = SharedFactsLog::new();
log.append("shared.note", JsonValue::string("v1"), 100)
.unwrap();
assert_eq!(log.history_len(), 1);
log.reset();
assert_eq!(log.history_len(), 0);
assert_eq!(log.version(), 0);
}
#[test]
fn test_mark_as_rollup_filters_from_prefix_query() {
let log = SharedFactsLog::new();
log.append(
"shared.ns.sessions.s1.summary",
JsonValue::string("s1"),
100,
)
.unwrap();
log.append(
"shared.ns.sessions.s2.summary",
JsonValue::string("s2"),
100,
)
.unwrap();
log.append(
"shared.ns.sessions.s3.summary",
JsonValue::string("s3"),
100,
)
.unwrap();
log.mark_as_rollup(&[FactId(1)]);
let result = log.facts_by_path_prefix("shared.ns.sessions.");
assert_eq!(result.len(), 2);
assert_eq!(result[0].fact_id, FactId(2));
assert_eq!(result[1].fact_id, FactId(3));
}
#[test]
fn test_mark_as_rollup_fact_by_id_still_visible() {
let log = SharedFactsLog::new();
log.append(
"shared.ns.sessions.s1.summary",
JsonValue::string("s1"),
100,
)
.unwrap();
log.mark_as_rollup(&[FactId(1)]);
let fact = log.fact_by_id(FactId(1));
assert!(fact.is_some());
assert_eq!(fact.unwrap().path, "shared.ns.sessions.s1.summary");
}
#[test]
fn test_recover_restores_facts_and_metadata() {
use tempfile::tempdir;
let dir = tempdir().unwrap();
let wal_path = dir.path().join("shared.wal");
let meta_path = dir.path().join("shared.json");
{
let log = SharedFactsLog::recover(&wal_path, &meta_path).unwrap();
log.append("shared.ns.stable.key1", JsonValue::string("v1"), 100)
.unwrap();
log.append(
"shared.ns.sessions.s1.summary",
JsonValue::string("s1"),
100,
)
.unwrap();
log.mark_as_rollup(&[FactId(2)]); }
{
let log = SharedFactsLog::recover(&wal_path, &meta_path).unwrap();
assert_eq!(log.history_len(), 2); assert_eq!(log.version(), 2);
assert_eq!(log.source_session_id(FactId(1)), Some(100));
let result = log.facts_by_path_prefix("shared.ns.");
assert_eq!(result.len(), 1);
assert_eq!(result[0].fact_id, FactId(1));
assert!(log.fact_by_id(FactId(2)).is_some());
let v = log
.append("shared.ns.stable.key2", JsonValue::string("v2"), 200)
.unwrap();
assert_eq!(v, 3);
}
}
#[test]
fn test_recover_metadata_not_found_uses_defaults() {
use tempfile::tempdir;
let dir = tempdir().unwrap();
let wal_path = dir.path().join("shared.wal");
let meta_path = dir.path().join("nonexistent.json");
let log = SharedFactsLog::recover(&wal_path, &meta_path).unwrap();
assert_eq!(log.history_len(), 0);
assert_eq!(log.version(), 0);
}
#[test]
fn test_record_used_at_startup_persists_across_recover() {
use tempfile::tempdir;
let dir = tempdir().unwrap();
let wal_path = dir.path().join("shared.wal");
let meta_path = dir.path().join("shared_meta.json");
{
let log = SharedFactsLog::recover(&wal_path, &meta_path).unwrap();
log.append("shared.ns.key", JsonValue::string("v1"), 100)
.unwrap();
log.record_used_at_startup(42, &[FactId(1), FactId(2)]);
log.record_used_at_startup(43, &[FactId(1)]);
}
let log2 = SharedFactsLog::recover(&wal_path, &meta_path).unwrap();
assert_eq!(
log2.get_used_at_startup(42),
Some(vec![FactId(1), FactId(2)]),
"重启后 used_at_startup(42) 丢失:P4-BUG 回归失败"
);
assert_eq!(log2.get_used_at_startup(43), Some(vec![FactId(1)]));
assert_eq!(
log2.get_sessions_using_fact(FactId(2)),
vec![42],
"重启后消费关系查询结果不完整"
);
}
#[test]
fn test_recover_removes_orphan_tmp_file() {
use tempfile::tempdir;
let dir = tempdir().unwrap();
let wal_path = dir.path().join("shared.wal");
let meta_path = dir.path().join("shared_meta.json");
let tmp_path = meta_path.with_extension("tmp");
std::fs::write(&tmp_path, b"{orphan").unwrap();
let _log = SharedFactsLog::recover(&wal_path, &meta_path).unwrap();
assert!(
!tmp_path.exists(),
"孤儿 metadata.tmp 未被 recover 清理:P4-D 回归失败"
);
}
#[test]
fn test_recover_detects_orphan_facts_from_crash_window() {
use tempfile::tempdir;
let dir = tempdir().unwrap();
let wal_path = dir.path().join("shared.wal");
let meta_path = dir.path().join("shared_meta.json");
{
let log = SharedFactsLog::recover(&wal_path, &meta_path).unwrap();
log.append("shared.ns.a", JsonValue::string("v1"), 100)
.unwrap();
log.append("shared.ns.b", JsonValue::string("v2"), 101)
.unwrap();
}
let degraded = serde_json::json!({
"fact_sources": {"1": 100},
"next_fact_id": 2,
"rolled_up": [],
"used_at_startup": {}
});
std::fs::write(&meta_path, degraded.to_string()).unwrap();
let log2 = SharedFactsLog::recover(&wal_path, &meta_path).unwrap();
let orphans = log2.verify_causal_consistency();
assert_eq!(
orphans,
vec![2],
"崩溃窗口产生的孤立 fact 未被检出:F1 回归失败"
);
}
#[test]
fn test_verify_causal_consistency_no_false_positive_when_consistent() {
use tempfile::tempdir;
let dir = tempdir().unwrap();
let wal_path = dir.path().join("shared.wal");
let meta_path = dir.path().join("shared_meta.json");
let log = SharedFactsLog::recover(&wal_path, &meta_path).unwrap();
log.append("shared.ns.a", JsonValue::string("v1"), 100)
.unwrap();
log.append("shared.ns.b", JsonValue::string("v2"), 101)
.unwrap();
assert!(
log.verify_causal_consistency().is_empty(),
"一致的 WAL+metadata 被误报为孤立:F1 对照测试失败"
);
}
}