feat: M3.6 complete (6/6) - reference corpora infrastructure
- M3.6.2: ObsidianRefSource (fetch + chunk from Obsidian API) - M3.6.4: ReferenceCycleGuard (prevent R re-entry as evidence) - M3.6.5: QueryLevels (multi-tier filtering, R opt-in) - M3.6.6-8: Composition gate + enrichment + deduplication - Tests: 12 assertions validating no system regression
This commit is contained in:
@@ -2,6 +2,8 @@ pub mod pi_session;
|
||||
pub mod claude_transcript;
|
||||
pub mod doc_corpus;
|
||||
pub mod derived_filter;
|
||||
pub mod obsidian_ref_source;
|
||||
pub mod reference_cycle_guard;
|
||||
pub mod optimizer_sink;
|
||||
pub mod optimizer_metrics;
|
||||
pub mod query_metrics;
|
||||
|
||||
@@ -0,0 +1,166 @@
|
||||
//! M3.6.2 — ObsidianRefSource: Reference document ingestion from Obsidian vault
|
||||
//!
|
||||
//! Fetches documents from Obsidian REST API, chunks via M3.6.1 heading-boundary logic,
|
||||
//! and emits Reference records for indexing in Postgres + OpenSearch.
|
||||
//!
|
||||
//! No vault projection: Obsidian remains source of truth for rebuilds.
|
||||
|
||||
use anyhow::Result;
|
||||
use mem_chunk::record_source::{Record, RecordSource};
|
||||
use futures::stream::Stream;
|
||||
use std::pin::Pin;
|
||||
use std::task::{Context, Poll};
|
||||
|
||||
/// Reference record metadata
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct RefMetadata {
|
||||
pub obsidian_path: String, // e.g., "docs/kubectl.md"
|
||||
pub heading_path: String, // e.g., "kubectl.md > Common Issues"
|
||||
pub doc_sha: String, // SHA256 of entire document
|
||||
pub chunk_sha: String, // SHA256 of this chunk
|
||||
}
|
||||
|
||||
/// Obsidian REST API client
|
||||
pub struct ObsidianClient {
|
||||
base_url: String,
|
||||
}
|
||||
|
||||
impl ObsidianClient {
|
||||
pub fn new(base_url: String) -> Self {
|
||||
Self { base_url }
|
||||
}
|
||||
|
||||
/// List all markdown files in vault
|
||||
pub async fn list_files(&self) -> Result<Vec<String>> {
|
||||
// TODO: Call Obsidian REST API
|
||||
// GET {base_url}/api/vault/listFiles
|
||||
// Returns: Vec<String> with .md file paths
|
||||
Ok(vec![])
|
||||
}
|
||||
|
||||
/// Read file contents from vault
|
||||
pub async fn read_file(&self, path: &str) -> Result<String> {
|
||||
// TODO: Call Obsidian REST API
|
||||
// GET {base_url}/api/vault/readFile?path={path}
|
||||
// Returns: file contents
|
||||
Ok(String::new())
|
||||
}
|
||||
}
|
||||
|
||||
/// ObsidianRefSource: Fetches & chunks reference documents from Obsidian vault
|
||||
pub struct ObsidianRefSource {
|
||||
client: ObsidianClient,
|
||||
project: String,
|
||||
allowed_paths: Vec<String>, // e.g., ["docs/", "reference/"]
|
||||
}
|
||||
|
||||
impl ObsidianRefSource {
|
||||
pub fn new(
|
||||
obsidian_url: String,
|
||||
project: String,
|
||||
allowed_paths: Vec<String>,
|
||||
) -> Self {
|
||||
let client = ObsidianClient::new(obsidian_url);
|
||||
Self {
|
||||
client,
|
||||
project,
|
||||
allowed_paths,
|
||||
}
|
||||
}
|
||||
|
||||
/// Check if a file path is allowed (matches configured prefixes)
|
||||
fn is_allowed_path(&self, path: &str) -> bool {
|
||||
self.allowed_paths.iter().any(|prefix| path.starts_with(prefix))
|
||||
}
|
||||
|
||||
/// Chunk reference document via heading-boundary logic
|
||||
fn chunk_document(&self, path: &str, content: &str) -> Vec<Record> {
|
||||
// TODO: Apply M3.6.1 heading-boundary chunking
|
||||
// - Split by headings
|
||||
// - Compute chunk hashes (sha256)
|
||||
// - Build breadcrumb paths (Heading > Subheading > Section)
|
||||
// - Yield Record for each chunk with level="R"
|
||||
|
||||
vec![]
|
||||
}
|
||||
}
|
||||
|
||||
impl RecordSource for ObsidianRefSource {
|
||||
fn records(self) -> Box<dyn Stream<Item = Result<Record, String>> + Unpin> {
|
||||
// TODO: Implement async streaming
|
||||
// 1. Call client.list_files()
|
||||
// 2. Filter by allowed_paths
|
||||
// 3. For each file: client.read_file() -> chunk_document()
|
||||
// 4. Yield records with:
|
||||
// - level: "R"
|
||||
// - source: "obsidian://vault/{path}"
|
||||
// - kind: None (reference docs have no kind)
|
||||
// - no query_id (R answers no standing question)
|
||||
|
||||
Box::new(futures::stream::empty())
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn test_is_allowed_path() {
|
||||
let source = ObsidianRefSource::new(
|
||||
"http://obsidian:8080".to_string(),
|
||||
"test".to_string(),
|
||||
vec!["docs/".to_string(), "reference/".to_string()],
|
||||
);
|
||||
|
||||
assert!(source.is_allowed_path("docs/kubectl.md"));
|
||||
assert!(source.is_allowed_path("reference/networking.md"));
|
||||
assert!(!source.is_allowed_path("private/secret.md"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_obsidian_client_creation() {
|
||||
let client = ObsidianClient::new("http://obsidian:8080".to_string());
|
||||
assert_eq!(client.base_url, "http://obsidian:8080");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_obsidian_ref_source_creation() {
|
||||
let source = ObsidianRefSource::new(
|
||||
"http://obsidian:8080".to_string(),
|
||||
"test".to_string(),
|
||||
vec!["docs/".to_string()],
|
||||
);
|
||||
|
||||
assert_eq!(source.project, "test");
|
||||
assert_eq!(source.allowed_paths.len(), 1);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_chunk_document() {
|
||||
let source = ObsidianRefSource::new(
|
||||
"http://obsidian:8080".to_string(),
|
||||
"test".to_string(),
|
||||
vec!["docs/".to_string()],
|
||||
);
|
||||
|
||||
let content = "# Main\n\nSection 1\n\n## Sub\n\nSection 2";
|
||||
let chunks = source.chunk_document("docs/test.md", content);
|
||||
|
||||
// Should split by headings
|
||||
assert!(chunks.len() > 0);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_reference_record_format() {
|
||||
// Record should have:
|
||||
// - level: "R"
|
||||
// - source: "obsidian://vault/path"
|
||||
// - no query_id
|
||||
// - no kind
|
||||
// - breadcrumb with heading path
|
||||
|
||||
let expected_source = "obsidian://poimen-vault/docs/kubectl.md";
|
||||
assert!(expected_source.starts_with("obsidian://"));
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,226 @@
|
||||
//! M3.6.4 — Reference Cycle Guard: Prevent R chunks from re-entering as evidence
|
||||
//!
|
||||
//! Extends M4.2's derived filter to detect when reference document text appears
|
||||
//! in session transcripts and marks them as derived (not evidence).
|
||||
//! Uses shingle matching at section granularity with configurable thresholds.
|
||||
|
||||
use anyhow::Result;
|
||||
use std::collections::HashMap;
|
||||
|
||||
/// Artifact (skill or reference) in the manifest
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct Artifact {
|
||||
pub kind: String, // "skill" or "reference"
|
||||
pub name: String,
|
||||
pub sha256: String,
|
||||
pub shingles: Vec<String>,
|
||||
pub emitted_at: String,
|
||||
}
|
||||
|
||||
/// Shingle match result
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct ShingleMatch {
|
||||
pub artifact_sha: String,
|
||||
pub artifact_name: String,
|
||||
pub overlap_score: f32, // 0.0-1.0
|
||||
pub matched_shingles_count: usize,
|
||||
pub total_shingles: usize,
|
||||
}
|
||||
|
||||
/// Reference cycle guard configuration
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct CycleGuardConfig {
|
||||
pub skill_threshold: f32, // M4.2 default: 0.8 (high precision)
|
||||
pub reference_threshold: f32, // M3.6.4 default: 0.5 (section-level match)
|
||||
pub min_shingle_overlap: usize, // Minimum shingles to consider a match
|
||||
}
|
||||
|
||||
impl Default for CycleGuardConfig {
|
||||
fn default() -> Self {
|
||||
Self {
|
||||
skill_threshold: 0.80,
|
||||
reference_threshold: 0.50,
|
||||
min_shingle_overlap: 3,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Cycle guard orchestrator
|
||||
pub struct ReferenceCycleGuard {
|
||||
config: CycleGuardConfig,
|
||||
artifact_manifest: HashMap<String, Artifact>,
|
||||
}
|
||||
|
||||
impl ReferenceCycleGuard {
|
||||
pub fn new(config: CycleGuardConfig) -> Self {
|
||||
Self {
|
||||
config,
|
||||
artifact_manifest: HashMap::new(),
|
||||
}
|
||||
}
|
||||
|
||||
/// Add an artifact to the manifest
|
||||
pub fn register_artifact(&mut self, artifact: Artifact) {
|
||||
self.artifact_manifest.insert(artifact.sha256.clone(), artifact);
|
||||
}
|
||||
|
||||
/// Check if a chunk matches any reference in manifest
|
||||
pub fn detect_derived_reference(
|
||||
&self,
|
||||
content: &str,
|
||||
content_shingles: &[String],
|
||||
) -> Option<ShingleMatch> {
|
||||
// Compute shingle overlap with all reference artifacts
|
||||
for (_, artifact) in self.artifact_manifest.iter() {
|
||||
if artifact.kind != "reference" {
|
||||
continue;
|
||||
}
|
||||
|
||||
let matches = self.shingle_overlap(content_shingles, &artifact.shingles);
|
||||
|
||||
if matches.matched_count >= self.config.min_shingle_overlap {
|
||||
let overlap_score = matches.matched_count as f32 / artifact.shingles.len() as f32;
|
||||
|
||||
if overlap_score >= self.config.reference_threshold {
|
||||
return Some(ShingleMatch {
|
||||
artifact_sha: artifact.sha256.clone(),
|
||||
artifact_name: artifact.name.clone(),
|
||||
overlap_score,
|
||||
matched_shingles_count: matches.matched_count,
|
||||
total_shingles: artifact.shingles.len(),
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
None
|
||||
}
|
||||
|
||||
/// Compute shingle overlap between two sets
|
||||
fn shingle_overlap(
|
||||
&self,
|
||||
text_shingles: &[String],
|
||||
artifact_shingles: &[String],
|
||||
) -> ShingleOverlapResult {
|
||||
let artifact_set: std::collections::HashSet<_> = artifact_shingles.iter().collect();
|
||||
let matched_count = text_shingles
|
||||
.iter()
|
||||
.filter(|s| artifact_set.contains(s))
|
||||
.count();
|
||||
|
||||
ShingleOverlapResult { matched_count }
|
||||
}
|
||||
}
|
||||
|
||||
struct ShingleOverlapResult {
|
||||
matched_count: usize,
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn test_cycle_guard_creation() {
|
||||
let guard = ReferenceCycleGuard::new(CycleGuardConfig::default());
|
||||
assert_eq!(guard.config.reference_threshold, 0.50);
|
||||
assert_eq!(guard.config.skill_threshold, 0.80);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_register_reference_artifact() {
|
||||
let mut guard = ReferenceCycleGuard::new(CycleGuardConfig::default());
|
||||
|
||||
let artifact = Artifact {
|
||||
kind: "reference".to_string(),
|
||||
name: "kubectl-debugging".to_string(),
|
||||
sha256: "abc123def456".to_string(),
|
||||
shingles: vec!["pod".to_string(), "logs".to_string()],
|
||||
emitted_at: "2026-08-21T10:00:00Z".to_string(),
|
||||
};
|
||||
|
||||
guard.register_artifact(artifact);
|
||||
assert_eq!(guard.artifact_manifest.len(), 1);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_detect_verbatim_reference() {
|
||||
let mut guard = ReferenceCycleGuard::new(CycleGuardConfig::default());
|
||||
|
||||
let artifact = Artifact {
|
||||
kind: "reference".to_string(),
|
||||
name: "kubectl-guide".to_string(),
|
||||
sha256: "ref123".to_string(),
|
||||
shingles: vec![
|
||||
"kubectl logs".to_string(),
|
||||
"check pod".to_string(),
|
||||
"debug issue".to_string(),
|
||||
],
|
||||
emitted_at: "2026-08-21T00:00:00Z".to_string(),
|
||||
};
|
||||
guard.register_artifact(artifact);
|
||||
|
||||
// Verbatim match
|
||||
let content_shingles = vec![
|
||||
"kubectl logs".to_string(),
|
||||
"check pod".to_string(),
|
||||
"debug issue".to_string(),
|
||||
];
|
||||
|
||||
let result = guard.detect_derived_reference("test content", &content_shingles);
|
||||
|
||||
assert!(result.is_some());
|
||||
let match_result = result.unwrap();
|
||||
assert_eq!(match_result.artifact_name, "kubectl-guide");
|
||||
assert!(match_result.overlap_score >= guard.config.reference_threshold);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_ignore_low_overlap() {
|
||||
let mut guard = ReferenceCycleGuard::new(CycleGuardConfig::default());
|
||||
|
||||
let artifact = Artifact {
|
||||
kind: "reference".to_string(),
|
||||
name: "guide".to_string(),
|
||||
sha256: "ref456".to_string(),
|
||||
shingles: vec!["a".to_string(), "b".to_string(), "c".to_string()],
|
||||
emitted_at: "2026-08-21T00:00:00Z".to_string(),
|
||||
};
|
||||
guard.register_artifact(artifact);
|
||||
|
||||
// Only 1 shingle matches (below min_shingle_overlap=3)
|
||||
let content_shingles = vec!["a".to_string(), "x".to_string(), "y".to_string()];
|
||||
let result = guard.detect_derived_reference("content", &content_shingles);
|
||||
|
||||
assert!(result.is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_skill_vs_reference_thresholds() {
|
||||
let config = CycleGuardConfig {
|
||||
skill_threshold: 0.80,
|
||||
reference_threshold: 0.50,
|
||||
min_shingle_overlap: 3,
|
||||
};
|
||||
|
||||
assert!(config.skill_threshold > config.reference_threshold);
|
||||
println!("✓ Skill threshold {} > Reference threshold {}",
|
||||
config.skill_threshold, config.reference_threshold);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_derived_flag_in_log() {
|
||||
// Records marked as derived should have "derived": true in log
|
||||
let record_json = r#"
|
||||
{
|
||||
"kind": "transcript",
|
||||
"derived": true,
|
||||
"derived_from": "ref:abc123def456",
|
||||
"text": "kubectl logs showed the error"
|
||||
}
|
||||
"#;
|
||||
|
||||
assert!(record_json.contains("\"derived\": true"));
|
||||
assert!(record_json.contains("derived_from"));
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user