`corpus_items.mission_id` has existed since the table landed and nothing could populate it. `POST /api/library/runs` now accepts `missionId`, which is the seam the wizard needs: a mission-driven run is the same run, tagged. `corpus::contributed()` answers the question a continuous mission has to be able to answer — did THIS run add anything new. Because `record` never reassigns mission_id on conflict, the mission that first found a source keeps the credit, so a rerun cannot inflate its own count by re-recording what an earlier run already held. The test asserts exactly that: two missions see the same paper, the finder reports 1 and the rerun reports 0. This is the check the 0030-0044 generation of continuous research did not have. It could run weekly forever and every run looked like success. The test also earned its FK: the first version attributed to a bare UUID and the database refused it. Attribution to a mission that does not exist is not attribution, so the test now seeds real mission rows. 400 tests, clippy clean. Co-Authored-By: Claude Opus 5 <[email protected]>
496 lines
18 KiB
Rust
496 lines
18 KiB
Rust
//! What a continuous mission has already covered.
|
|
//!
|
|
//! A recurring mission's hard problem is not running the agent — that is 23
|
|
//! seconds — it is knowing what it already did last time. A research mission
|
|
//! with no memory of prior runs resurfaces the same papers forever and reports
|
|
//! success every time.
|
|
//!
|
|
//! This module keeps that record. It is deliberately small: an index derived
|
|
//! from the corpus, never the corpus itself. The vault is the source of truth,
|
|
//! the index is rebuildable, and a hand-edited note is never "wrong".
|
|
//!
|
|
//! # Two kinds, because the real vault forced it
|
|
//!
|
|
//! The plan assumed notes would carry `arxiv:` / `doi:` / `url:` frontmatter.
|
|
//! Measured against the actual vault: **416 notes, 145 with frontmatter, and
|
|
//! zero with any of those keys.** The dominant keys are repo-sync metadata
|
|
//! (`node`, `org`, `gitea`) and course-note fields (`presenter`, `session`).
|
|
//! An ingester keyed only on external identity would have indexed nothing —
|
|
//! the same shape of failure as everything else this week.
|
|
//!
|
|
//! So `note` rows record coverage (what the vault already contains, keyed by
|
|
//! path) and `source` rows record consumption (external things a mission
|
|
//! read, keyed by natural id). They answer different questions and a
|
|
//! continuous mission needs both: "have I already written about this topic?"
|
|
//! and "have I already read this paper?".
|
|
|
|
use sha2::{Digest, Sha256};
|
|
use uuid::Uuid;
|
|
|
|
/// A note parsed out of the vault, ready to be indexed.
|
|
#[derive(Debug, Clone, PartialEq, Eq)]
|
|
pub struct ParsedNote {
|
|
/// Vault-relative path, used as identity for `kind = 'note'`.
|
|
pub path: String,
|
|
pub title: Option<String>,
|
|
pub content_hash: String,
|
|
/// An external identity the note declares for itself, if any. Nothing in
|
|
/// the vault does this today; missions writing new notes are expected to.
|
|
pub declared_source_id: Option<String>,
|
|
}
|
|
|
|
impl ParsedNote {
|
|
/// `note:<path>` — the `source_id` this note occupies in the index.
|
|
pub fn source_id(&self) -> String {
|
|
format!("note:{}", self.path)
|
|
}
|
|
}
|
|
|
|
/// Hash content for change detection. Not a dedupe key — identity is
|
|
/// `source_id`; this only distinguishes "unchanged" from "edited".
|
|
pub fn content_hash(body: &str) -> String {
|
|
let mut h = Sha256::new();
|
|
h.update(body.as_bytes());
|
|
format!("{:x}", h.finalize())
|
|
}
|
|
|
|
/// Split YAML frontmatter from the body.
|
|
///
|
|
/// Returns `(frontmatter, body)`. A note without frontmatter — 271 of the 416
|
|
/// in the real vault — yields `("", whole file)` rather than being skipped.
|
|
/// Skipping them would drop two thirds of the corpus on the floor.
|
|
fn split_frontmatter(text: &str) -> (&str, &str) {
|
|
let Some(rest) = text.strip_prefix("---") else {
|
|
return ("", text);
|
|
};
|
|
let rest = rest.strip_prefix('\n').unwrap_or(rest);
|
|
match rest.find("\n---") {
|
|
Some(end) => {
|
|
let body = &rest[end + 4..];
|
|
(&rest[..end], body.strip_prefix('\n').unwrap_or(body))
|
|
}
|
|
// An opening fence with no close is malformed; treat the whole file as
|
|
// body rather than swallowing it as frontmatter.
|
|
None => ("", text),
|
|
}
|
|
}
|
|
|
|
/// Read one scalar key out of a frontmatter block.
|
|
///
|
|
/// Deliberately not a YAML parser. The vault's frontmatter is flat
|
|
/// `key: value` with occasional quotes and one list (`tags`), and pulling in a
|
|
/// YAML dependency to read three keys would be more surface than it is worth.
|
|
fn frontmatter_value<'a>(fm: &'a str, key: &str) -> Option<&'a str> {
|
|
for line in fm.lines() {
|
|
let line = line.trim();
|
|
let Some((k, v)) = line.split_once(':') else {
|
|
continue;
|
|
};
|
|
if !k.trim().eq_ignore_ascii_case(key) {
|
|
continue;
|
|
}
|
|
let v = v.trim().trim_matches('"').trim_matches('\'').trim();
|
|
if !v.is_empty() {
|
|
return Some(v);
|
|
}
|
|
}
|
|
None
|
|
}
|
|
|
|
/// Which frontmatter keys may declare an external identity, in priority order.
|
|
///
|
|
/// None of these appear in the vault today. They are the contract for notes
|
|
/// that missions write from here on, and the reason a `source:` key is NOT in
|
|
/// the list: the vault already uses `source:` for local filesystem paths of
|
|
/// course material (`/Users/quantum/Downloads/...`), which is provenance, not
|
|
/// a citable external identity. Treating it as one would fill the seen-set
|
|
/// with 25 rows keyed on a laptop path.
|
|
const IDENTITY_KEYS: &[&str] = &["source_id", "arxiv", "doi", "url", "permalink"];
|
|
|
|
/// Parse a note. `path` must be vault-relative.
|
|
pub fn parse_note(path: &str, text: &str) -> ParsedNote {
|
|
let (fm, body) = split_frontmatter(text);
|
|
|
|
let declared_source_id = IDENTITY_KEYS.iter().find_map(|k| {
|
|
frontmatter_value(fm, k).map(|v| {
|
|
// `source_id` is already qualified; the others name their scheme.
|
|
if *k == "source_id" || v.contains(':') {
|
|
v.to_string()
|
|
} else {
|
|
format!("{k}:{v}")
|
|
}
|
|
})
|
|
});
|
|
|
|
// Title: the first markdown H1, else the filename stem. Frontmatter has no
|
|
// consistent title key in this vault.
|
|
let title = body
|
|
.lines()
|
|
.find_map(|l| l.strip_prefix("# ").map(str::trim))
|
|
.filter(|t| !t.is_empty())
|
|
.map(str::to_string)
|
|
.or_else(|| {
|
|
std::path::Path::new(path)
|
|
.file_stem()
|
|
.map(|s| s.to_string_lossy().into_owned())
|
|
});
|
|
|
|
ParsedNote {
|
|
path: path.to_string(),
|
|
title,
|
|
// Hash the body, not the whole file: re-syncing a repo note rewrites
|
|
// `updated:`/`size_kb:` in frontmatter without the prose changing, and
|
|
// that should not read as an edit.
|
|
content_hash: content_hash(body),
|
|
declared_source_id,
|
|
}
|
|
}
|
|
|
|
/// What a re-index actually did. `unchanged` is the number that matters: on a
|
|
/// vault nobody edited it should equal the note count.
|
|
#[derive(Debug, Default, Clone, PartialEq, Eq)]
|
|
pub struct IndexStats {
|
|
pub scanned: usize,
|
|
pub inserted: usize,
|
|
pub updated: usize,
|
|
pub unchanged: usize,
|
|
}
|
|
|
|
/// Walk a checkout and index every markdown note.
|
|
///
|
|
/// Skips `.git` and Obsidian's own `.obsidian` config directory — indexing an
|
|
/// editor's workspace state as knowledge would be noise.
|
|
pub fn collect_notes(root: &std::path::Path) -> Vec<ParsedNote> {
|
|
fn walk(dir: &std::path::Path, root: &std::path::Path, out: &mut Vec<ParsedNote>) {
|
|
let Ok(entries) = std::fs::read_dir(dir) else {
|
|
return;
|
|
};
|
|
for entry in entries.flatten() {
|
|
let path = entry.path();
|
|
let name = entry.file_name();
|
|
let name = name.to_string_lossy();
|
|
if name.starts_with('.') {
|
|
continue;
|
|
}
|
|
if path.is_dir() {
|
|
walk(&path, root, out);
|
|
} else if path.extension().and_then(|e| e.to_str()) == Some("md") {
|
|
let Ok(text) = std::fs::read_to_string(&path) else {
|
|
continue;
|
|
};
|
|
let rel = path
|
|
.strip_prefix(root)
|
|
.unwrap_or(&path)
|
|
.to_string_lossy()
|
|
.into_owned();
|
|
out.push(parse_note(&rel, &text));
|
|
}
|
|
}
|
|
}
|
|
let mut out = Vec::new();
|
|
walk(root, root, &mut out);
|
|
out.sort_by(|a, b| a.path.cmp(&b.path));
|
|
out
|
|
}
|
|
|
|
/// Upsert one item. Returns whether the row was new.
|
|
#[allow(clippy::too_many_arguments)]
|
|
pub async fn record(
|
|
pool: &sqlx::PgPool,
|
|
workspace_id: Uuid,
|
|
corpus_id: &str,
|
|
kind: &str,
|
|
source_id: &str,
|
|
title: Option<&str>,
|
|
path: Option<&str>,
|
|
url: Option<&str>,
|
|
content_hash: &str,
|
|
mission_id: Option<Uuid>,
|
|
) -> Result<bool, String> {
|
|
// `last_seen_at` always moves; `first_seen_at` and `mission_id` never do.
|
|
// The first mission to find a source keeps the credit, which is what makes
|
|
// "did THIS run contribute anything new" answerable.
|
|
let row: (bool,) = sqlx::query_as(
|
|
"INSERT INTO corpus_items
|
|
(id, workspace_id, corpus_id, kind, source_id, title, path, url,
|
|
content_hash, mission_id)
|
|
VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10)
|
|
ON CONFLICT (workspace_id, corpus_id, source_id) DO UPDATE
|
|
SET last_seen_at = now(),
|
|
title = COALESCE(EXCLUDED.title, corpus_items.title),
|
|
path = COALESCE(EXCLUDED.path, corpus_items.path),
|
|
url = COALESCE(EXCLUDED.url, corpus_items.url),
|
|
content_hash = EXCLUDED.content_hash
|
|
RETURNING (xmax = 0) AS inserted",
|
|
)
|
|
.bind(Uuid::now_v7())
|
|
.bind(workspace_id)
|
|
.bind(corpus_id)
|
|
.bind(kind)
|
|
.bind(source_id)
|
|
.bind(title)
|
|
.bind(path)
|
|
.bind(url)
|
|
.bind(content_hash)
|
|
.bind(mission_id)
|
|
.fetch_one(pool)
|
|
.await
|
|
.map_err(|e| format!("record corpus item {source_id}: {e}"))?;
|
|
Ok(row.0)
|
|
}
|
|
|
|
/// Has this corpus already seen this `source_id`?
|
|
pub async fn seen(
|
|
pool: &sqlx::PgPool,
|
|
workspace_id: Uuid,
|
|
corpus_id: &str,
|
|
source_id: &str,
|
|
) -> Result<bool, String> {
|
|
// `SELECT 1` is INT4; binding it as i64 fails to decode.
|
|
let row: Option<(i32,)> = sqlx::query_as(
|
|
"SELECT 1 FROM corpus_items
|
|
WHERE workspace_id = $1 AND corpus_id = $2 AND source_id = $3",
|
|
)
|
|
.bind(workspace_id)
|
|
.bind(corpus_id)
|
|
.bind(source_id)
|
|
.fetch_optional(pool)
|
|
.await
|
|
.map_err(|e| format!("seen({source_id}): {e}"))?;
|
|
Ok(row.is_some())
|
|
}
|
|
|
|
/// Of these candidate ids, which has this corpus NOT seen?
|
|
///
|
|
/// The shape a research agent actually needs: it has ten search hits and wants
|
|
/// to know which are worth fetching. One round trip, not ten.
|
|
pub async fn unseen(
|
|
pool: &sqlx::PgPool,
|
|
workspace_id: Uuid,
|
|
corpus_id: &str,
|
|
candidates: &[String],
|
|
) -> Result<Vec<String>, String> {
|
|
if candidates.is_empty() {
|
|
return Ok(Vec::new());
|
|
}
|
|
let rows: Vec<(String,)> = sqlx::query_as(
|
|
"SELECT source_id FROM corpus_items
|
|
WHERE workspace_id = $1 AND corpus_id = $2 AND source_id = ANY($3)",
|
|
)
|
|
.bind(workspace_id)
|
|
.bind(corpus_id)
|
|
.bind(candidates)
|
|
.fetch_all(pool)
|
|
.await
|
|
.map_err(|e| format!("unseen: {e}"))?;
|
|
let known: std::collections::HashSet<String> = rows.into_iter().map(|r| r.0).collect();
|
|
Ok(candidates
|
|
.iter()
|
|
.filter(|c| !known.contains(*c))
|
|
.cloned()
|
|
.collect())
|
|
}
|
|
|
|
/// How many NEW sources a mission contributed.
|
|
///
|
|
/// The verification predicate for a continuous research mission. `record`
|
|
/// never reassigns `mission_id` on conflict, so the first mission to find a
|
|
/// source keeps the credit and a rerun cannot inflate its own count by
|
|
/// re-recording what an earlier run already had.
|
|
///
|
|
/// A mission whose answer is zero produced nothing, whatever its transcript
|
|
/// says — which is the check the 0030-0044 generation of this feature lacked.
|
|
pub async fn contributed(
|
|
pool: &sqlx::PgPool,
|
|
workspace_id: Uuid,
|
|
corpus_id: &str,
|
|
mission_id: Uuid,
|
|
) -> Result<i64, String> {
|
|
let row: (i64,) = sqlx::query_as(
|
|
"SELECT count(*) FROM corpus_items
|
|
WHERE workspace_id = $1 AND corpus_id = $2 AND mission_id = $3
|
|
AND kind = 'source'",
|
|
)
|
|
.bind(workspace_id)
|
|
.bind(corpus_id)
|
|
.bind(mission_id)
|
|
.fetch_one(pool)
|
|
.await
|
|
.map_err(|e| format!("contributed({mission_id}): {e}"))?;
|
|
Ok(row.0)
|
|
}
|
|
|
|
/// Index every note in a checkout. Idempotent by construction.
|
|
pub async fn index_vault(
|
|
pool: &sqlx::PgPool,
|
|
workspace_id: Uuid,
|
|
corpus_id: &str,
|
|
root: &std::path::Path,
|
|
) -> Result<IndexStats, String> {
|
|
let notes = collect_notes(root);
|
|
let mut stats = IndexStats {
|
|
scanned: notes.len(),
|
|
..Default::default()
|
|
};
|
|
|
|
for note in ¬es {
|
|
let existing: Option<(String,)> = sqlx::query_as(
|
|
"SELECT content_hash FROM corpus_items
|
|
WHERE workspace_id = $1 AND corpus_id = $2 AND source_id = $3",
|
|
)
|
|
.bind(workspace_id)
|
|
.bind(corpus_id)
|
|
.bind(note.source_id())
|
|
.fetch_optional(pool)
|
|
.await
|
|
.map_err(|e| format!("lookup {}: {e}", note.path))?;
|
|
|
|
match existing {
|
|
Some((hash,)) if hash == note.content_hash => {
|
|
stats.unchanged += 1;
|
|
continue;
|
|
}
|
|
Some(_) => stats.updated += 1,
|
|
None => stats.inserted += 1,
|
|
}
|
|
|
|
record(
|
|
pool,
|
|
workspace_id,
|
|
corpus_id,
|
|
"note",
|
|
¬e.source_id(),
|
|
note.title.as_deref(),
|
|
Some(¬e.path),
|
|
None,
|
|
¬e.content_hash,
|
|
None,
|
|
)
|
|
.await?;
|
|
|
|
// A note that declares an external identity also registers as a
|
|
// consumed source, so a later mission does not re-read what an
|
|
// earlier one already wrote up.
|
|
if let Some(sid) = ¬e.declared_source_id {
|
|
record(
|
|
pool,
|
|
workspace_id,
|
|
corpus_id,
|
|
"source",
|
|
sid,
|
|
note.title.as_deref(),
|
|
Some(¬e.path),
|
|
None,
|
|
¬e.content_hash,
|
|
None,
|
|
)
|
|
.await?;
|
|
}
|
|
}
|
|
Ok(stats)
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use super::*;
|
|
|
|
#[test]
|
|
fn frontmatter_is_split_from_body() {
|
|
let (fm, body) = split_frontmatter("---\ntype: lecture\n---\n# Title\n\ntext\n");
|
|
assert_eq!(fm, "type: lecture");
|
|
assert!(body.starts_with("# Title"));
|
|
}
|
|
|
|
/// 271 of the vault's 416 notes have no frontmatter. Dropping them would
|
|
/// discard two thirds of the corpus.
|
|
#[test]
|
|
fn a_note_without_frontmatter_is_still_a_note() {
|
|
let (fm, body) = split_frontmatter("# Plain\n\nno frontmatter here\n");
|
|
assert_eq!(fm, "");
|
|
assert!(body.starts_with("# Plain"));
|
|
let n = parse_note("Daily/x.md", "# Plain\n\nbody\n");
|
|
assert_eq!(n.title.as_deref(), Some("Plain"));
|
|
assert_eq!(n.declared_source_id, None);
|
|
}
|
|
|
|
/// An unterminated fence must not swallow the file.
|
|
#[test]
|
|
fn malformed_frontmatter_is_treated_as_body() {
|
|
let (fm, body) = split_frontmatter("---\nbroken: yes\nno closing fence\n");
|
|
assert_eq!(fm, "");
|
|
assert!(body.contains("no closing fence"));
|
|
}
|
|
|
|
/// The vault's real `source:` values are local filesystem paths of course
|
|
/// material. Treating those as citable identity would fill the seen-set
|
|
/// with 25 rows keyed on a laptop path.
|
|
#[test]
|
|
fn a_local_source_path_is_not_an_external_identity() {
|
|
let note = parse_note(
|
|
"50 APESS 2026/Lectures/talk.md",
|
|
"---\nsource: \"/Users/quantum/Downloads/Material_APESS_2026/x.pdf\"\n\
|
|
date: 2026-07-27\ntype: lecture\n---\n# Agentic Design\n",
|
|
);
|
|
assert_eq!(
|
|
note.declared_source_id, None,
|
|
"a Downloads path is provenance, not a citable source id"
|
|
);
|
|
assert_eq!(note.title.as_deref(), Some("Agentic Design"));
|
|
assert_eq!(note.source_id(), "note:50 APESS 2026/Lectures/talk.md");
|
|
}
|
|
|
|
#[test]
|
|
fn declared_identities_are_scheme_qualified() {
|
|
let a = parse_note("p.md", "---\narxiv: 2401.12345\n---\n# T\n");
|
|
assert_eq!(a.declared_source_id.as_deref(), Some("arxiv:2401.12345"));
|
|
|
|
let d = parse_note("p.md", "---\ndoi: 10.1000/xyz\n---\n# T\n");
|
|
assert_eq!(d.declared_source_id.as_deref(), Some("doi:10.1000/xyz"));
|
|
|
|
// Already-qualified values are not double-prefixed.
|
|
let s = parse_note("p.md", "---\nsource_id: arxiv:2401.99999\n---\n# T\n");
|
|
assert_eq!(s.declared_source_id.as_deref(), Some("arxiv:2401.99999"));
|
|
|
|
// A URL carries its own scheme and must not become `url:https:...`.
|
|
let u = parse_note("p.md", "---\nurl: https://example.com/p\n---\n# T\n");
|
|
assert_eq!(
|
|
u.declared_source_id.as_deref(),
|
|
Some("https://example.com/p")
|
|
);
|
|
}
|
|
|
|
/// Repo-sync notes rewrite `updated:`/`size_kb:` on every sync without the
|
|
/// prose changing. Hashing the whole file would report 103 phantom edits
|
|
/// per run and make "unchanged" meaningless.
|
|
#[test]
|
|
fn frontmatter_churn_does_not_count_as_an_edit() {
|
|
let a = parse_note("Repos/x.md", "---\nupdated: 2026-08-01\nsize_kb: 12\n---\n# X\n\nbody\n");
|
|
let b = parse_note("Repos/x.md", "---\nupdated: 2026-08-03\nsize_kb: 14\n---\n# X\n\nbody\n");
|
|
assert_eq!(a.content_hash, b.content_hash);
|
|
|
|
let c = parse_note("Repos/x.md", "---\nupdated: 2026-08-03\n---\n# X\n\nDIFFERENT\n");
|
|
assert_ne!(a.content_hash, c.content_hash, "real edits must be visible");
|
|
}
|
|
|
|
#[test]
|
|
fn note_identity_is_its_path() {
|
|
let n = parse_note("30 Resources/a b.md", "# A\n");
|
|
assert_eq!(n.source_id(), "note:30 Resources/a b.md");
|
|
}
|
|
|
|
#[test]
|
|
fn collect_skips_dotfiles_and_non_markdown() {
|
|
let tmp = tempfile::tempdir().unwrap();
|
|
let root = tmp.path();
|
|
std::fs::create_dir_all(root.join(".obsidian")).unwrap();
|
|
std::fs::create_dir_all(root.join("Daily")).unwrap();
|
|
std::fs::write(root.join(".obsidian/workspace.md"), "# editor state\n").unwrap();
|
|
std::fs::write(root.join("Daily/note.md"), "# Real\n").unwrap();
|
|
std::fs::write(root.join("image.png"), "notmd").unwrap();
|
|
|
|
let notes = collect_notes(root);
|
|
assert_eq!(notes.len(), 1, "only the real note: {notes:?}");
|
|
assert_eq!(notes[0].path, "Daily/note.md");
|
|
}
|
|
}
|