use ijima_core::{NamespaceId, Result, Store};
use ijima_miner::{Extraction, MiningContext};
#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct MiningReport {
pub archived: usize,
pub queued: usize,
}
pub async fn ingest_extractions(
store: &dyn Store,
ns: &NamespaceId,
extractions: Vec<Extraction>,
) -> Result<MiningReport> {
let mut report = MiningReport::default();
for extraction in extractions {
match extraction {
Extraction::Auto(memory) => {
store.store_memory(ns, memory).await?;
report.archived += 1;
}
Extraction::PendingReview(memory) => {
store.enqueue_extraction(ns, memory, 0.5).await?;
report.queued += 1;
}
Extraction::Nothing => {}
}
}
Ok(report)
}
pub fn mining_context(
session_id: &str,
project: &str,
harness: ijima_core::harness::Harness,
) -> MiningContext {
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs().to_string())
.unwrap_or_default();
MiningContext {
session_id: session_id.to_string(),
project: project.to_string(),
harness,
now,
}
}
#[cfg(test)]
mod tests {
use super::*;
use ijima_core::{Memory, MemoryId, MemorySource};
use ijima_miner::mine;
fn sample_memory(id: &str, content: &str) -> Memory {
Memory {
id: MemoryId(id.into()),
content: content.into(),
project: "ijima".into(),
topic: "decisions".into(),
source: MemorySource::Mined,
harness: ijima_core::harness::Harness::Pi,
session_id: Some("sess_1".into()),
origin: ijima_core::InstanceId::local(),
authority: ijima_core::AuthorityScope::local(),
importance: 0.7,
created_at: "0".into(),
}
}
#[tokio::test]
async fn ingest_archives_auto_and_queues_pending() {
let store = crate::SurrealStore::open_embedded().await.expect("open");
let ns = NamespaceId::new("ns_trigger");
let extractions = vec![
Extraction::Auto(sample_memory("m1", "auto decision")),
Extraction::PendingReview(sample_memory("m2", "pending fact")),
Extraction::Nothing,
];
let report = ingest_extractions(&store, &ns, extractions)
.await
.expect("ingest");
assert_eq!(
report,
MiningReport {
archived: 1,
queued: 1
}
);
let got = store
.recall_memory(&ns, &MemoryId("m1".into()))
.await
.expect("recall");
assert!(got.is_some());
let pending = store.list_pending(&ns, 10).await.expect("list");
assert_eq!(pending.len(), 1);
}
#[tokio::test]
async fn full_pipeline_mine_then_ingest() {
let store = crate::SurrealStore::open_embedded().await.expect("open");
let ns = NamespaceId::new("ns_e2e");
let ctx = mining_context("sess_1", "ijima", ijima_core::harness::Harness::Pi);
let turns = vec![
"We decided to use SurrealDB for storage.".to_string(),
"See https://example.com/loader for the candle loader.".to_string(),
];
let extractions = mine(&turns, &ctx).expect("mine");
assert_eq!(extractions.len(), 2);
let report = ingest_extractions(&store, &ns, extractions)
.await
.expect("ingest");
assert_eq!(report.archived, 2);
assert_eq!(report.queued, 0);
}
}